diff --git a/.superpowers/sdd/predeploy-fix-report.md b/.superpowers/sdd/predeploy-fix-report.md new file mode 100644 index 00000000..ee9c497d --- /dev/null +++ b/.superpowers/sdd/predeploy-fix-report.md @@ -0,0 +1,359 @@ +# Pre-deployment Fix Wave Report + +Date: 2026-07-14 +Worktree: `/home/chirone/ThothII/.worktrees/activity-log-cte-layout` +Base: `e5366d14a6da8fb331d94be60b8929cefb1fe3e0` + +## Outcome + +All three reviewed findings are implemented in one coherent backend/frontend wave: + +1. Resume leaves the prior selection, Zustand state, document panel, and EventSource untouched + until `POST /resume` succeeds. Cold Resume changes state and reconnects only after backend + clear/rebind; already-active same-session Resume preserves the existing binding; failure is a + no-op apart from the fixed toast. +2. SSE uses monotonically increasing per-session ids, cursor-filtered replay, native and manual + reconnect cursors, id continuity across `hub.clear`, and descriptor-id pending-gate + idempotence at both backend and frontend layers. +3. Generic Pi system events and readiness errors are projected through explicit public + allowlists. Sentinel URLs, paths, tokens, stderr, commands, and extra fields do not reach HTTP + or SSE. + +No harness, workflow, persistence, model, CTE viewer, CTE card, or shared Card file changed. + +## Interfaces + +- Frontend `resumeSession(id)` now returns + `Promise<{ id: string; alreadyActive: boolean }>` via `ResumeSessionResult`. +- Backend successful Resume always returns the same shape: + - running/waiting runtime: `{ id, alreadyActive: true }` + - validated cold runtime: `{ id, alreadyActive: false }` +- `SseHub.publish(sessionId, event, data): number` returns the assigned SSE id. +- `SseHub.subscribe(sessionId, send, { afterId, pending })` calls + `send(event, data, id)` for replay/live frames with `id > afterId`. +- `GET /sessions/:id/events` accepts native `Last-Event-ID` and manual + `?lastEventId=`; when both are valid it uses the greater cursor. +- Every emitted SSE frame is `id: \nevent: \ndata: \n\n`. +- Public readiness failure is exactly: + `Session services are not ready. Check configuration and connectivity, then try again.` +- Generic Pi system events are exactly `{ type: "system_event", event }`, and `event` must be a + non-empty string. + +## Files + +Backend production: + +- `backend/src/bridge/session-bridge.ts` +- `backend/src/routes/sessions.ts` +- `backend/src/sse/sse-hub.ts` + +Backend tests: + +- `backend/test/routes-sessions.test.ts` +- `backend/test/session-bridge.test.ts` +- `backend/test/sse-hub.test.ts` +- `backend/test/sse-route.test.ts` (new) + +Frontend production/support: + +- `frontend/src/api/sessions.ts` +- `frontend/src/api/types.ts` +- `frontend/src/shell/AppShell.tsx` +- `frontend/src/store/sessionStore.ts` +- `frontend/src/stream/useSessionStream.ts` +- `frontend/src/test/fakeEventSource.ts` + +Frontend tests: + +- `frontend/src/api/sessions.test.ts` +- `frontend/src/shell/AppShell.session-mgmt.test.tsx` +- `frontend/src/store/sessionStore.test.ts` +- `frontend/src/stream/useSessionStream.test.tsx` + +## TDD RED/GREEN evidence + +### 1. Backend Resume result and client-boundary allowlists + +RED command: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts test/session-bridge.test.ts +``` + +RED output (exit 1): + +```text +Test Files 2 failed (2) +Tests 6 failed | 35 passed (41) + +expected { id: 's1' } to deeply equal { id: 's1', alreadyActive: false } +expected raw readiness URL/token/path to equal the fixed public message +expected three raw generic system events to equal [{ type: 'system_event', event: 'session_exit' }] +``` + +GREEN command: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts test/session-bridge.test.ts +``` + +GREEN output (exit 0): + +```text +✓ test/session-bridge.test.ts (14 tests) +✓ test/routes-sessions.test.ts (27 tests) +Test Files 2 passed (2) +Tests 41 passed (41) +``` + +### 2. Backend exact-once SseHub and route framing + +RED command: + +```text +cd backend && npx vitest run test/sse-hub.test.ts test/sse-route.test.ts +``` + +RED output (exit 1): + +```text +Test Files 2 failed (2) +Tests 7 failed (7) + +expected [undefined, undefined, undefined] to deeply equal [1, 2, 3] +expected unconditional replay not to contain "one" / "two" +expected one buffered pending gate, received replay plus a second pending emission +``` + +GREEN command: + +```text +cd backend && npx vitest run test/sse-hub.test.ts test/sse-route.test.ts +``` + +GREEN output (exit 0): + +```text +✓ test/sse-hub.test.ts (4 tests) +✓ test/sse-route.test.ts (3 tests) +Test Files 2 passed (2) +Tests 7 passed (7) +``` + +### 3. Frontend cursor tracking and gate idempotence + +RED command: + +```text +cd frontend && npx vitest run src/stream/useSessionStream.test.tsx src/store/sessionStore.test.ts +``` + +RED output (exit 1): + +```text +Test Files 2 failed (2) +Tests 2 failed | 28 passed (30) + +expected /sessions/s1/events to be /sessions/s1/events?lastEventId=7 +expected duplicate gate pendingWidget to remain null, received gate-1 +``` + +GREEN command: + +```text +cd frontend && npx vitest run src/stream/useSessionStream.test.tsx src/store/sessionStore.test.ts +``` + +GREEN output (exit 0): + +```text +✓ src/store/sessionStore.test.ts (23 tests) +✓ src/stream/useSessionStream.test.tsx (7 tests) +Test Files 2 passed (2) +Tests 30 passed (30) +``` + +### 4. Frontend typed Resume and AppShell ordering/preservation + +Typed API RED command: + +```text +cd frontend && npx tsc -b +``` + +Typed API RED output (exit 1): + +```text +src/api/sessions.test.ts(43,9): error TS2322: Type 'void' is not assignable to type +'{ id: string; alreadyActive: boolean; }'. +``` + +Lifecycle RED command: + +```text +cd frontend && npx vitest run src/api/sessions.test.ts src/shell/AppShell.session-mgmt.test.tsx +``` + +Lifecycle RED output (exit 1): + +```text +✓ src/api/sessions.test.ts (9 tests) +❯ src/shell/AppShell.session-mgmt.test.tsx (15 tests | 4 failed) +Test Files 1 failed | 1 passed (2) +Tests 4 failed | 20 passed (24) + +already-active same-session Resume created two EventSources instead of one +deferred cold Resume closed the document panel before POST completion +failed same-session Resume closed the prior EventSource +failed Resume with no active session opened an EventSource +``` + +GREEN commands: + +```text +cd frontend && npx vitest run src/api/sessions.test.ts src/shell/AppShell.session-mgmt.test.tsx +cd frontend && npx tsc -b +``` + +GREEN output (exit 0): + +```text +✓ src/api/sessions.test.ts (9 tests) +✓ src/shell/AppShell.session-mgmt.test.tsx (15 tests) +Test Files 2 passed (2) +Tests 24 passed (24) +TypeScript: no output, exit 0 +``` + +The AppShell cold-reconnect test additionally proves that the old source accepts an event while +Resume is pending, the replacement URL carries `lastEventId=8`, the replacement receives one +post-resume transcript/activity row, and two deliveries of the same descriptor id yield one gate. + +## Affected verification + +Backend command: + +```text +cd backend && npx vitest run test/routes-sessions.test.ts test/session-bridge.test.ts \ + test/sse-hub.test.ts test/sse-route.test.ts test/health.test.ts test/e2e-f1.test.ts +``` + +Output (exit 0): + +```text +Test Files 6 passed (6) +Tests 52 passed (52) +``` + +Backend typecheck: + +```text +cd backend && npx tsc --noEmit -p . +``` + +Output: no output, exit 0. + +Frontend command: + +```text +cd frontend && npx vitest run src/api/sessions.test.ts src/store/sessionStore.test.ts \ + src/stream/useSessionStream.test.tsx src/shell/AppShell.session-mgmt.test.tsx \ + src/shell/CentralStatus.test.tsx src/shell/ModelActivityPanel.test.tsx \ + src/shell/f1-loop.test.tsx src/shell/AppShell.new-session.test.tsx +``` + +Output (exit 0): + +```text +Test Files 8 passed (8) +Tests 74 passed (74) +``` + +Frontend typecheck: + +```text +cd frontend && npx tsc -b +``` + +Output: no output, exit 0. + +## Full verification + +Backend full suite: + +```text +cd backend && npx vitest run +``` + +```text +Test Files 22 passed (22) +Tests 177 passed (177) +``` + +Frontend full suite: + +```text +cd frontend && npx vitest run +``` + +```text +Test Files 43 passed (43) +Tests 271 passed (271) +``` + +Backend production build: + +```text +cd backend && npm run build +> tsc -p tsconfig.json +exit 0 +``` + +Frontend production build: + +```text +cd frontend && npm run build +> tsc -b && vite build +✓ 4835 modules transformed. +✓ built in 8.25s +exit 0 +``` + +Final whitespace verification: + +```text +git diff --check +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. +- 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. +- 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. +- Scope: `git diff` contains no CTE/Card/harness/workflow/persistence/model changes. Four pre-existing + modified `.superpowers/sdd/{progress,task-2-report,task-3-report,task-4-report}.md` files are user + work and are excluded from staging. + +## Remaining concerns + +- 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. +- 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 + 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. diff --git a/backend/src/bridge/session-bridge.ts b/backend/src/bridge/session-bridge.ts index 54a58778..e277916c 100644 --- a/backend/src/bridge/session-bridge.ts +++ b/backend/src/bridge/session-bridge.ts @@ -13,7 +13,7 @@ export type ClientEvent = | { type: "activity_delta"; text: string } | { type: "activity_event"; activity: ToolActivity } | { type: "info"; [k: string]: any } - | { type: "system_event"; [k: string]: any }; + | { type: "system_event"; event: string }; export type TurnState = "idle" | "running" | "waiting" | "failed"; @@ -61,7 +61,9 @@ export class SessionBridge { this.emitToolActivity(m, "end"); // Tool updates and raw payloads remain intentionally dropped. } else if (m.type === "system_event") { - this.fan(m as ClientEvent); + if (typeof m.event === "string" && m.event.trim() !== "") { + this.fan({ type: "system_event", event: m.event }); + } } else if (m.type === "agent_end") { if (this.state !== "failed" && !this.pending) this.state = "idle"; this.fan({ type: "system_event", event: "agent_end" }); diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 10d370f2..3835c761 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -8,6 +8,18 @@ import type { ReadinessManager } from "../runtime/readiness-manager.js"; 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."; + +function eventCursor(...values: unknown[]): number { + let cursor = 0; + for (const value of values.flatMap((item) => Array.isArray(item) ? item : [item])) { + if (typeof value !== "string" || !/^\d+$/.test(value)) continue; + const parsed = Number(value); + if (Number.isSafeInteger(parsed)) cursor = Math.max(cursor, parsed); + } + return cursor; +} export function sessionRoutes( app: FastifyInstance, @@ -57,7 +69,7 @@ export function sessionRoutes( const b = req.body as { question: string; name?: string }; const s = d.getSettings(); const ensure = await d.readiness.ensure(s.workspace ?? ""); - if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" }); + if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE }); // Settings (global) supply workspace/provider/model/thinking. The new-question // form sends only the question text. `workspace` selects the tht `-c `. const { id } = await d.tht.sessionNew({ @@ -115,7 +127,7 @@ export function sessionRoutes( } const settings = d.getSettings(); const ensure = await d.readiness.ensure(settings.workspace ?? ""); - if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" }); + 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, @@ -131,7 +143,7 @@ export function sessionRoutes( 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 }); + return reply.code(200).send({ id, alreadyActive: false }); }); app.post("/sessions/:id/close", async (req) => { const id = (req.params as { id: string }).id; @@ -146,6 +158,10 @@ export function sessionRoutes( app.get("/sessions/:id/events", (req, reply) => { const id = (req.params as any).id; const rt = d.mgr.get(id); + const afterId = eventCursor( + req.headers["last-event-id"], + (req.query as { lastEventId?: unknown }).lastEventId, + ); // Add CORS headers manually: reply.raw.writeHead bypasses Fastify's onSend hooks // (where @fastify/cors injects headers), so we must set them explicitly here. const origin = (req.headers.origin as string | undefined) ?? "*"; @@ -160,8 +176,12 @@ export function sessionRoutes( // Send the handshake immediately. Without this, Node waits for the first event body and // proxies/clients cannot establish an idle SSE subscription or inspect its headers. reply.raw.flushHeaders(); - const send = (event: string, data: object) => reply.raw.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`); - const off = d.hub.subscribe(id, send, rt?.bridge.pendingWidget() ?? null); + const send = (event: string, data: object, eventId: number) => + reply.raw.write(`id: ${eventId}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`); + const off = d.hub.subscribe(id, send, { + afterId, + pending: rt?.bridge.pendingWidget() ?? null, + }); req.raw.on("close", off); }); app.post("/sessions/:id/rename", async (req, reply) => { diff --git a/backend/src/sse/sse-hub.ts b/backend/src/sse/sse-hub.ts index 174edb05..1e63abc7 100644 --- a/backend/src/sse/sse-hub.ts +++ b/backend/src/sse/sse-hub.ts @@ -1,48 +1,92 @@ -type Send = (event: string, data: object) => void; +type Send = (event: string, data: object, id: number) => void; /** * Per-session ring buffer of recent events. * - * When the browser opens the SSE connection slightly after session creation - * (React re-render, navigation, etc.), events produced by the Pi process in - * that gap would be lost. The buffer replays them to late subscribers. - * - * `system_event` with `event: "agent_end"` acts as the natural sentinel: - * once a subscriber sees it, the buffer for that session is safe to clear. + * Event ids are transport identity: they are monotonically increasing for the lifetime of a + * session id, including across a cold Resume that clears the old buffer and subscribers. */ const BUFFER_LIMIT = 200; -interface BufferedEvent { event: string; data: object } +interface BufferedEvent { + id: number; + event: string; + data: object; +} + +interface SubscribeOptions { + afterId?: number; + pending?: object | null; +} + +function descriptorId(value: unknown): string | null { + if (!value || typeof value !== "object") return null; + const id = (value as { id?: unknown }).id; + return typeof id === "string" && id.trim() !== "" ? id : null; +} + +function bufferedGateId(item: BufferedEvent): string | null { + if (item.event !== "ui_request") return null; + const request = (item.data as { ui_request?: unknown }).ui_request; + return descriptorId(request); +} export class SseHub { private subs = new Map>(); private buffers = new Map(); + private lastIds = new Map(); - subscribe(sessionId: string, send: Send, pending?: object | null): () => void { + subscribe( + sessionId: string, + send: Send, + options: SubscribeOptions = {}, + ): () => void { if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set()); this.subs.get(sessionId)!.add(send); - // Replay buffered events so late subscribers don't miss the session lifecycle. - const buf = this.buffers.get(sessionId); - if (buf) for (const { event, data } of buf) send(event, data); + const afterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0 + ? options.afterId ?? 0 + : 0; + const buf = this.buffers.get(sessionId) ?? []; + for (const item of buf) { + if (item.id > afterId) send(item.event, item.data, item.id); + } + + // A pending gate normally already exists in the buffer. If ring-buffer eviction removed it, + // give the fallback a fresh transport id and buffer it for future cursor-based reconnects. + // Matching is by stable descriptor identity only, never by event content. + const pendingId = descriptorId(options.pending); + const pendingBuffered = pendingId !== null && buf.some((item) => bufferedGateId(item) === pendingId); + if (pendingId !== null && !pendingBuffered) { + const data = { type: "ui_request", ui_request: options.pending! }; + const item = this.buffer(sessionId, "ui_request", data); + send(item.event, item.data, item.id); + } - // Match the live-event shape (hub.publish sends the full ClientEvent): - // both carry { type, ui_request } so re-emit and live widgets are identical. - if (pending) send("ui_request", { type: "ui_request", ui_request: pending }); return () => this.subs.get(sessionId)?.delete(send); } - publish(sessionId: string, event: string, data: object): void { - // Buffer the event for late subscribers. - let buf = this.buffers.get(sessionId); - if (!buf) { buf = []; this.buffers.set(sessionId, buf); } - buf.push({ event, data }); - if (buf.length > BUFFER_LIMIT) buf.splice(0, buf.length - BUFFER_LIMIT); - - for (const s of this.subs.get(sessionId) ?? []) s(event, data); + 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); + return item.id; } - /** Clear buffer and subscribers for a finished session. */ + private buffer(sessionId: string, event: string, data: object): BufferedEvent { + const id = (this.lastIds.get(sessionId) ?? 0) + 1; + this.lastIds.set(sessionId, id); + let buf = this.buffers.get(sessionId); + if (!buf) { + buf = []; + this.buffers.set(sessionId, buf); + } + const item = { id, event, data }; + buf.push(item); + if (buf.length > BUFFER_LIMIT) buf.splice(0, buf.length - BUFFER_LIMIT); + return item; + } + + /** Clear buffered/runtime bindings while retaining session event-id monotonicity. */ clear(sessionId: string): void { this.buffers.delete(sessionId); this.subs.delete(sessionId); diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 9935ce5b..7f40980c 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -185,7 +185,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" }); + expect(response.json()).toEqual({ id: "s1", alreadyActive: false }); expect(order).toEqual(["teardown:s1", "clear:s1", "reopen", "create", "start"]); }, ); @@ -223,7 +223,7 @@ test("POST resume without a runtime clears stale SSE state before cold start", a const response = await app.inject({ method: "POST", url: "/sessions/crashed/resume" }); await new Promise((resolve) => setImmediate(resolve)); - expect(response.json()).toEqual({ id: "crashed" }); + expect(response.json()).toEqual({ id: "crashed", alreadyActive: false }); expect(order).toEqual(["clear:crashed", "reopen", "create", "start"]); expect(createOptions).toMatchObject({ provider: "local-qwen", model: "qwen3.6-35b-a3b", thinking: "low", mode: "resume", @@ -359,11 +359,13 @@ test("POST resume on an archived session is refused with 409", async () => { expect(res.statusCode).toBe(409); }); -test("POST /sessions refuses with 503 when ollamaEnsure fails (no session created)", async () => { +test("POST /sessions readiness failure returns one fixed public message without raw diagnostics", async () => { let createdCalled = false; + const rawFailure = + "connect https://secret.invalid/ready?token=DO_NOT_LEAK using /srv/private/model-key"; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { - ollamaEnsure: async () => ({ ok: false, stage: "model", error: "modello non installato" }), + ollamaEnsure: async () => ({ ok: false, stage: "model", error: rawFailure }), searchPack: async () => {}, sessionNew: async () => { createdCalled = true; return { id: "s1" }; }, } as any, @@ -372,7 +374,10 @@ test("POST /sessions refuses with 503 when ollamaEnsure fails (no session create }); const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); expect(res.statusCode).toBe(503); - expect(res.json().error).toContain("non installato"); + expect(res.json()).toEqual({ + error: "Session services are not ready. Check configuration and connectivity, then try again.", + }); + expect(res.body).not.toMatch(/secret\.invalid|DO_NOT_LEAK|\/srv\/private\/model-key/); expect(createdCalled).toBe(false); }); @@ -403,10 +408,12 @@ test("POST /sessions/:id/resume returns 409 for a read-only session without call expect(ensureCalled).toBe(false); }); -test("POST /sessions/:id/resume refuses with 503 when ollamaEnsure fails", async () => { +test("POST /sessions/:id/resume readiness failure returns the same fixed public message", async () => { + const rawFailure = + "stderr https://secret.invalid/resume?api_key=DO_NOT_LEAK /srv/private/resume-key"; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { - ollamaEnsure: async () => ({ ok: false, error: "Ollama down" }), + ollamaEnsure: async () => ({ ok: false, error: rawFailure }), sessionShow: async () => ({ status: "open", archived: false }), } as any, getSettings: () => ({ workspace: "psd" }) as any, @@ -414,6 +421,10 @@ test("POST /sessions/:id/resume refuses with 503 when ollamaEnsure fails", async }); const res = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); expect(res.statusCode).toBe(503); + expect(res.json()).toEqual({ + error: "Session services are not ready. Check configuration and connectivity, then try again.", + }); + expect(res.body).not.toMatch(/secret\.invalid|DO_NOT_LEAK|\/srv\/private\/resume-key/); }); test("POST /runtime/prewarm returns 202 without awaiting readiness", async () => { diff --git a/backend/test/session-bridge.test.ts b/backend/test/session-bridge.test.ts index efd64cca..a60d7348 100644 --- a/backend/test/session-bridge.test.ts +++ b/backend/test/session-bridge.test.ts @@ -164,6 +164,26 @@ test("agent_end di Pi diventa un system_event agent_end per il FE", () => { expect(seen).toEqual([{ type: "system_event", event: "agent_end" }]); }); +test("generic Pi system events expose only an allowlisted non-empty event name", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + const seen: any[] = []; + bridge.onClientEvent((event) => seen.push(event)); + + fire({ + type: "system_event", + event: "session_exit", + command: "curl https://secret.invalid/?token=DO_NOT_LEAK", + result: { path: "/srv/private/key" }, + }); + fire({ type: "system_event", event: " ", stderr: "DO_NOT_LEAK_STDERR" }); + fire({ type: "system_event", event: 42, args: "DO_NOT_LEAK_ARGS" }); + + expect(seen).toEqual([{ type: "system_event", event: "session_exit" }]); + expect(Object.keys(seen[0])).toEqual(["type", "event"]); + expect(JSON.stringify(seen)).not.toMatch(/DO_NOT_LEAK|command|result|path|stderr|args/); +}); + test("steer invia un comando steer e riattiva il turno", () => { const { rpc, sent, fire } = fakeRpc(); const bridge = new SessionBridge(rpc); diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts index 7f4f6fe0..581b7a36 100644 --- a/backend/test/sse-hub.test.ts +++ b/backend/test/sse-hub.test.ts @@ -1,16 +1,93 @@ import { test, expect } from "vitest"; import { SseHub } from "../src/sse/sse-hub.js"; -test("re-emette il widget pendente alla sottoscrizione", () => { - const hub = new SseHub(); const sent: any[] = []; - hub.subscribe("s1", (ev, data) => sent.push({ ev, data }), { id: "u1", widget: "select" }); - expect(sent[0]).toEqual({ ev: "ui_request", data: { type: "ui_request", ui_request: { id: "u1", widget: "select" } } }); +test("publish assigns monotonic ids and a subscriber replays only ids newer than its cursor", () => { + const hub = new SseHub(); + const firstId = hub.publish("s1", "text_delta", { text: "one" }); + const secondId = hub.publish("s1", "text_delta", { text: "two" }); + const sent: any[] = []; + + hub.subscribe( + "s1", + (event, data, id) => sent.push({ event, data, id }), + { afterId: firstId }, + ); + const thirdId = hub.publish("s1", "text_delta", { text: "three" }); + + expect([firstId, secondId, thirdId]).toEqual([1, 2, 3]); + expect(sent).toEqual([ + { event: "text_delta", data: { text: "two" }, id: 2 }, + { event: "text_delta", data: { text: "three" }, id: 3 }, + ]); }); -test("publish raggiunge i subscriber e unsubscribe li stacca", () => { - const hub = new SseHub(); const sent: any[] = []; - const off = hub.subscribe("s1", (ev, data) => sent.push({ ev, data })); - hub.publish("s1", "text_delta", { text: "x" }); - off(); hub.publish("s1", "text_delta", { text: "y" }); - expect(sent).toEqual([{ ev: "text_delta", data: { text: "x" } }]); +test("clear drops stale subscribers and buffers but preserves 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.clear("s1"); + expect(hub.publish("s1", "info", { text: "after" })).toBe(2); + const resumed: any[] = []; + hub.subscribe( + "s1", + (event, data, id) => resumed.push({ event, data, id }), + { afterId: 1 }, + ); + + expect(stale).toEqual([{ event: "info", data: { text: "before" }, id: 1 }]); + expect(resumed).toEqual([{ event: "info", data: { text: "after" }, id: 2 }]); +}); + +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" }; + hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor }); + const replayed: any[] = []; + + hub.subscribe( + "s1", + (event, data, id) => replayed.push({ event, data, id }), + { afterId: 0, pending: { ...descriptor, title: "Runtime copy" } }, + ); + + expect(replayed).toEqual([{ + event: "ui_request", + data: { type: "ui_request", ui_request: descriptor }, + id: 1, + }]); + + const alreadySeen: any[] = []; + hub.subscribe( + "s1", + (event, data, id) => alreadySeen.push({ event, data, id }), + { afterId: 1, pending: descriptor }, + ); + expect(alreadySeen).toEqual([]); +}); + +test("an unbuffered pending gate receives one fresh buffered id", () => { + const hub = new SseHub(); + hub.publish("s1", "info", { text: "ready" }); + const first: any[] = []; + + hub.subscribe( + "s1", + (event, data, id) => first.push({ event, data, id }), + { afterId: 1, pending: { id: "gate-1", widget: "select" } }, + ); + expect(first).toEqual([{ + event: "ui_request", + data: { type: "ui_request", ui_request: { id: "gate-1", widget: "select" } }, + id: 2, + }]); + + const reconnect: any[] = []; + hub.subscribe( + "s1", + (event, data, id) => reconnect.push({ event, data, id }), + { afterId: 1, pending: { id: "gate-1", widget: "select" } }, + ); + expect(reconnect).toEqual(first); }); diff --git a/backend/test/sse-route.test.ts b/backend/test/sse-route.test.ts new file mode 100644 index 00000000..2c7f865c --- /dev/null +++ b/backend/test/sse-route.test.ts @@ -0,0 +1,82 @@ +import { test, expect } from "vitest"; +import { buildApp } from "../src/app.js"; +import { loadConfig } from "../src/config.js"; +import { SseHub } from "../src/sse/sse-hub.js"; + +async function readUntil( + reader: ReadableStreamDefaultReader, + predicate: (text: string) => boolean, +): Promise { + const decoder = new TextDecoder(); + let text = ""; + const deadline = Date.now() + 1_000; + while (Date.now() < deadline) { + const result = await Promise.race([ + reader.read(), + new Promise<{ timeout: true }>((resolve) => + setTimeout(() => resolve({ timeout: true }), Math.max(0, deadline - Date.now())), + ), + ]); + if ("timeout" in result || result.done) break; + text += decoder.decode(result.value); + if (predicate(text)) break; + } + return text; +} + +async function captureReplay(options: { query?: string; lastEventId?: string }): Promise { + const hub = new SseHub(); + hub.publish("s1", "info", { type: "info", text: "one" }); + hub.publish("s1", "info", { type: "info", text: "two" }); + hub.publish("s1", "info", { type: "info", text: "three" }); + 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 controller = new AbortController(); + + try { + const response = await fetch( + `http://127.0.0.1:${port}/sessions/s1/events${options.query ?? ""}`, + { + headers: options.lastEventId ? { "Last-Event-ID": options.lastEventId } : undefined, + signal: controller.signal, + }, + ); + const reader = response.body!.getReader(); + const body = await readUntil(reader, (text) => text.includes('"three"')); + await reader.cancel(); + return body; + } finally { + controller.abort(); + await app.close(); + } +} + +test("SSE emits ids and honors the native Last-Event-ID replay cursor", async () => { + const body = await captureReplay({ lastEventId: "1" }); + + expect(body).not.toContain('"one"'); + expect(body).toContain("id: 2\nevent: info\n"); + expect(body).toContain('data: {"type":"info","text":"two"}'); + expect(body).toContain("id: 3\nevent: info\n"); + expect(body).toContain('data: {"type":"info","text":"three"}'); +}); + +test("SSE honors the manual lastEventId query cursor", async () => { + const body = await captureReplay({ query: "?lastEventId=2" }); + + expect(body).not.toContain('"one"'); + expect(body).not.toContain('"two"'); + expect(body).toContain("id: 3\nevent: info\n"); + expect(body).toContain('data: {"type":"info","text":"three"}'); +}); + +test("SSE uses the newer valid cursor when header and query are both present", async () => { + const body = await captureReplay({ query: "?lastEventId=2", lastEventId: "1" }); + + expect(body).not.toContain('"two"'); + expect(body).toContain("id: 3\nevent: info\n"); +}); diff --git a/frontend/src/api/sessions.test.ts b/frontend/src/api/sessions.test.ts index a681079e..58664284 100644 --- a/frontend/src/api/sessions.test.ts +++ b/frontend/src/api/sessions.test.ts @@ -1,6 +1,6 @@ import { http, HttpResponse } from "msw"; import { server } from "../test/msw"; -import { createSession, listSessions, prewarmRuntime } from "./sessions"; +import { createSession, listSessions, prewarmRuntime, resumeSession } from "./sessions"; import { renameSession, setSessionGroup, archiveSession, unarchiveSession, deleteSession, getSessionDocuments, @@ -36,6 +36,14 @@ test("listSessions GETs the array", async () => { expect(rows[0].id).toBe("s1"); }); +test("resumeSession returns the typed runtime disposition", async () => { + server.use(http.post("http://localhost:8787/sessions/s1/resume", () => + HttpResponse.json({ id: "s1", alreadyActive: false }))); + + const result: { id: string; alreadyActive: boolean } = await resumeSession("s1"); + expect(result).toEqual({ id: "s1", alreadyActive: false }); +}); + test("renameSession POSTs {name}", async () => { let body: unknown = null; server.use(http.post("http://localhost:8787/sessions/s1/rename", async ({ request }) => { diff --git a/frontend/src/api/sessions.ts b/frontend/src/api/sessions.ts index 6d947291..53e3cff9 100644 --- a/frontend/src/api/sessions.ts +++ b/frontend/src/api/sessions.ts @@ -1,5 +1,5 @@ import { apiFetch } from "./client"; -import type { SessionSummary, SessionDocument, UiResponse } from "./types"; +import type { ResumeSessionResult, SessionSummary, SessionDocument, UiResponse } from "./types"; export const createSession = (i: { question: string; name?: string }) => apiFetch<{ id: string }>("/sessions", { method: "POST", body: JSON.stringify(i) }); @@ -29,7 +29,7 @@ export const closeSession = (id: string) => apiFetch(`/sessions/${id}/close`, { method: "POST" }); export const resumeSession = (id: string) => - apiFetch(`/sessions/${id}/resume`, { method: "POST" }); + apiFetch(`/sessions/${id}/resume`, { method: "POST" }); export const renameSession = (id: string, name: string) => apiFetch(`/sessions/${id}/rename`, { method: "POST", body: JSON.stringify({ name }) }); diff --git a/frontend/src/api/types.ts b/frontend/src/api/types.ts index faf7a68a..e1d9a6c8 100644 --- a/frontend/src/api/types.ts +++ b/frontend/src/api/types.ts @@ -89,7 +89,7 @@ export type StreamEvent = }; } | { type: "info"; level?: "info" | "warning" | "error"; text: string } - | { type: "system_event"; event: string; [k: string]: unknown }; + | { type: "system_event"; event: string }; export interface SessionSummary { id: string; @@ -104,6 +104,11 @@ export interface SessionSummary { archived: boolean; } +export interface ResumeSessionResult { + id: string; + alreadyActive: boolean; +} + export interface SessionDocument { phase: string; key: string; diff --git a/frontend/src/shell/AppShell.session-mgmt.test.tsx b/frontend/src/shell/AppShell.session-mgmt.test.tsx index 34a9d6c7..ef7b4978 100644 --- a/frontend/src/shell/AppShell.session-mgmt.test.tsx +++ b/frontend/src/shell/AppShell.session-mgmt.test.tsx @@ -17,9 +17,13 @@ const LIST = [ { id: "s2", status: "finalized", question: "Archiviata due", summary: null, created_at: "2026-01-01T00:00:00Z", updated_at: null, author: null, name: null, group: null, archived: true }, ]; +const resumeResult = (id: string, alreadyActive = false) => + HttpResponse.json({ id, alreadyActive }); + beforeEach(() => { FakeEventSource.instances = []; (globalThis as any).EventSource = FakeEventSource; + useSessionStore.getState().resetSession(); server.use( http.get("http://localhost:8787/sessions", () => HttpResponse.json(LIST)), http.get("http://localhost:8787/sessions/:id/documents", () => HttpResponse.json([ @@ -74,7 +78,7 @@ test("Resume from the panel activates the session and closes the panel", async ( server.use( http.post("http://localhost:8787/sessions/:id/resume", ({ params }) => { resumed = params.id as string; - return new HttpResponse(null, { status: 204 }); + return resumeResult(resumed); }), ); wrap(); @@ -88,10 +92,10 @@ test("Resume from the panel activates the session and closes the panel", async ( await waitFor(() => expect(screen.queryByText("Domanda originale")).not.toBeInTheDocument()); // panel closed }); -test("Resume paints the re-entry phase from the manifest (optimistic, before the first gate)", async () => { +test("Resume paints the re-entry phase from the manifest after the cold Resume succeeds", async () => { useSessionStore.getState().resetSession(); server.use( - http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 204 })), + http.post("http://localhost:8787/sessions/:id/resume", () => resumeResult("s1")), http.get("http://localhost:8787/sessions/:id", () => HttpResponse.json({ id: "s1", status: "open", phase: 4 })), ); @@ -102,10 +106,13 @@ test("Resume paints the re-entry phase from the manifest (optimistic, before the await waitFor(() => expect(useSessionStore.getState().currentPhase).toBe("F4")); }); -test("resuming the active session reconnects its EventSource", async () => { +test("an already-active same-session Resume preserves its EventSource and store", async () => { + let resumeCalls = 0; server.use( - http.post("http://localhost:8787/sessions/:id/resume", () => - new HttpResponse(null, { status: 204 })), + http.post("http://localhost:8787/sessions/:id/resume", () => { + resumeCalls += 1; + return resumeResult("s1", resumeCalls > 1); + }), http.get("http://localhost:8787/sessions/:id", () => HttpResponse.json({ id: "s1", status: "open", phase: 1 })), ); @@ -114,18 +121,23 @@ test("resuming the active session reconnects its EventSource", async () => { 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: "Live state" }, "4")); + const before = useSessionStore.getState().activityLog.map((entry) => ({ ...entry })); await userEvent.click(screen.getByText("Attiva uno")); + await screen.findByText("Domanda originale"); await userEvent.click(await screen.findByRole("button", { name: /resume/i })); - await waitFor(() => expect(FakeEventSource.instances).toHaveLength(2)); - expect(first.closed).toBe(true); + await waitFor(() => expect(screen.queryByText("Domanda originale")).not.toBeInTheDocument()); + expect(FakeEventSource.instances).toHaveLength(1); + expect(first.closed).toBe(false); + expect(useSessionStore.getState().activityLog).toEqual(before); }); -test("same-session Resume reconnects only after the deferred POST succeeds", async () => { +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", () => - new HttpResponse(null, { status: 204 })), + resumeResult("s1")), http.get("http://localhost:8787/sessions/:id", () => HttpResponse.json({ id: "s1", status: "open", phase: 1 })), ); @@ -134,6 +146,7 @@ test("same-session Resume reconnects only after the deferred POST succeeds", asy 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")); let releaseResume!: () => void; let markStarted!: () => void; @@ -142,24 +155,62 @@ test("same-session Resume reconnects only after the deferred POST succeeds", asy server.use(http.post("http://localhost:8787/sessions/:id/resume", async () => { markStarted(); await resumeReleased; - return new HttpResponse(null, { status: 204 }); + return resumeResult("s1"); })); await userEvent.click(screen.getByText("Attiva uno")); + await screen.findByText("Domanda originale"); await userEvent.click(await screen.findByRole("button", { name: /resume/i })); await resumeStarted; expect(FakeEventSource.instances).toHaveLength(1); expect(first.closed).toBe(false); + expect(screen.getByText("Domanda originale")).toBeInTheDocument(); + 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")); + expect(useSessionStore.getState().transcript).toEqual([ + { role: "assistant", text: "Still attached" }, + ]); releaseResume(); 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(useSessionStore.getState().transcript).toEqual([]); + expect(useSessionStore.getState().activityLog).toEqual([ + { kind: "lifecycle", phase: null, text: "Resuming session" }, + ]); + + act(() => replacement.emitNamed( + "text_delta", + { type: "text_delta", text: "Post-resume delivery" }, + "9", + )); + expect(useSessionStore.getState().transcript).toEqual([ + { role: "assistant", text: "Post-resume delivery" }, + ]); + expect(useSessionStore.getState().activityLog.filter((entry) => entry.kind === "assistant")) + .toEqual([{ kind: "assistant", phase: "F1", text: "Post-resume delivery" }]); + + const gate = { + type: "ui_request" as const, + ui_request: { id: "post-resume-gate", widget: "select", title: "Review once" }, + }; + act(() => { + replacement.emitNamed("ui_request", gate, "10"); + replacement.emitNamed("ui_request", gate, "11"); + }); + expect(useSessionStore.getState().pendingWidget?.id).toBe("post-resume-gate"); + expect(useSessionStore.getState().activityLog.filter((entry) => entry.kind === "gate")) + .toEqual([{ kind: "gate", phase: "F1", text: "Review once" }]); }); -test("a failed same-session Resume does not create a replacement EventSource", async () => { +test("a failed same-session Resume preserves its source, activity, and document panel", async () => { server.use( http.post("http://localhost:8787/sessions/:id/resume", () => - new HttpResponse(null, { status: 204 })), + resumeResult("s1")), http.get("http://localhost:8787/sessions/:id", () => HttpResponse.json({ id: "s1", status: "open", phase: 1 })), ); @@ -167,6 +218,9 @@ test("a failed same-session Resume does not create a replacement EventSource", a await userEvent.click(await screen.findByText("Attiva uno")); 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: "Keep me" }, "4")); + const before = useSessionStore.getState().activityLog.map((entry) => ({ ...entry })); server.use(http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 409 }))); @@ -175,17 +229,21 @@ test("a failed same-session Resume does not create a replacement EventSource", a await waitFor(() => expect(screen.getByText("Domanda originale")).toBeInTheDocument()); expect(FakeEventSource.instances).toHaveLength(1); + expect(first.closed).toBe(false); + expect(useSessionStore.getState().activityLog).toEqual(before); + act(() => first.emitNamed("text_delta", { type: "text_delta", text: "After failure" }, "5")); + expect(useSessionStore.getState().transcript.at(-1)?.text).toBe("After failure"); }); -test("resuming a different session opens exactly one new EventSource", async () => { +test("resuming a different already-active session binds it only after success", async () => { const other = { ...LIST[0], id: "s3", question: "Attiva tre", group: null, created_at: "2026-01-03T00:00:00Z", }; server.use( http.get("http://localhost:8787/sessions", () => HttpResponse.json([LIST[0], other])), - http.post("http://localhost:8787/sessions/:id/resume", () => - new HttpResponse(null, { status: 204 })), + http.post("http://localhost:8787/sessions/:id/resume", ({ params }) => + resumeResult(params.id as string, params.id === "s3")), http.get("http://localhost:8787/sessions/:id", ({ params }) => HttpResponse.json({ id: params.id, status: "open", phase: params.id === "s3" ? 3 : 1 })), ); @@ -193,28 +251,42 @@ test("resuming a different session opens exactly one new EventSource", async () await userEvent.click(await screen.findByText("Attiva uno")); 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: "Old target" }, "3")); 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).toHaveLength(2); + expect(first.closed).toBe(true); expect(FakeEventSource.instances[1].url).toContain("/sessions/s3/events"); + expect(FakeEventSource.instances[1].url).not.toContain("lastEventId"); + expect(useSessionStore.getState().activityLog).not.toContainEqual(expect.objectContaining({ + text: "Old target", + })); }); -test("a failed resume keeps the panel open and does not activate the session", async () => { +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({ + activityLog: [{ kind: "status", phase: "F7", text: "Preserve activity", level: "info" }], + }); wrap(); await userEvent.click(await screen.findByText("Attiva uno")); await screen.findByText("Domanda originale"); // panel open await userEvent.click(screen.getByRole("button", { name: /resume/i })); // panel stays open (document still visible) after the failed resume await waitFor(() => expect(screen.getByText("Domanda originale")).toBeInTheDocument()); + expect(FakeEventSource.instances).toHaveLength(0); + expect(useSessionStore.getState().activityLog).toEqual([ + { kind: "status", phase: "F7", text: "Preserve activity", level: "info" }, + ]); }); test("closing and reopening Model activity preserves the complete activity log", async () => { wrap(); - server.use(http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 204 }))); + server.use(http.post("http://localhost:8787/sessions/:id/resume", () => resumeResult("s1"))); await userEvent.click(await screen.findByText("Attiva uno")); await userEvent.click(await screen.findByRole("button", { name: /resume/i })); act(() => { @@ -257,7 +329,7 @@ test("session finalization shows the completion banner and returns to landing", useSessionStore.getState().resetSession(); let finalized = false; server.use( - http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 204 })), + http.post("http://localhost:8787/sessions/:id/resume", () => resumeResult("s1")), http.get("http://localhost:8787/sessions", () => HttpResponse.json(finalized ? [{ ...LIST[0], status: "finalized" }] : LIST)), ); @@ -299,7 +371,7 @@ test("renaming a group reassigns its members via setSessionGroup", async () => { test("opening Model activity replaces the session rail with a 40/60 activity and chat layout", async () => { - server.use(http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 204 }))); + server.use(http.post("http://localhost:8787/sessions/:id/resume", () => resumeResult("s1"))); wrap(); await userEvent.click(await screen.findByText("Attiva uno")); await userEvent.click(await screen.findByRole("button", { name: /resume/i })); diff --git a/frontend/src/shell/AppShell.tsx b/frontend/src/shell/AppShell.tsx index 9f77a0c4..6b2ecd72 100644 --- a/frontend/src/shell/AppShell.tsx +++ b/frontend/src/shell/AppShell.tsx @@ -100,22 +100,29 @@ export function AppShell() { }); } async function doResume(id: string) { - const s = sessions.find((x) => x.id === id) ?? null; const reconnectSameSession = activeSessionId === id; - // Optimistic switch: change to the session view IMMEDIATELY so the click feels - // instant (the resume POST spawns a Pi process and can take seconds). The - // working spinner shows straight away; the backend calls run after. - setPanelSession(null); - resetSession(); - recordLifecycle("Resuming session"); - setAwaitingQuestion(false); - // Optimistic: the resume POST is about to hand the ball to the harness. - setAgentActive(true); - setActiveSessionId(id); try { - await resumeSession(id); - if (reconnectSameSession) setStreamGeneration((value) => value + 1); - // Optimistic phase paint: colour the re-entry phase before the first gate. + const result = await resumeSession(id); + setPanelSession(null); + setAwaitingQuestion(false); + + // A running/waiting runtime for the currently selected session is already bound to this + // store and EventSource. Reopening it would replay state and can lose in-flight delivery. + if (result.alreadyActive && reconnectSameSession) return; + + resetSession(); + if (!result.alreadyActive) { + 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. + if (!result.alreadyActive && reconnectSameSession) { + setStreamGeneration((value) => value + 1); + } + + // 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 }; @@ -126,10 +133,6 @@ export function AppShell() { /* non-fatal: the first gate will set the phase */ } } catch { - // Revert the optimistic switch and restore the panel. - resetSession(); - setActiveSessionId(null); - if (s) setPanelSession(s); toast.error("Failed to resume session."); } } diff --git a/frontend/src/store/sessionStore.test.ts b/frontend/src/store/sessionStore.test.ts index 4f659644..091c0b79 100644 --- a/frontend/src/store/sessionStore.test.ts +++ b/frontend/src/store/sessionStore.test.ts @@ -8,6 +8,29 @@ test("ui_request sets pendingWidget", () => { expect(useSessionStore.getState().pendingWidget?.id).toBe("u1"); }); +test("the same gate descriptor id is handled once even after pending is cleared", () => { + const store = useSessionStore.getState(); + const gate = { + type: "ui_request" as const, + ui_request: { id: "gate-1", widget: "select", title: "Review tables" }, + }; + + store.applyEvent(gate); + store.applyEvent({ + type: "ui_request", + ui_request: { ...gate.ui_request, title: "Duplicate transport copy" }, + }); + store.clearPending(); + store.applyEvent(gate); + + const state = useSessionStore.getState(); + expect(state.pendingWidget).toBeNull(); + expect(state.activityLog.filter((entry) => entry.kind === "gate")).toEqual([ + { kind: "gate", phase: null, text: "Review tables" }, + ]); + expect(state.seenGateIds).toEqual(new Set(["gate-1"])); +}); + test("text_delta accumulates into transcript", () => { const s = useSessionStore.getState(); s.applyEvent({ type: "text_delta", text: "Ana" }); diff --git a/frontend/src/store/sessionStore.ts b/frontend/src/store/sessionStore.ts index 9573ce1e..3db2317d 100644 --- a/frontend/src/store/sessionStore.ts +++ b/frontend/src/store/sessionStore.ts @@ -16,6 +16,7 @@ interface SessionState { lastSystemEvent: StreamEvent | null; currentPhase: string | null; phaseError: string | null; + seenGateIds: Set; // True while a Pi turn is in flight. Turned off by the backend-forwarded `agent_end` // system event — the only end-of-turn signal on the final workflow step, where no // follow-up gate arrives to release the spinner. @@ -62,6 +63,7 @@ const empty = { lastSystemEvent: null, currentPhase: null as string | null, phaseError: null as string | null, + seenGateIds: new Set(), agentActive: false, }; @@ -70,9 +72,11 @@ export const useSessionStore = create((set) => ({ applyEvent: (e) => set((st) => { if (e.type === "ui_request") { + if (st.seenGateIds.has(e.ui_request.id)) return {}; const currentPhase = phaseOf(e.ui_request.phase, st.currentPhase); return { pendingWidget: e.ui_request, + seenGateIds: new Set([...st.seenGateIds, e.ui_request.id]), // The turn is blocked on the gate, not ended: keep the agent marked active. agentActive: true, currentPhase, diff --git a/frontend/src/stream/useSessionStream.test.tsx b/frontend/src/stream/useSessionStream.test.tsx index 5fab088b..f5c15d2b 100644 --- a/frontend/src/stream/useSessionStream.test.tsx +++ b/frontend/src/stream/useSessionStream.test.tsx @@ -70,16 +70,46 @@ test("closes the stream on unmount", () => { expect(es.closed).toBe(true); }); -test("reconnects the same session when the generation changes", () => { +test("a same-session generation reconnect includes the last consumed SSE id", () => { const { rerender } = renderHook( ({ generation }) => useSessionStream("s1", generation), { initialProps: { generation: 0 } }, ); const first = FakeEventSource.instances[0]; + act(() => first.emitNamed("text_delta", { type: "text_delta", text: "hello" }, "7")); rerender({ generation: 1 }); expect(first.closed).toBe(true); expect(FakeEventSource.instances).toHaveLength(2); - expect(FakeEventSource.instances[1].url).toBe(first.url); + expect(FakeEventSource.instances[1].url).toBe( + "http://localhost:8787/sessions/s1/events?lastEventId=7", + ); + + act(() => + FakeEventSource.instances[1].emitNamed( + "text_delta", + { type: "text_delta", text: " world" }, + "8", + ), + ); + expect(useSessionStore.getState().transcript).toEqual([ + { role: "assistant", text: "hello world" }, + ]); + expect(useSessionStore.getState().activityLog.filter((entry) => entry.kind === "assistant")) + .toEqual([{ kind: "assistant", phase: null, text: "hello world" }]); +}); + +test("changing the session resets the manual reconnect cursor", () => { + const { rerender } = renderHook( + ({ sessionId }) => useSessionStream(sessionId), + { initialProps: { sessionId: "s1" as string | null } }, + ); + const first = FakeEventSource.instances[0]; + act(() => first.emitNamed("info", { type: "info", text: "ready" }, "12")); + + rerender({ sessionId: "s2" }); + + expect(first.closed).toBe(true); + expect(FakeEventSource.instances[1].url).toBe("http://localhost:8787/sessions/s2/events"); }); diff --git a/frontend/src/stream/useSessionStream.ts b/frontend/src/stream/useSessionStream.ts index a54ca0fd..76d3086a 100644 --- a/frontend/src/stream/useSessionStream.ts +++ b/frontend/src/stream/useSessionStream.ts @@ -1,4 +1,4 @@ -import { useEffect, useState } from "react"; +import { useEffect, useRef, useState } from "react"; import { BASE } from "../api/client"; import { joinBackendPath } from "../api/runtime-config"; import { useSessionStore } from "../store/sessionStore"; @@ -7,15 +7,24 @@ import type { StreamEvent } from "../api/types"; export function useSessionStream(sessionId: string | null, generation = 0) { const [connected, setConnected] = useState(false); const applyEvent = useSessionStore((s) => s.applyEvent); + const cursor = useRef({ sessionId: null as string | null, lastEventId: "" }); + + if (cursor.current.sessionId !== sessionId) { + cursor.current = { sessionId, lastEventId: "" }; + } useEffect(() => { if (!sessionId) return; - const es = new EventSource(joinBackendPath(BASE, `/sessions/${sessionId}/events`)); + 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 handle = (ev: MessageEvent) => { + if (ev.lastEventId) cursor.current.lastEventId = ev.lastEventId; try { applyEvent(JSON.parse(ev.data) as StreamEvent); } catch { diff --git a/frontend/src/test/fakeEventSource.ts b/frontend/src/test/fakeEventSource.ts index 2865d087..b239e777 100644 --- a/frontend/src/test/fakeEventSource.ts +++ b/frontend/src/test/fakeEventSource.ts @@ -1,27 +1,27 @@ export class FakeEventSource { static instances: FakeEventSource[] = []; - onmessage: ((e: { data: string }) => void) | null = null; + onmessage: ((e: { data: string; lastEventId: string }) => void) | null = null; onopen: (() => void) | null = null; onerror: (() => void) | null = null; closed = false; - private listeners = new Map void>>(); + private listeners = new Map void>>(); constructor(public url: string) { FakeEventSource.instances.push(this); } /** Dispatch to the default onmessage handler (unnamed `event: message`). */ - emit(obj: unknown) { - this.onmessage?.({ data: JSON.stringify(obj) }); + emit(obj: unknown, lastEventId = "") { + this.onmessage?.({ data: JSON.stringify(obj), lastEventId }); } /** Dispatch to handlers registered for a NAMED event (e.g. "ui_request"). */ - emitNamed(type: string, obj: unknown) { + emitNamed(type: string, obj: unknown, lastEventId = "") { const data = JSON.stringify(obj); - for (const h of this.listeners.get(type) ?? []) h({ data }); + for (const h of this.listeners.get(type) ?? []) h({ data, lastEventId }); } - addEventListener(type: string, handler: (e: { data: string }) => void) { + addEventListener(type: string, handler: (e: { data: string; lastEventId: string }) => void) { if (!this.listeners.has(type)) this.listeners.set(type, new Set()); this.listeners.get(type)!.add(handler); } - removeEventListener(type: string, handler: (e: { data: string }) => void) { + removeEventListener(type: string, handler: (e: { data: string; lastEventId: string }) => void) { this.listeners.get(type)?.delete(handler); } close() {