From 0bcec1591ba46a946f3a965a91230f7166ba6a77 Mon Sep 17 00:00:00 2001 From: User Date: Wed, 15 Jul 2026 02:08:11 +0200 Subject: [PATCH] fix: close delete resume and SSE replay races --- backend/src/sse/sse-hub.ts | 8 +- backend/test/sse-hub.test.ts | 27 ++++ backend/test/sse-route.test.ts | 36 +++++ .../src/shell/AppShell.session-mgmt.test.tsx | 134 ++++++++++++++++++ frontend/src/shell/AppShell.tsx | 8 +- 5 files changed, 208 insertions(+), 5 deletions(-) diff --git a/backend/src/sse/sse-hub.ts b/backend/src/sse/sse-hub.ts index c805c69b..6577e622 100644 --- a/backend/src/sse/sse-hub.ts +++ b/backend/src/sse/sse-hub.ts @@ -52,9 +52,15 @@ export class SseHub { 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 + const requestedAfterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0 ? options.afterId ?? 0 : 0; + // Native EventSource persists Last-Event-ID across a backend process restart. A cursor newer + // than anything this hub generation has produced cannot identify an event in this process, + // so replay the fresh generation from its beginning instead of suppressing every low id. + const afterId = requestedAfterId > (this.lastIds.get(sessionId) ?? 0) + ? 0 + : requestedAfterId; const buf = this.buffers.get(sessionId) ?? []; for (const item of buf) { if (item.id > afterId) send(item.event, item.data, item.id); diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts index a61ba2d6..35343eb8 100644 --- a/backend/test/sse-hub.test.ts +++ b/backend/test/sse-hub.test.ts @@ -151,3 +151,30 @@ test("an unbuffered pending gate receives one fresh buffered id", () => { ); expect(reconnect).toEqual(first); }); + +test("a cursor newer than this hub generation replays low ids and a buffered gate once", () => { + const hub = new SseHub(); + const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" }; + hub.publish("s1", "info", { type: "info", text: "fresh generation" }); + 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: 900, pending: { ...descriptor } }, + ); + + expect(replayed).toEqual([ + { + event: "info", + data: { type: "info", text: "fresh generation" }, + id: 1, + }, + { + event: "ui_request", + data: { type: "ui_request", ui_request: descriptor }, + id: 2, + }, + ]); +}); diff --git a/backend/test/sse-route.test.ts b/backend/test/sse-route.test.ts index a1ecdb00..60c3cc98 100644 --- a/backend/test/sse-route.test.ts +++ b/backend/test/sse-route.test.ts @@ -92,6 +92,42 @@ test("SSE uses the newer valid cursor when header and query are both present", a expect(body).toContain("id: 3\nevent: info\n"); }); +test("SSE resets a cursor from an older process and emits a buffered pending gate once", async () => { + const hub = new SseHub(); + const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" }; + hub.publish("s1", "info", { type: "info", text: "fresh generation" }); + hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor }); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), { + hub, + thtRunner: {} as any, + mgr: { + get: () => ({ bridge: { pendingWidget: () => ({ ...descriptor }) } }), + } 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`, { + headers: { "Last-Event-ID": "900" }, + signal: controller.signal, + }); + const reader = response.body!.getReader(); + const body = await readUntil(reader, (text) => text.includes('"gate-fresh"')); + await reader.cancel(); + + expect(body).toContain("id: 1\nevent: info\n"); + expect(body).toContain('data: {"type":"info","text":"fresh generation"}'); + expect(body.match(/event: ui_request/g)).toHaveLength(1); + expect(body).toContain("id: 2\nevent: ui_request\n"); + expect(body).toContain('"ui_request":{"id":"gate-fresh","widget":"select","title":"Choose"}'); + } finally { + controller.abort(); + await app.close(); + } +}); + 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" }), { diff --git a/frontend/src/shell/AppShell.session-mgmt.test.tsx b/frontend/src/shell/AppShell.session-mgmt.test.tsx index 3c9c2d5e..77e20f88 100644 --- a/frontend/src/shell/AppShell.session-mgmt.test.tsx +++ b/frontend/src/shell/AppShell.session-mgmt.test.tsx @@ -29,6 +29,7 @@ function deferred() { beforeEach(() => { FakeEventSource.instances = []; (globalThis as any).EventSource = FakeEventSource; + (globalThis as any).PointerEvent = MouseEvent; useSessionStore.getState().resetSession(); server.use( http.get("http://localhost:8787/sessions", () => HttpResponse.json(LIST)), @@ -489,6 +490,139 @@ test("starting a new question invalidates a pending Resume intent", async () => expect(screen.getByText(/type your question/i)).toBeInTheDocument(); }); +test("a successful Delete invalidates an earlier pending Resume for the same target", async () => { + const resumeGate = deferred(); + const resumeStarted = deferred(); + const deleteCompleted = deferred(); + server.use( + http.post("http://localhost:8787/sessions/:id/resume", async () => { + resumeStarted.resolve(); + await resumeGate.promise; + return resumeResult("s1"); + }), + http.delete("http://localhost:8787/sessions/:id", ({ params }) => { + expect(params.id).toBe("s1"); + deleteCompleted.resolve(); + return new HttpResponse(null, { status: 204 }); + }), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await resumeStarted.promise; + screen.getByRole("checkbox", { name: "Select Attiva uno" }).focus(); + await userEvent.keyboard(" "); + await userEvent.click(await screen.findByRole("button", { name: "Delete 1 selected sessions" })); + await deleteCompleted.promise; + + resumeGate.resolve(); + await new Promise((resolve) => setImmediate(resolve)); + + expect(FakeEventSource.instances).toHaveLength(0); + expect(screen.getByTestId("session-item-s1")).toHaveAttribute("data-active", "false"); +}); + +test("a successful Delete detaches a Resume that commits while Delete is pending", async () => { + const deleteGate = deferred(); + const deleteStarted = deferred(); + server.use( + http.delete("http://localhost:8787/sessions/:id", async ({ params }) => { + expect(params.id).toBe("s1"); + deleteStarted.resolve(); + await deleteGate.promise; + return 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: 1 })), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + screen.getByRole("checkbox", { name: "Select Attiva uno" }).focus(); + await userEvent.keyboard(" "); + await userEvent.click(await screen.findByRole("button", { name: "Delete 1 selected sessions" })); + await deleteStarted.promise; + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + expect(screen.getByTestId("session-item-s1")).toHaveAttribute("data-active", "true"); + + deleteGate.resolve(); + + await waitFor(() => { + expect(screen.getByTestId("session-item-s1")).toHaveAttribute("data-active", "false"); + }); + expect(FakeEventSource.instances[0].closed).toBe(true); +}); + +test("deleting another session does not invalidate a pending Resume", async () => { + const resumeGate = deferred(); + const resumeStarted = deferred(); + const deleteCompleted = deferred(); + server.use( + http.post("http://localhost:8787/sessions/:id/resume", async () => { + resumeStarted.resolve(); + await resumeGate.promise; + return resumeResult("s1"); + }), + http.delete("http://localhost:8787/sessions/:id", ({ params }) => { + expect(params.id).toBe("s2"); + deleteCompleted.resolve(); + return new HttpResponse(null, { status: 204 }); + }), + ); + 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: /archive/i })); + (await screen.findByRole("checkbox", { name: "Select Archiviata due" })).focus(); + await userEvent.keyboard(" "); + await userEvent.click(await screen.findByRole("button", { name: "Delete 1 selected sessions" })); + await deleteCompleted.promise; + + resumeGate.resolve(); + + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + expect(FakeEventSource.instances[0].url).toContain("/sessions/s1/events"); + expect(screen.getByTestId("session-item-s1")).toHaveAttribute("data-active", "true"); +}); + +test("a failed Delete does not invalidate a pending Resume for its target", async () => { + const resumeGate = deferred(); + const resumeStarted = deferred(); + const deleteFailed = deferred(); + server.use( + http.post("http://localhost:8787/sessions/:id/resume", async () => { + resumeStarted.resolve(); + await resumeGate.promise; + return resumeResult("s1"); + }), + http.delete("http://localhost:8787/sessions/:id", ({ params }) => { + expect(params.id).toBe("s1"); + deleteFailed.resolve(); + return new HttpResponse(null, { status: 500 }); + }), + ); + wrap(); + + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await resumeStarted.promise; + screen.getByRole("checkbox", { name: "Select Attiva uno" }).focus(); + await userEvent.keyboard(" "); + await userEvent.click(await screen.findByRole("button", { name: "Delete 1 selected sessions" })); + await deleteFailed.promise; + + resumeGate.resolve(); + + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + expect(FakeEventSource.instances[0].url).toContain("/sessions/s1/events"); + expect(screen.getByTestId("session-item-s1")).toHaveAttribute("data-active", "true"); +}); + 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 d066f783..58036829 100644 --- a/frontend/src/shell/AppShell.tsx +++ b/frontend/src/shell/AppShell.tsx @@ -240,16 +240,16 @@ 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( targets.filter((_, index) => results[index].status === "fulfilled").map((session) => session.id), ); + const deletedActiveSession = deletedIds.has(activeSessionIdRef.current ?? ""); + const deletedResumeTarget = deletedIds.has(latestResumeIntentRef.current?.id ?? ""); + if (deletedActiveSession || deletedResumeTarget) invalidateResumeIntent(); if (deletedIds.has(panelSession?.id ?? "")) setPanelSession(null); - if (deletedIds.has(activeSessionId ?? "")) { resetSession(); selectActiveSession(null); } + if (deletedActiveSession) { resetSession(); selectActiveSession(null); } setSelectedSessionIds((current) => new Set([...current].filter((id) => !deletedIds.has(id)))); refresh(); if (deletedIds.size !== targets.length) {