From 6747f6f8a0e28b0554032e9ff1b2a09c970840b7 Mon Sep 17 00:00:00 2001 From: User Date: Wed, 15 Jul 2026 01:37:42 +0200 Subject: [PATCH] fix: serialize session lifecycle transitions --- .superpowers/sdd/predeploy-fix-report.md | 313 +++++++++ backend/src/pi/pi-process-manager.ts | 15 +- backend/src/routes/sessions.ts | 106 ++- backend/src/sse/sse-hub.ts | 25 +- backend/test/pi-process-manager.test.ts | 21 + backend/test/routes-sessions.test.ts | 603 +++++++++++++++++- backend/test/sse-hub.test.ts | 63 +- backend/test/sse-route.test.ts | 62 ++ .../src/shell/AppShell.session-mgmt.test.tsx | 230 ++++++- frontend/src/shell/AppShell.tsx | 65 +- 10 files changed, 1402 insertions(+), 101 deletions(-) diff --git a/.superpowers/sdd/predeploy-fix-report.md b/.superpowers/sdd/predeploy-fix-report.md index 93f2ccb1..98968698 100644 --- a/.superpowers/sdd/predeploy-fix-report.md +++ b/.superpowers/sdd/predeploy-fix-report.md @@ -496,6 +496,9 @@ no output, exit 0 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 in-process cold same-id Resume cannot reuse ids. Permanent DELETE removes the counter via `forget`. +- The Delete-then-Resume adversarial route test proves the deleted session is not resurrected but + currently receives the runner's generic HTTP 500 when `sessionShow` can no longer find it. A + future API cleanup can normalize that missing-session response to 404 or 409. - Frontend tests still print pre-existing MSW unhandled-request and React ref/`act` warnings even 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. @@ -614,3 +617,313 @@ 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. + +## Stale-bootstrap, lifecycle-lock, and competing-Resume hardening + +Date: 2026-07-15 +Base: `08b1f4909e8eb7538156cecc2e7a6cafb46ddfc7` + +This follow-up closes asynchronous identity/order and multi-client transport gaps found in the +pre-deployment review: + +- `PiProcessManager.teardownIfCurrent(id, runtime)` makes teardown an identity-checked operation. + Bootstrap re-checks identity after configuration/retrieval and before both the public + `Starting model` event and model start. Its failure continuation acquires the same session + lifecycle lock, claims only its own runtime identity, and holds serialization through persisted + failure and the public terminal sequence. A continuation left behind by Close or DELETE cannot + target a replacement or recreate forgotten SSE state. +- The former Resume-only promise tail is now a per-session lifecycle lock shared by Resume, Close, + and DELETE. Each route reads the current runtime inside the lock immediately before replacement + or removal and uses identity-checked teardown. Deferred route tests prove both orderings: + Resume then Close/Delete finishes removed with no post-removal bootstrap event; Close then Resume + creates only after Close completes; DELETE then Resume cannot recreate a deleted session. +- AppShell assigns each Resume invocation a monotonic token and records the latest target. A + completion for a different, superseding session id cannot reset the store, select a source, close + the panel, or repaint phase from a late manifest. Same-id invocations are per-target single-flight + operations through the POST and local binding commit: repeated pre-commit clicks update the + shared operation's latest token but issue no second POST or commit path. The operation becomes + joinable again before its manifest fetch, whose repaint remains token/id/selection guarded. Start + new, Stop, streamed session exit, and active-session deletion invalidate pending Resume work. + This prevents stale-source preservation and reverse/non-Resume intent overwrite without allowing + a slow manifest to suppress a later explicit rebind. +- `SseHub` subscriber registrations now carry idempotent transport-close callbacks. `clear` and + `forget` snapshot and actively close every response before discarding runtime transport state; + callback-driven unsubscription during that iteration is safe. The SSE route ends its response so + native EventSource reconnects with `Last-Event-ID`. Post-clear events retain monotonic ids and are + buffered for replay; `forget` additionally resets the id state. + +Production files: + +- `backend/src/pi/pi-process-manager.ts` +- `backend/src/routes/sessions.ts` +- `backend/src/sse/sse-hub.ts` +- `frontend/src/shell/AppShell.tsx` + +Regression tests: + +- `backend/test/pi-process-manager.test.ts` +- `backend/test/routes-sessions.test.ts` +- `backend/test/sse-hub.test.ts` +- `backend/test/sse-route.test.ts` +- `frontend/src/shell/AppShell.session-mgmt.test.tsx` + +### TDD RED/GREEN evidence + +Runtime identity API RED: + +```text +cd backend && npx vitest run test/pi-process-manager.test.ts -t "identity-checked teardown" + +Test Files 1 failed (1) +Tests 1 failed | 39 skipped (40) +TypeError: mgr.teardownIfCurrent is not a function +``` + +Runtime identity API GREEN: + +```text +Test Files 1 passed (1) +Tests 1 passed | 39 skipped (40) +``` + +Deferred bootstrap RED: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts \ + -t "stale bootstrap|bootstrap that" + +Test Files 1 failed (1) +Tests 6 failed | 35 skipped (41) + +close/delete + replacement: stale continuation removed the replacement runtime +delete without replacement: stale continuation called failSession after forget +``` + +Deferred bootstrap GREEN: + +```text +Test Files 1 passed (1) +Tests 6 passed | 35 skipped (41) +``` + +Shared lifecycle ordering RED: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts \ + -t "Resume followed|Close followed|Delete followed" + +Test Files 1 failed (1) +Tests 4 failed | 41 skipped (45) + +All four deferred assertions observed the competing route settle before the first lifecycle +operation released. +``` + +Bootstrap plus lifecycle GREEN: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts \ + -t "Resume followed|Close followed|Delete followed|stale bootstrap|bootstrap that" + +Test Files 1 passed (1) +Tests 10 passed | 35 skipped (45) +``` + +Competing frontend Resume RED: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + -t "competing Resume|stale Resume manifest" + +Test Files 1 failed (1) +Tests 2 failed | 16 skipped (18) + +reverse POST completion opened a second, stale EventSource +late s1 manifest repainted the selected s3 phase from F3 to F7 +``` + +Competing and same-id Resume GREEN: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + -t "competing Resume|stale Resume manifest|false then true" + +Test Files 1 passed (1) +Tests 3 passed | 15 skipped (18) +``` + +### Independent-review hardening RED/GREEN + +The first final review reported no Critical findings and three Important edge cases: bootstrap +could start during an in-progress Close; bootstrap-owned failure was persisted twice; and an older +same-id result could overwrite newer state. The integrated reviewer also required non-Resume +navigation to invalidate pending Resume work. The final main review tightened the same-ID contract +to true single-flight so a second same-target click cannot preserve a dead pre-restart source. + +Backend review RED: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts \ + -t "Close suppresses|bootstrap failure persists once" + +Test Files 1 failed (1) +Tests 2 failed | 45 skipped (47) + +deferred configure started Pi while closeSession was still pending +bootstrap/public failure called failSession twice +``` + +Backend review GREEN: + +```text +Test Files 1 passed (1) +Tests 2 passed | 45 skipped (47) +``` + +Same-id single-flight RED: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + -t "share one cold request" + +Test Files 1 failed (1) +Tests 1 failed | 18 skipped (19) + +two concurrent same-ID invocations issued two cold POSTs (three total including initial activation) +``` + +Non-Resume invalidation RED: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + -t "starting a new question invalidates" + +Test Files 1 failed (1) +Tests 1 failed | 19 skipped (20) + +the late Resume opened an EventSource after Start new returned to the landing state +``` + +Frontend review GREEN: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + -t "share one cold request|competing Resume|stale Resume manifest|starting a new question invalidates" + +Test Files 1 passed (1) +Tests 4 passed | 15 skipped (19) +``` + +Post-commit single-flight lifetime RED: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + -t "releases same-id single-flight" + +Test Files 1 failed (1) +Tests 1 failed | 19 skipped (20) + +s1 committed and waited on its manifest; after s3 superseded it, a new s1 Resume reused the old +operation and issued no second s1 POST (expected 2, received 1). +``` + +Same-id and manifest lifetime GREEN: + +```text +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx -t "same-id|manifest" +Test Files 1 passed (1) +Tests 4 passed | 16 skipped (20) + +cd frontend && npx tsc -b +no output, exit 0 +``` + +Multi-client SSE disconnect RED: + +```text +cd backend && npx vitest run test/sse-hub.test.ts test/sse-route.test.ts + +Test Files 2 failed (2) +Tests 3 failed | 6 passed (9) + +clear/forget invoked zero of two registered close callbacks, and two live HTTP SSE responses timed +out instead of reaching EOF after clear. +``` + +Multi-client SSE disconnect GREEN: + +```text +cd backend && npx vitest run test/sse-hub.test.ts test/sse-route.test.ts +Test Files 2 passed (2) +Tests 9 passed (9) + +cd backend && npx tsc --noEmit -p . +no output, exit 0 +``` + +The Hub tests use two subscribers whose close callbacks immediately unsubscribe themselves, proving +safe snapshot iteration and exactly-once closure. The live-route test opens two HTTP streams, proves +both receive EOF on clear, publishes a new event and gate, then reconnects after id 1 and replays +exactly ids 2 and 3. The forget test closes both subscribers and proves the next id resets to 1. + +Close now removes the observed runtime identity before awaiting persistence. Failure persistence is +claimed once per runtime and lifecycle-serialized; bootstrap's public `session_failed` cannot start +a duplicate. A per-target in-flight map owns the only same-ID POST and commit while its mutable +latest token keeps s1→s2→s1 ordering correct; it is removed immediately after the binding commit, +before awaiting the independently guarded manifest. One shared invalidation helper is called when +active deletion, streamed exit, Start new, or Stop begins. + +### Focused verification + +```text +cd backend && npx vitest run test/routes-sessions.test.ts test/pi-process-manager.test.ts \ + test/sse-hub.test.ts test/sse-route.test.ts +Test Files 4 passed (4) +Tests 96 passed (96) + +cd backend && npx tsc --noEmit -p . +no output, exit 0 + +cd frontend && npx vitest run src/shell/AppShell.session-mgmt.test.tsx \ + src/shell/AppShell.new-session.test.tsx src/stream/useSessionStream.test.tsx +Test Files 3 passed (3) +Tests 36 passed (36) + +cd frontend && npx tsc -b +no output, exit 0 +``` + +### Full verification + +```text +cd backend && npx vitest run +Test Files 22 passed (22) +Tests 202 passed (202) + +cd frontend && npx vitest run +Test Files 43 passed (43) +Tests 280 passed (280) + +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 frontend suite/build retain the previously documented MSW, React ref/`act`, experimental type +stripping, and large-chunk warnings. No warning class was introduced by this wave. No harness, +workflow, persistence, SQL/CTE viewer, model-selection, or deployment file changed. The four +pre-existing modified `.superpowers/sdd/{progress,task-2-report,task-3-report,task-4-report}.md` +files remain excluded from staging. + +### Final independent-review verdict + +After the multi-client transport fix, the independent reviewer reported no Critical, Important, or +Minor findings. Its own focused verification passed 96 backend transport/lifecycle tests, 31 +frontend Resume/stream tests, both TypeScript checks, and `git diff --check`. Final assessment: +**Ready to deploy: Yes.** diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index ae9d3461..23fbbc16 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -174,9 +174,16 @@ export class PiProcessManager { teardown(id: string): void { const rt = this.runtimes.get(id); - if (rt) { - rt.child.kill(); - this.runtimes.delete(id); - } + if (rt) this.teardownIfCurrent(id, rt); + } + + /** Remove only the runtime identity the caller observed. */ + teardownIfCurrent(id: string, expected: SessionRuntime): boolean { + if (this.runtimes.get(id) !== expected) return false; + // Delete before signalling the child so its asynchronous exit cannot be mistaken for a + // crash, and so a replacement installed by a later lifecycle operation is never targeted. + this.runtimes.delete(id); + expected.child.kill(); + return true; } } diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 7dac44b4..66f47b88 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -27,24 +27,27 @@ export function sessionRoutes( app: FastifyInstance, d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager }, ) { - const resumeTails = new Map>(); + const lifecycleTails = new Map>(); const boundRuntimes = new Map< string, ReturnType >(); + const failurePersistenceClaimed = new WeakSet< + ReturnType + >(); - const withResumeLock = async (id: string, work: () => Promise): Promise => { - const previous = resumeTails.get(id) ?? Promise.resolve(); + const withSessionLifecycle = async (id: string, work: () => Promise): Promise => { + const previous = lifecycleTails.get(id) ?? Promise.resolve(); let release!: () => void; const gate = new Promise((resolve) => { release = resolve; }); const tail = previous.then(() => gate); - resumeTails.set(id, tail); + lifecycleTails.set(id, tail); await previous; try { return await work(); } finally { release(); - if (resumeTails.get(id) === tail) resumeTails.delete(id); + if (lifecycleTails.get(id) === tail) lifecycleTails.delete(id); } }; @@ -61,7 +64,17 @@ export function sessionRoutes( // 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); + if (!failurePersistenceClaimed.has(rt)) { + failurePersistenceClaimed.add(rt); + void withSessionLifecycle(id, async () => { + // agent_end can release the old binding before this queued work acquires the lock. + // Undefined means no replacement; a different identity means Resume won and the + // old failure must not touch its manifest. + const bound = boundRuntimes.get(id); + if (bound !== undefined && bound !== rt) return; + await d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined); + }).catch(() => undefined); + } } d.hub.publish(id, e.type, e); if ( @@ -91,14 +104,25 @@ export function sessionRoutes( try { if (retrieval) info(id, "Preparing retrieval context"); await Promise.all([configure, retrieval]); + if (d.mgr.get(id) !== rt) return; info(id, "Starting model"); + if (d.mgr.get(id) !== rt) return; start(); } catch { - d.mgr.teardown(id); - void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined); - rt.bridge.emitClientEvent({ type: "info", level: "error", text: BOOTSTRAP_FAILURE_MESSAGE }); - rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" }); - rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" }); + void withSessionLifecycle(id, async () => { + // A bootstrap continuation can settle after Close/Delete or after a replacement was + // installed. Claim only the runtime identity that actually failed; holding the same + // lifecycle lock through persistence prevents a Resume from becoming that failure's + // accidental target. + if (d.mgr.get(id) !== rt || !d.mgr.teardownIfCurrent(id, rt)) return; + if (!failurePersistenceClaimed.has(rt)) { + failurePersistenceClaimed.add(rt); + await d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined); + } + rt.bridge.emitClientEvent({ type: "info", level: "error", text: BOOTSTRAP_FAILURE_MESSAGE }); + rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" }); + rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" }); + }); } })(); }; @@ -158,7 +182,7 @@ export function sessionRoutes( }); app.post("/sessions/:id/resume", async (req, reply) => { const id = (req.params as any).id; - return withResumeLock(id, async () => { + return withSessionLifecycle(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); @@ -192,18 +216,31 @@ export function sessionRoutes( return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE }); } + // Re-check immediately before the replacement commit. A lifecycle operation that ran + // before this request acquired the lock may have changed or removed the runtime. + const current = d.mgr.get(id); + if (current) { + const state = current.bridge.turnState(); + if (state === "running" || state === "waiting") { + return reply.code(200).send({ id, alreadyActive: true }); + } + } + let rt: ReturnType | undefined; try { - if (existing) { - boundRuntimes.delete(id); - d.mgr.teardown(id); + if (current) { + if (boundRuntimes.get(id) === current) boundRuntimes.delete(id); + d.mgr.teardownIfCurrent(id, current); } 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); + if (rt) { + if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id); + d.mgr.teardownIfCurrent(id, rt); + } return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE }); } @@ -217,14 +254,19 @@ export function sessionRoutes( }); 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 { + return withSessionLifecycle(id, async () => { + // Invalidate the live generation before persistence can yield. Otherwise its deferred + // bootstrap may start Pi while Close is already in progress. + const current = d.mgr.get(id); boundRuntimes.delete(id); - d.mgr.teardown(id); - d.hub.clear(id); - } - return { closed: true }; + if (current) d.mgr.teardownIfCurrent(id, current); + try { + await d.tht.closeSession(id, d.getSettings().workspace); + } finally { + d.hub.clear(id); + } + return { closed: true }; + }); }); app.get("/sessions/:id/events", (req, reply) => { const id = (req.params as any).id; @@ -252,6 +294,11 @@ export function sessionRoutes( const off = d.hub.subscribe(id, send, { afterId, pending: rt?.bridge.pendingWidget() ?? null, + // clear()/forget() end every old transport so native EventSource reconnects with its + // Last-Event-ID instead of remaining attached to a subscriber callback that no longer exists. + close: () => { + if (!reply.raw.writableEnded) reply.raw.end(); + }, }); req.raw.on("close", off); }); @@ -273,11 +320,14 @@ 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(); + return withSessionLifecycle(id, async () => { + const current = d.mgr.get(id); + boundRuntimes.delete(id); + if (current) d.mgr.teardownIfCurrent(id, current); + 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 367b1e03..c805c69b 100644 --- a/backend/src/sse/sse-hub.ts +++ b/backend/src/sse/sse-hub.ts @@ -1,5 +1,11 @@ type Send = (event: string, data: object, id: number) => void; +interface Subscriber { + send: Send; + close: () => void; + closed: boolean; +} + /** * Per-session ring buffer of recent events. * @@ -17,6 +23,7 @@ interface BufferedEvent { interface SubscribeOptions { afterId?: number; pending?: object | null; + close?: () => void; } function descriptorId(value: unknown): string | null { @@ -32,7 +39,7 @@ function bufferedGateId(item: BufferedEvent): string | null { } export class SseHub { - private subs = new Map>(); + private subs = new Map>(); private buffers = new Map(); private lastIds = new Map(); @@ -42,7 +49,8 @@ export class SseHub { options: SubscribeOptions = {}, ): () => void { if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set()); - this.subs.get(sessionId)!.add(send); + const subscriber = { send, close: options.close ?? (() => undefined), closed: false }; + this.subs.get(sessionId)!.add(subscriber); const afterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0 ? options.afterId ?? 0 @@ -63,12 +71,14 @@ export class SseHub { send(item.event, item.data, item.id); } - return () => this.subs.get(sessionId)?.delete(send); + return () => this.subs.get(sessionId)?.delete(subscriber); } publish(sessionId: string, event: string, data: object): number { const item = this.buffer(sessionId, event, data); - for (const send of this.subs.get(sessionId) ?? []) send(event, data, item.id); + for (const subscriber of this.subs.get(sessionId) ?? []) { + subscriber.send(event, data, item.id); + } return item.id; } @@ -88,6 +98,13 @@ export class SseHub { /** Clear buffered/runtime bindings while retaining session event-id monotonicity. */ clear(sessionId: string): void { + // Snapshot because ending an HTTP response can synchronously/asynchronously unsubscribe it. + // Mark before invoking callbacks so even a re-entrant clear cannot close a response twice. + for (const subscriber of [...(this.subs.get(sessionId) ?? [])]) { + if (subscriber.closed) continue; + subscriber.closed = true; + try { subscriber.close(); } catch { /* disconnect every remaining subscriber */ } + } this.buffers.delete(sessionId); this.subs.delete(sessionId); } diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index 41a99085..f5aa1c9e 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -66,6 +66,27 @@ test("l'exit del VECCHIO child non elimina il nuovo runtime (exit identity-check mgr.teardown("respawn-id"); }); +test("identity-checked teardown cannot kill a replacement runtime", () => { + const cfg = loadConfig({ THT_HARNESS_DIR: "../harness" }); + const firstChild = recordingChild(); + const secondChild = recordingChild(); + firstChild.kill = vi.fn(); + secondChild.kill = vi.fn(); + const children = [firstChild, secondChild]; + const mgr = new PiProcessManager(cfg, { spawnFn: () => children.shift() as any }); + const first = mgr.createFor("replace-id", {}); + mgr.teardown("replace-id"); + const second = mgr.createFor("replace-id", {}); + + expect(mgr.teardownIfCurrent("replace-id", first)).toBe(false); + expect(mgr.get("replace-id")).toBe(second); + expect(secondChild.kill).not.toHaveBeenCalled(); + + expect(mgr.teardownIfCurrent("replace-id", second)).toBe(true); + expect(mgr.get("replace-id")).toBeUndefined(); + expect(secondChild.kill).toHaveBeenCalledOnce(); +}); + 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/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 79eac85e..f5ce8718 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -18,6 +18,16 @@ function mutApp(thtRunner: any) { }); } +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((onResolve, onReject) => { + resolve = onResolve; + reject = onReject; + }); + return { promise, resolve, reject }; +} + test("POST /sessions usa i settings (workspace/provider/model/thinking) e crea+avvia", async () => { const modelKey = path.join(os.tmpdir(), `thoth-model-key-${process.pid}`); writeFileSync(modelKey, "test-model-key", { mode: 0o600 }); @@ -52,11 +62,17 @@ test("POST /sessions configura Pi con il thinking globale selezionato", async () let configured: any; const bridge = { onClientEvent: () => {}, emitClientEvent: () => {} }; const runtime = { bridge } as any; + let current: any; const mgr = { - createFor: () => runtime, + get: () => current, + createFor: () => { current = runtime; return runtime; }, configure: async (_rt: any, options: any) => { configured = options; }, start: () => {}, - teardown: () => {}, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr, @@ -78,12 +94,17 @@ test("POST /sessions/:id/resume configura Pi con il thinking persistito", async let configured: any; const bridge = { onClientEvent: () => {}, emitClientEvent: () => {} }; const runtime = { bridge } as any; + let current: any; const mgr = { - get: () => undefined, - createFor: () => runtime, + get: () => current, + createFor: () => { current = runtime; return runtime; }, configure: async (_rt: any, options: any) => { configured = options; }, start: () => {}, - teardown: () => {}, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr, @@ -111,12 +132,17 @@ test("POST /sessions/:id/resume usa il thinking globale se manca nel manifest", let configured: any; const bridge = { onClientEvent: () => {}, emitClientEvent: () => {} }; const runtime = { bridge } as any; + let current: any; const mgr = { - get: () => undefined, - createFor: () => runtime, + get: () => current, + createFor: () => { current = runtime; return runtime; }, configure: async (_rt: any, options: any) => { configured = options; }, start: () => {}, - teardown: () => {}, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr, @@ -162,11 +188,17 @@ test.each(["idle", "failed"])( const order: string[] = []; const oldRuntime = { bridge: { turnState: () => state } } as any; const newRuntime = { bridge: { onClientEvent: () => {}, emitClientEvent: () => {} } } as any; + let current: any = oldRuntime; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { - get: () => oldRuntime, - teardown: (id: string) => order.push(`teardown:${id}`), - createFor: () => { order.push("create"); return newRuntime; }, + get: () => current, + teardownIfCurrent: (id: string, expected: any) => { + if (current !== expected) return false; + order.push(`teardown:${id}`); + current = undefined; + return true; + }, + createFor: () => { order.push("create"); current = newRuntime; return newRuntime; }, configure: async () => {}, start: () => order.push("start"), } as any, @@ -195,12 +227,14 @@ test("POST resume without a runtime clears stale SSE state before cold start", a const order: string[] = []; let createOptions: any; const newRuntime = { bridge: { onClientEvent: () => {}, emitClientEvent: () => {} } } as any; + let current: any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { - get: () => undefined, + get: () => current, createFor: (_id: string, options: any) => { createOptions = options; order.push("create"); + current = newRuntime; return newRuntime; }, configure: async () => {}, @@ -327,10 +361,15 @@ test("POST resume sanitizes runtime creation failure and preserves the old hub a } as any; hub.subscribe("s1", (_event: string, data: any) => delivered.push(data.text)); hub.publish("s1", "info", { text: "before" }); + let current: any = { bridge: { turnState: () => "failed" } }; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { - get: () => ({ bridge: { turnState: () => "failed" } }), - teardown: () => {}, + get: () => current, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, createFor: () => { createCalls += 1; throw new Error("CREATE_SENTINEL"); }, } as any, hub, @@ -424,6 +463,490 @@ test("concurrent cold Resume requests serialize and create one runtime", async ( }); }); +test.each([ + ["close", "resolve"], + ["close", "reject"], + ["delete", "resolve"], + ["delete", "reject"], +] as const)( + "a bootstrap that %s left behind cannot %s against a replacement runtime", + async (lifecycle, outcome) => { + const oldConfigure = deferred(); + const published: any[] = []; + const starts: string[] = []; + let failed = 0; + let configureCalls = 0; + let current: any; + const runtime = (name: string) => { + let listener: ((event: any) => void) | undefined; + return { + name, + bridge: { + turnState: () => "running", + onClientEvent: (next: (event: any) => void) => { listener = next; }, + emitClientEvent: (event: any) => listener?.(event), + }, + }; + }; + const oldRuntime = runtime("old"); + const replacement = runtime("replacement"); + const runtimes = [oldRuntime, replacement]; + const mgr = { + get: () => current, + createFor: () => { + current = runtimes.shift(); + return current; + }, + configure: async () => { + configureCalls += 1; + if (configureCalls === 1) await oldConfigure.promise; + }, + start: (_id: string, expected: any) => { + if (current !== expected) throw new Error("stale start"); + starts.push(expected.name); + }, + teardown: () => { current = undefined; }, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, + }; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: mgr as any, + hub: { + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + clear: () => {}, + forget: () => {}, + } as any, + thtRunner: { + sessionNew: async () => ({ id: "s1" }), + searchPack: async () => {}, + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + closeSession: 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: "old" } }); + if (lifecycle === "close") { + await app.inject({ method: "POST", url: "/sessions/s1/close" }); + await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + } else { + await app.inject({ method: "DELETE", url: "/sessions/s1" }); + await app.inject({ method: "POST", url: "/sessions", payload: { question: "replacement" } }); + } + await new Promise((resolve) => setImmediate(resolve)); + expect(current).toBe(replacement); + published.length = 0; + starts.length = 0; + failed = 0; + + if (outcome === "resolve") oldConfigure.resolve(); + else oldConfigure.reject(new Error("old bootstrap failed")); + await new Promise((resolve) => setImmediate(resolve)); + + expect(current).toBe(replacement); + expect(starts).toEqual([]); + expect(failed).toBe(0); + expect(published).toEqual([]); + }, +); + +test.each(["resolve", "reject"] as const)( + "a deleted session stays forgotten when its stale bootstrap later %s", + async (outcome) => { + const oldConfigure = deferred(); + const published: any[] = []; + let failed = 0; + let current: any; + let forgotten = false; + let listener: ((event: any) => void) | undefined; + const runtime = { + bridge: { + onClientEvent: (next: (event: any) => void) => { listener = next; }, + emitClientEvent: (event: any) => listener?.(event), + }, + }; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = runtime; return runtime; }, + configure: async () => oldConfigure.promise, + start: () => { throw new Error("deleted runtime must not start"); }, + teardown: () => { current = undefined; }, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, + } as any, + hub: { + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + forget: () => { forgotten = true; }, + } 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: "old" } }); + await app.inject({ method: "DELETE", url: "/sessions/s1" }); + published.length = 0; + + if (outcome === "resolve") oldConfigure.resolve(); + else oldConfigure.reject(new Error("old bootstrap failed")); + await new Promise((resolve) => setImmediate(resolve)); + + expect(current).toBeUndefined(); + expect(forgotten).toBe(true); + expect(failed).toBe(0); + expect(published).toEqual([]); + }, +); + +test.each(["close", "delete"] as const)( + "Resume followed by %s leaves no runtime or post-removal bootstrap events", + async (lifecycle) => { + const reopen = deferred(); + const configure = deferred(); + let markReopenStarted!: () => void; + const reopenStarted = new Promise((resolve) => { markReopenStarted = resolve; }); + let current: any; + let removalSettled = false; + const published: any[] = []; + const removalEvents: string[] = []; + const replacement = { + bridge: { + turnState: () => "running", + onClientEvent: () => {}, + emitClientEvent: () => {}, + }, + }; + const teardown = (expected?: any) => { + if (expected && current !== expected) return false; + current = undefined; + return true; + }; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = replacement; return replacement; }, + configure: async () => configure.promise, + start: () => { throw new Error("removed runtime must not start"); }, + teardown: () => { teardown(); }, + teardownIfCurrent: (_id: string, expected: any) => teardown(expected), + } as any, + hub: { + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + clear: () => { removalEvents.push("clear"); }, + forget: () => { removalEvents.push("forget"); }, + } as any, + thtRunner: { + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => { + markReopenStarted(); + await reopen.promise; + }, + closeSession: async () => { removalEvents.push("close"); }, + deleteSession: async () => { removalEvents.push("delete"); }, + failSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const resumeResponse = app.inject({ method: "POST", url: "/sessions/s1/resume" }); + await reopenStarted; + const removalResponse = app.inject({ + method: lifecycle === "close" ? "POST" : "DELETE", + url: `/sessions/s1/${lifecycle === "close" ? "close" : ""}`.replace(/\/$/, ""), + }).then((response) => { + removalSettled = true; + return response; + }); + await new Promise((resolve) => setImmediate(resolve)); + expect(removalSettled).toBe(false); + + reopen.resolve(); + expect((await resumeResponse).json()).toEqual({ id: "s1", alreadyActive: false }); + await removalResponse; + published.length = 0; + configure.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(current).toBeUndefined(); + expect(published).toEqual([]); + expect(removalEvents.at(-1)).toBe(lifecycle === "close" ? "clear" : "forget"); + }, +); + +test("Close followed by Resume installs a fresh runtime only after Close finishes", async () => { + const close = deferred(); + let markCloseStarted!: () => void; + const closeStarted = new Promise((resolve) => { markCloseStarted = resolve; }); + const oldRuntime = { bridge: { turnState: () => "running" } }; + const replacement = { + bridge: { + turnState: () => "running", + onClientEvent: () => {}, + emitClientEvent: () => {}, + }, + }; + let current: any = oldRuntime; + let resumeSettled = false; + const teardown = (expected?: any) => { + if (expected && current !== expected) return false; + current = undefined; + return true; + }; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = replacement; return replacement; }, + configure: async () => {}, + start: () => {}, + teardown: () => { teardown(); }, + teardownIfCurrent: (_id: string, expected: any) => teardown(expected), + } as any, + hub: { publish: () => 1, clear: () => {} } as any, + thtRunner: { + closeSession: async () => { markCloseStarted(); await close.promise; }, + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const closeResponse = app.inject({ method: "POST", url: "/sessions/s1/close" }); + await closeStarted; + const resumeResponse = app.inject({ method: "POST", url: "/sessions/s1/resume" }) + .then((response) => { resumeSettled = true; return response; }); + await new Promise((resolve) => setImmediate(resolve)); + expect(resumeSettled).toBe(false); + + close.resolve(); + await closeResponse; + const resumed = await resumeResponse; + await new Promise((resolve) => setImmediate(resolve)); + + expect(resumed.json()).toEqual({ id: "s1", alreadyActive: false }); + expect(current).toBe(replacement); +}); + +test("Close suppresses a bootstrap that settles while close persistence is pending", async () => { + const configure = deferred(); + const close = deferred(); + let markCloseStarted!: () => void; + const closeStarted = new Promise((resolve) => { markCloseStarted = resolve; }); + const published: any[] = []; + const starts: any[] = []; + let current: any; + let listener: ((event: any) => void) | undefined; + const runtime = { + bridge: { + onClientEvent: (next: (event: any) => void) => { listener = next; }, + emitClientEvent: (event: any) => listener?.(event), + }, + }; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = runtime; return runtime; }, + configure: async () => configure.promise, + start: (_id: string, expected: any) => { starts.push(expected); }, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, + } as any, + hub: { + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + clear: () => {}, + } as any, + thtRunner: { + sessionNew: async () => ({ id: "s1" }), + searchPack: async () => {}, + closeSession: async () => { markCloseStarted(); await close.promise; }, + failSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); + const closeResponse = app.inject({ method: "POST", url: "/sessions/s1/close" }); + await closeStarted; + published.length = 0; + configure.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(starts).toEqual([]); + expect(published).toEqual([]); + + close.resolve(); + await closeResponse; + expect(current).toBeUndefined(); +}); + +test("bootstrap failure persists once and keeps Resume serialized behind that persistence", async () => { + const oldConfigure = deferred(); + const failurePersistence = deferred(); + let markFailureStarted!: () => void; + const failureStarted = new Promise((resolve) => { markFailureStarted = resolve; }); + const oldRuntime = { + bridge: { + onClientEvent: undefined as ((next: (event: any) => void) => void) | undefined, + emitClientEvent: undefined as ((event: any) => void) | undefined, + }, + } as any; + const replacement = { + bridge: { + turnState: () => "running", + onClientEvent: () => {}, + emitClientEvent: () => {}, + }, + }; + let oldListener: ((event: any) => void) | undefined; + oldRuntime.bridge.onClientEvent = (next: (event: any) => void) => { oldListener = next; }; + oldRuntime.bridge.emitClientEvent = (event: any) => oldListener?.(event); + const runtimes = [oldRuntime, replacement]; + let current: any; + let configureCalls = 0; + let failCalls = 0; + let resumeSettled = false; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = runtimes.shift(); return current; }, + configure: async () => { + configureCalls += 1; + if (configureCalls === 1) await oldConfigure.promise; + }, + start: () => {}, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, + } as any, + hub: { publish: () => 1, clear: () => {} } as any, + thtRunner: { + sessionNew: async () => ({ id: "s1" }), + searchPack: async () => {}, + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + failSession: async () => { + failCalls += 1; + if (failCalls === 1) { + markFailureStarted(); + await failurePersistence.promise; + } + }, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); + oldConfigure.reject(new Error("configure failed")); + await failureStarted; + const resumeResponse = app.inject({ method: "POST", url: "/sessions/s1/resume" }) + .then((response) => { resumeSettled = true; return response; }); + await new Promise((resolve) => setImmediate(resolve)); + expect(resumeSettled).toBe(false); + + failurePersistence.resolve(); + expect((await resumeResponse).json()).toEqual({ id: "s1", alreadyActive: false }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(failCalls).toBe(1); + expect(current).toBe(replacement); +}); + +test("Delete followed by Resume cannot resurrect the deleted session", async () => { + const deletion = deferred(); + let markDeleteStarted!: () => void; + const deleteStarted = new Promise((resolve) => { markDeleteStarted = resolve; }); + const oldRuntime = { bridge: { turnState: () => "running" } }; + let current: any = oldRuntime; + let deleted = false; + let resumeSettled = false; + let createCalls = 0; + const teardown = (expected?: any) => { + if (expected && current !== expected) return false; + current = undefined; + return true; + }; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { + createCalls += 1; + current = { bridge: { onClientEvent: () => {}, emitClientEvent: () => {} } }; + return current; + }, + configure: async () => {}, + start: () => {}, + teardown: () => { teardown(); }, + teardownIfCurrent: (_id: string, expected: any) => teardown(expected), + } as any, + hub: { publish: () => 1, clear: () => {}, forget: () => {} } as any, + thtRunner: { + deleteSession: async () => { + markDeleteStarted(); + await deletion.promise; + deleted = true; + }, + sessionShow: async () => { + if (deleted) throw new Error("session deleted"); + return { status: "open", archived: false }; + }, + reopenSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const deleteResponse = app.inject({ method: "DELETE", url: "/sessions/s1" }); + await deleteStarted; + const resumeResponse = app.inject({ method: "POST", url: "/sessions/s1/resume" }) + .then((response) => { resumeSettled = true; return response; }); + await new Promise((resolve) => setImmediate(resolve)); + expect(resumeSettled).toBe(false); + + deletion.resolve(); + await deleteResponse; + const resumed = await resumeResponse; + + expect(resumed.statusCode).toBe(500); + expect(current).toBeUndefined(); + expect(createCalls).toBe(0); +}); + test("a replaced runtime cannot publish or fail the newly resumed session", async () => { const controlledBridge = (initialState: string) => { let state = initialState; @@ -448,7 +971,11 @@ test("a replaced runtime cannot publish or fail the newly resumed session", asyn createFor: () => { current = runtimes.shift(); return current; }, configure: async () => {}, start: () => {}, - teardown: () => { current = undefined; }, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any, hub: { clear: () => {}, @@ -505,7 +1032,11 @@ test("a deleted runtime cannot repopulate or fail the forgotten session", async createFor: () => { current = runtime; return runtime; }, configure: async () => {}, start: () => {}, - teardown: () => { current = undefined; }, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any, hub: { publish: (_id: string, event: string, data: object) => { @@ -604,7 +1135,12 @@ test("POST resume tears down a created runtime when bridge binding fails", async const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { get: () => current, - teardown: () => { teardownCalls += 1; current = undefined; }, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + teardownCalls += 1; + current = undefined; + return true; + }, createFor: () => { current = { bridge: { @@ -688,9 +1224,17 @@ test("DELETE /sessions/:id calls deleteSession", async () => { test("DELETE /sessions/:id tears down the runtime before deleting on disk", async () => { const order: string[] = []; + const runtime = {} as any; 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, + mgr: { + get: () => runtime, + teardownIfCurrent: (id: string, expected: any) => { + expect(expected).toBe(runtime); + order.push(`teardown:${id}`); + return true; + }, + } as any, hub: { forget: (id: string) => { order.push(`forget:${id}`); } } as any, getSettings: () => ({ workspace: "w" }) as any, spawnFn: () => nodeSpawn("node", [FAKE, SCRIPT]) as any, @@ -704,7 +1248,7 @@ test("POST /sessions/:id/close clears transient hub state without forgetting its const calls: string[] = []; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { closeSession: async () => {} } as any, - mgr: { teardown: () => {} } as any, + mgr: { get: () => undefined } as any, hub: { clear: (id: string) => { calls.push(`clear:${id}`); }, forget: (id: string) => { calls.push(`forget:${id}`); }, @@ -833,14 +1377,20 @@ test("POST /sessions returns after bridge attachment but starts only after retri emitClientEvent: () => {}, }; const runtime = { bridge } as any; + let current: any; const mgr = { - createFor: () => runtime, + get: () => current, + createFor: () => { current = runtime; return runtime; }, configure: async () => {}, start: () => { expect(bridgeAttached).toBe(true); started = true; }, - teardown: () => {}, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr, @@ -871,13 +1421,18 @@ test("POST /sessions bootstrap failure emits only a fixed recovery message", asy const runtime = { bridge } as any; const rawFailure = "connect https://secret.invalid/bootstrap?token=DO_NOT_LEAK using /srv/private/model-key"; + let current: any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { - get: () => undefined, - createFor: () => runtime, + get: () => current, + createFor: () => { current = runtime; return runtime; }, configure: async () => { throw new Error(rawFailure); }, start: () => {}, - teardown: () => {}, + teardownIfCurrent: (_id: string, expected: any) => { + if (current !== expected) return false; + current = undefined; + return true; + }, } as any, hub: { publish: (_id: string, event: string, data: any) => published.push({ event, data }), diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts index 24627c9e..a61ba2d6 100644 --- a/backend/test/sse-hub.test.ts +++ b/backend/test/sse-hub.test.ts @@ -21,14 +21,34 @@ test("publish assigns monotonic ids and a subscriber replays only ids newer than ]); }); -test("clear drops stale subscribers and buffers but preserves the per-session id sequence", () => { +test("clear closes every stale subscriber once and replays post-clear events from the cursor", () => { const hub = new SseHub(); - const stale: any[] = []; - hub.subscribe("s1", (event, data, id) => stale.push({ event, data, id })); + const firstStale: any[] = []; + const secondStale: any[] = []; + const closeCalls = [0, 0]; + let offFirst: () => void = () => undefined; + let offSecond: () => void = () => undefined; + offFirst = hub.subscribe( + "s1", + (event, data, id) => firstStale.push({ event, data, id }), + { close: () => { closeCalls[0] += 1; offFirst(); } }, + ); + offSecond = hub.subscribe( + "s1", + (event, data, id) => secondStale.push({ event, data, id }), + { close: () => { closeCalls[1] += 1; offSecond(); } }, + ); expect(hub.publish("s1", "info", { text: "before" })).toBe(1); hub.clear("s1"); + offFirst(); + offSecond(); + expect(closeCalls).toEqual([1, 1]); expect(hub.publish("s1", "info", { text: "after" })).toBe(2); + expect(hub.publish("s1", "ui_request", { + type: "ui_request", + ui_request: { id: "gate-1", widget: "select" }, + })).toBe(3); const resumed: any[] = []; hub.subscribe( "s1", @@ -36,22 +56,47 @@ test("clear drops stale subscribers and buffers but preserves the per-session id { afterId: 1 }, ); - expect(stale).toEqual([{ event: "info", data: { text: "before" }, id: 1 }]); - expect(resumed).toEqual([{ event: "info", data: { text: "after" }, id: 2 }]); + const before = [{ event: "info", data: { text: "before" }, id: 1 }]; + expect(firstStale).toEqual(before); + expect(secondStale).toEqual(before); + expect(resumed).toEqual([ + { event: "info", data: { text: "after" }, id: 2 }, + { + event: "ui_request", + data: { type: "ui_request", ui_request: { id: "gate-1", widget: "select" } }, + id: 3, + }, + ]); }); -test("forget drops subscribers, buffers, and the per-session id sequence", () => { +test("forget closes every subscriber once and resets the per-session id sequence", () => { const hub = new SseHub(); - const stale: any[] = []; - hub.subscribe("s1", (event, data, id) => stale.push({ event, data, id })); + const firstStale: any[] = []; + const secondStale: any[] = []; + const closeCalls = [0, 0]; + let offFirst: () => void = () => undefined; + let offSecond: () => void = () => undefined; + offFirst = hub.subscribe( + "s1", + (event, data, id) => firstStale.push({ event, data, id }), + { close: () => { closeCalls[0] += 1; offFirst(); } }, + ); + offSecond = hub.subscribe( + "s1", + (event, data, id) => secondStale.push({ event, data, id }), + { close: () => { closeCalls[1] += 1; offSecond(); } }, + ); expect(hub.publish("s1", "info", { text: "before" })).toBe(1); hub.forget("s1"); + expect(closeCalls).toEqual([1, 1]); 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 }]); + const before = [{ event: "info", data: { text: "before" }, id: 1 }]; + expect(firstStale).toEqual(before); + expect(secondStale).toEqual(before); expect(fresh).toEqual([{ event: "info", data: { text: "after" }, id: 1 }]); }); diff --git a/backend/test/sse-route.test.ts b/backend/test/sse-route.test.ts index 2c7f865c..a1ecdb00 100644 --- a/backend/test/sse-route.test.ts +++ b/backend/test/sse-route.test.ts @@ -55,6 +55,17 @@ async function captureReplay(options: { query?: string; lastEventId?: string }): } } +async function expectStreamEnd(reader: ReadableStreamDefaultReader): Promise { + const result = await Promise.race([ + reader.read(), + new Promise<{ timeout: true }>((resolve) => + setTimeout(() => resolve({ timeout: true }), 1_000), + ), + ]); + expect(result).not.toEqual({ timeout: true }); + expect("done" in result && result.done).toBe(true); +} + test("SSE emits ids and honors the native Last-Event-ID replay cursor", async () => { const body = await captureReplay({ lastEventId: "1" }); @@ -80,3 +91,54 @@ test("SSE uses the newer valid cursor when header and query are both present", a expect(body).not.toContain('"two"'); expect(body).toContain("id: 3\nevent: info\n"); }); + +test("clear ends every live SSE response and a cursor reconnect replays post-clear events", async () => { + const hub = new SseHub(); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), { + hub, + thtRunner: {} as any, + }); + await app.listen({ port: 0, host: "127.0.0.1" }); + const port = (app.server.address() as { port: number }).port; + const controllers = [new AbortController(), new AbortController(), new AbortController()]; + + try { + const responses = await Promise.all(controllers.slice(0, 2).map((controller) => fetch( + `http://127.0.0.1:${port}/sessions/s1/events`, + { signal: controller.signal }, + ))); + const readers = responses.map((response) => response.body!.getReader()); + + expect(hub.publish("s1", "info", { type: "info", text: "before" })).toBe(1); + const initial = await Promise.all(readers.map((reader) => + readUntil(reader, (text) => text.includes('"before"')))); + expect(initial.every((body) => body.includes("id: 1\nevent: info\n"))).toBe(true); + + hub.clear("s1"); + await Promise.all(readers.map(expectStreamEnd)); + + expect(hub.publish("s1", "info", { type: "info", text: "after" })).toBe(2); + expect(hub.publish("s1", "ui_request", { + type: "ui_request", + ui_request: { id: "gate-1", widget: "select" }, + })).toBe(3); + + const reconnected = await fetch( + `http://127.0.0.1:${port}/sessions/s1/events`, + { + headers: { "Last-Event-ID": "1" }, + signal: controllers[2].signal, + }, + ); + const reconnectReader = reconnected.body!.getReader(); + const replay = await readUntil(reconnectReader, (text) => text.includes('"gate-1"')); + expect(replay).toContain("id: 2\nevent: info\n"); + expect(replay).toContain('data: {"type":"info","text":"after"}'); + expect(replay).toContain("id: 3\nevent: ui_request\n"); + expect(replay).toContain('"ui_request":{"id":"gate-1","widget":"select"}'); + await reconnectReader.cancel(); + } finally { + for (const controller of controllers) controller.abort(); + await app.close(); + } +}); diff --git a/frontend/src/shell/AppShell.session-mgmt.test.tsx b/frontend/src/shell/AppShell.session-mgmt.test.tsx index 8894996f..3c9c2d5e 100644 --- a/frontend/src/shell/AppShell.session-mgmt.test.tsx +++ b/frontend/src/shell/AppShell.session-mgmt.test.tsx @@ -20,6 +20,12 @@ const LIST = [ const resumeResult = (id: string, alreadyActive = false) => HttpResponse.json({ id, alreadyActive }); +function deferred() { + let resolve!: () => void; + const promise = new Promise((onResolve) => { resolve = onResolve; }); + return { promise, resolve }; +} + beforeEach(() => { FakeEventSource.instances = []; (globalThis as any).EventSource = FakeEventSource; @@ -134,58 +140,104 @@ 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 () => { +test("concurrent same-id Resume invocations share one cold request and replacement source", 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; }); + const coldGate = deferred(); + const coldStarted = deferred(); 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); + if (resumeCalls === 1) return resumeResult("s1", false); + if (resumeCalls === 2) coldStarted.resolve(); + await coldGate.promise; + return resumeResult("s1", false); }), 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(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + const oldSource = FakeEventSource.instances[0]; + act(() => oldSource.emitNamed("info", { type: "info", text: "Old generation" }, "900")); + await userEvent.click(screen.getByText("Attiva uno")); + const resumeButton = await screen.findByRole("button", { name: /resume/i }); await userEvent.click(resumeButton); - await coldStarted; + await coldStarted.promise; await userEvent.click(resumeButton); - await alreadyActiveStarted; + await new Promise((resolve) => setImmediate(resolve)); expect(resumeCalls).toBe(2); - releaseCold(); - await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); - const stream = FakeEventSource.instances[0]; - act(() => stream.emitNamed( + coldGate.resolve(); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(2)); + const replacement = FakeEventSource.instances[1]; + expect(oldSource.closed).toBe(true); + expect(replacement.url).toBe("http://localhost:8787/sessions/s1/events"); + expect(useSessionStore.getState().activityLog).toEqual([ + { kind: "lifecycle", phase: null, text: "Resuming session" }, + ]); + act(() => replacement.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("a committed Resume releases same-id single-flight before its manifest settles", async () => { + const other = { + ...LIST[0], id: "s3", question: "Attiva tre", group: null, + created_at: "2026-01-03T00:00:00Z", + }; + let s1ResumeCalls = 0; + let s1ManifestCalls = 0; + const firstS1ManifestGate = deferred(); + const firstS1ManifestStarted = deferred(); + server.use( + http.get("http://localhost:8787/sessions", () => HttpResponse.json([LIST[0], other])), + http.post("http://localhost:8787/sessions/:id/resume", ({ params }) => { + const id = params.id as string; + if (id === "s1") s1ResumeCalls += 1; + return resumeResult(id); + }), + http.get("http://localhost:8787/sessions/:id", async ({ params }) => { + if (params.id === "s1") { + s1ManifestCalls += 1; + if (s1ManifestCalls === 1) { + firstS1ManifestStarted.resolve(); + await firstS1ManifestGate.promise; + return HttpResponse.json({ id: "s1", status: "open", phase: 7 }); + } + return HttpResponse.json({ id: "s1", status: "open", phase: 4 }); + } + return HttpResponse.json({ id: "s3", status: "open", phase: 3 }); + }), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await firstS1ManifestStarted.promise; + expect(s1ResumeCalls).toBe(1); + + await userEvent.click(screen.getByText("Attiva tre")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(useSessionStore.getState().currentPhase).toBe("F3")); + + await userEvent.click(screen.getByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(s1ResumeCalls).toBe(2)); + await waitFor(() => expect(useSessionStore.getState().currentPhase).toBe("F4")); + expect(FakeEventSource.instances.at(-1)?.url).toContain("/sessions/s1/events"); + + firstS1ManifestGate.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(useSessionStore.getState().currentPhase).toBe("F4"); + expect(FakeEventSource.instances.at(-1)?.url).toContain("/sessions/s1/events"); }); test("cold same-session Resume keeps the old stream until success then receives post-clear events once", async () => { @@ -321,6 +373,122 @@ test("resuming a different already-active session binds it only after success", })); }); +test("competing Resume requests for different ids commit only the latest intent", async () => { + const other = { + ...LIST[0], id: "s3", question: "Attiva tre", group: null, + created_at: "2026-01-03T00:00:00Z", + }; + const s1Gate = deferred(); + const s3Gate = deferred(); + const s1Started = deferred(); + const s3Started = deferred(); + server.use( + http.get("http://localhost:8787/sessions", () => HttpResponse.json([LIST[0], other])), + http.post("http://localhost:8787/sessions/:id/resume", async ({ params }) => { + const id = params.id as string; + if (id === "s1") { + s1Started.resolve(); + await s1Gate.promise; + } else { + s3Started.resolve(); + await s3Gate.promise; + } + return resumeResult(id); + }), + http.get("http://localhost:8787/sessions/:id", ({ params }) => + HttpResponse.json({ id: params.id, status: "open", phase: params.id === "s3" ? 3 : 1 })), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await s1Started.promise; + await userEvent.click(screen.getByText("Attiva tre")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await s3Started.promise; + + s3Gate.resolve(); + await waitFor(() => expect(useSessionStore.getState().currentPhase).toBe("F3")); + expect(FakeEventSource.instances).toHaveLength(1); + const latestSource = FakeEventSource.instances[0]; + expect(latestSource.url).toContain("/sessions/s3/events"); + act(() => latestSource.emitNamed("info", { type: "info", text: "Latest target" }, "1")); + + s1Gate.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(FakeEventSource.instances).toHaveLength(1); + expect(latestSource.closed).toBe(false); + expect(useSessionStore.getState().currentPhase).toBe("F3"); + expect(useSessionStore.getState().activityLog).toContainEqual(expect.objectContaining({ + text: "Latest target", + })); +}); + +test("a stale Resume manifest cannot repaint the latest session phase", async () => { + const other = { + ...LIST[0], id: "s3", question: "Attiva tre", group: null, + created_at: "2026-01-03T00:00:00Z", + }; + const s1ManifestGate = deferred(); + const s1ManifestStarted = deferred(); + server.use( + http.get("http://localhost:8787/sessions", () => HttpResponse.json([LIST[0], other])), + http.post("http://localhost:8787/sessions/:id/resume", ({ params }) => + resumeResult(params.id as string)), + http.get("http://localhost:8787/sessions/:id", async ({ params }) => { + if (params.id === "s1") { + s1ManifestStarted.resolve(); + await s1ManifestGate.promise; + return HttpResponse.json({ id: "s1", status: "open", phase: 7 }); + } + return HttpResponse.json({ id: "s3", status: "open", phase: 3 }); + }), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await s1ManifestStarted.promise; + await userEvent.click(screen.getByText("Attiva tre")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(useSessionStore.getState().currentPhase).toBe("F3")); + expect(FakeEventSource.instances.at(-1)?.url).toContain("/sessions/s3/events"); + + s1ManifestGate.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(useSessionStore.getState().currentPhase).toBe("F3"); + expect(FakeEventSource.instances.at(-1)?.url).toContain("/sessions/s3/events"); +}); + +test("starting a new question invalidates a pending Resume intent", async () => { + const resumeGate = deferred(); + const resumeStarted = deferred(); + server.use( + http.post("http://localhost:8787/sessions/:id/resume", async () => { + resumeStarted.resolve(); + await resumeGate.promise; + return resumeResult("s1"); + }), + http.post("http://localhost:8787/runtime/prewarm", () => + HttpResponse.json({ status: "warming" }, { status: 202 })), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await resumeStarted.promise; + await userEvent.click(screen.getByRole("button", { name: /^new session$/i })); + expect(await screen.findByText(/type your question/i)).toBeInTheDocument(); + + resumeGate.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(FakeEventSource.instances).toHaveLength(0); + expect(screen.getByText(/type your question/i)).toBeInTheDocument(); +}); + test("a failed Resume with no active session preserves the panel and existing activity", async () => { server.use(http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 409 }))); useSessionStore.setState({ diff --git a/frontend/src/shell/AppShell.tsx b/frontend/src/shell/AppShell.tsx index ee454dfc..d066f783 100644 --- a/frontend/src/shell/AppShell.tsx +++ b/frontend/src/shell/AppShell.tsx @@ -33,6 +33,12 @@ import { useEffect, useMemo, useRef, useState } from "react"; export function AppShell() { const [activeSessionId, setActiveSessionId] = useState(null); const activeSessionIdRef = useRef(null); + const resumeInvocationRef = useRef(0); + const latestResumeIntentRef = useRef<{ token: number; id: string } | null>(null); + const resumeInFlightRef = useRef(new Map; + }>()); const [streamCursorResetEpoch, setStreamCursorResetEpoch] = useState(0); const [creatingSession, setCreatingSession] = useState(false); const [awaitingQuestion, setAwaitingQuestion] = useState(false); @@ -72,6 +78,11 @@ export function AppShell() { setActiveSessionId(id); } + function invalidateResumeIntent() { + resumeInvocationRef.current += 1; + latestResumeIntentRef.current = null; + } + // A background refresh can remove a session (for example from another browser). // Keep the local selection aligned with the authoritative list. useEffect(() => { @@ -107,8 +118,37 @@ export function AppShell() { }); } async function doResume(id: string) { + const token = ++resumeInvocationRef.current; + latestResumeIntentRef.current = { token, id }; + const inFlight = resumeInFlightRef.current.get(id); + if (inFlight) { + // Repeated intent for the same target shares one backend lifecycle operation and one + // commit path. Updating its token still lets s1→s2→s1 make the final s1 intent authoritative. + inFlight.latestToken = token; + return inFlight.promise; + } + + const operation = { + latestToken: token, + promise: Promise.resolve(), + }; + operation.promise = runResume(id, operation).finally(() => { + if (resumeInFlightRef.current.get(id) === operation) { + resumeInFlightRef.current.delete(id); + } + }); + resumeInFlightRef.current.set(id, operation); + return operation.promise; + } + + async function runResume( + id: string, + operation: { latestToken: number; promise: Promise }, + ) { try { const result = await resumeSession(id); + const latest = latestResumeIntentRef.current; + if (latest?.token !== operation.latestToken || latest.id !== id) return; const reconnectSameSession = activeSessionIdRef.current === id; setPanelSession(null); setAwaitingQuestion(false); @@ -129,10 +169,22 @@ export function AppShell() { setStreamCursorResetEpoch((value) => value + 1); } + // Single-flight covers only the backend Resume and its local binding commit. A slow + // manifest read must not prevent a later explicit intent from starting a new Resume. + if (resumeInFlightRef.current.get(id) === operation) { + resumeInFlightRef.current.delete(id); + } + // Paint the persisted re-entry phase while the replacement stream starts replaying. // The manifest's `phase` is the 1-based current phase (1..8). try { const m = (await getSession(id)) as { phase?: number }; + const latestAfterManifest = latestResumeIntentRef.current; + if ( + latestAfterManifest?.token !== operation.latestToken + || latestAfterManifest.id !== id + || activeSessionIdRef.current !== id + ) return; if (typeof m.phase === "number" && m.phase >= 1 && m.phase <= 8) { setPhase(`F${m.phase}`); } @@ -140,7 +192,12 @@ export function AppShell() { /* non-fatal: the first gate will set the phase */ } } catch { - toast.error("Failed to resume session."); + if ( + latestResumeIntentRef.current?.token === operation.latestToken + && latestResumeIntentRef.current.id === id + ) { + toast.error("Failed to resume session."); + } } } async function move(s: SessionSummary, group: string) { @@ -183,6 +240,9 @@ export function AppShell() { } async function deleteSessions(targets: SessionSummary[]) { + if (targets.some((session) => session.id === activeSessionIdRef.current)) { + invalidateResumeIntent(); + } try { const results = await Promise.allSettled(targets.map((session) => deleteSession(session.id))); const deletedIds = new Set( @@ -246,6 +306,7 @@ export function AppShell() { if (ev === "session_exit") { // Never let a streamed event terminate the managed Pi child. Only the // explicit “Stop & save” action is allowed to call /close. + invalidateResumeIntent(); resetSession(); selectActiveSession(null); setAwaitingQuestion(false); @@ -258,6 +319,7 @@ export function AppShell() { }, [lastSystemEvent]); function startNewSession() { + invalidateResumeIntent(); resetSession(); setAwaitingQuestion(true); setCreatingSession(false); @@ -290,6 +352,7 @@ export function AppShell() { async function stopSession() { if (!activeSessionId) return; + invalidateResumeIntent(); try { await closeSession(activeSessionId); } finally {