From 2958b32fd517d9ff159771e80be567dce374a960 Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 20 Jul 2026 01:32:47 +0200 Subject: [PATCH] =?UTF-8?q?fix(backend):=20generation-aware=20SSE=20event?= =?UTF-8?q?=20ids=20=E2=80=94=20stale=20cursors=20can=20no=20longer=20eat?= =?UTF-8?q?=20events?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Audit finding 3.1 (high, 3/3 reviewer consensus). Event ids restart at 1 when the backend restarts; a browser auto-reconnect carrying the old numeric Last-Event-ID was honored whenever the new process had already emitted that many events, silently suppressing fresh events (same ids, different content). The previous guard only caught cursor > lastId. Wire ids are now ":" (generation = per-hub instance token; seq = the existing per-session monotonic counter). The hub parses raw header/query candidates itself: other-generation and legacy bare- number cursors are stale → replay from the beginning; same-generation cursors keep the newest-valid-wins behavior. EventSource treats ids as opaque, so no frontend change. Finding 3.2 (eviction) resolved by NOT evicting: close keeps the seq counter on purpose (sessions reopen; monotonicity is what makes old cursors detectable) — documented at the call site; buffers are emptied by clear() and ring-bounded at 200. Co-Authored-By: Claude Fable 5 --- backend/src/routes/sessions.ts | 23 ++----- backend/src/sse/sse-hub.ts | 60 +++++++++++------ backend/test/sse-hub.test.ts | 115 +++++++++++++++++++++------------ backend/test/sse-route.test.ts | 76 ++++++++++++++-------- 4 files changed, 171 insertions(+), 103 deletions(-) diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 5c90f744..06601508 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -16,16 +16,6 @@ const RESUME_FAILURE_MESSAGE = const DWH_UNREACHABLE_MESSAGE = "Cannot start a session: the data warehouse is unreachable. Check the VPN connection and 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, d: { @@ -370,6 +360,9 @@ export function sessionRoutes( } catch { return storageFailure(reply); } finally { + // clear, NOT forget: a closed session can be reopened, and the per-session seq + // monotonicity is what keeps a browser's old cursor detectable. The buffer is + // emptied here; only delete discards the id counter. d.hub.clear(id); } return { closed: true }; @@ -383,10 +376,6 @@ export function sessionRoutes( if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" }); } catch { return storageFailure(reply); } 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) ?? "*"; @@ -401,10 +390,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, eventId: number) => + const send = (event: string, data: object, eventId: string) => reply.raw.write(`id: ${eventId}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`); const off = d.hub.subscribe(id, send, { - afterId, + // Raw cursor candidates: the hub parses ":" and treats any + // other-generation (or legacy numeric) cursor as stale → replay from the start. + after: [req.headers["last-event-id"], (req.query as { lastEventId?: unknown }).lastEventId], 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. diff --git a/backend/src/sse/sse-hub.ts b/backend/src/sse/sse-hub.ts index 6577e622..f7327f34 100644 --- a/backend/src/sse/sse-hub.ts +++ b/backend/src/sse/sse-hub.ts @@ -1,4 +1,4 @@ -type Send = (event: string, data: object, id: number) => void; +type Send = (event: string, data: object, id: string) => void; interface Subscriber { send: Send; @@ -9,8 +9,13 @@ interface Subscriber { /** * Per-session ring buffer of recent events. * - * 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. + * Event ids are transport identity: `:`, where seq is monotonically + * increasing for the lifetime of a session id in THIS process (including across a cold + * Resume that clears the old buffer and subscribers), and generation identifies the hub + * instance. After a backend restart seqs start over from 1: without the generation part, + * a browser auto-reconnect carrying an old numeric cursor would silently suppress the new + * process's first events (same ids, different content). A cursor whose generation does not + * match this instance is stale by definition and replays from the beginning. */ const BUFFER_LIMIT = 200; @@ -21,7 +26,8 @@ interface BufferedEvent { } interface SubscribeOptions { - afterId?: number; + /** Raw cursor candidates (Last-Event-ID header, query param) — parsed by the hub. */ + after?: unknown[]; pending?: object | null; close?: () => void; } @@ -42,6 +48,32 @@ export class SseHub { private subs = new Map>(); private buffers = new Map(); private lastIds = new Map(); + private readonly generation = + `${Date.now().toString(36)}${Math.random().toString(36).slice(2, 6)}`; + + private wireId(seq: number): string { + return `${this.generation}:${seq}`; + } + + /** + * Resolve the raw cursor candidates (header + query) to a seq in THIS generation, + * taking the newest valid one. A cursor from another generation (or the legacy + * bare-number format) is stale → 0, i.e. replay from the beginning. + */ + private parseAfter(sessionId: string, candidates: unknown[] = []): number { + let after = 0; + for (const value of candidates.flatMap((item) => (Array.isArray(item) ? item : [item]))) { + if (typeof value !== "string") continue; + const m = /^([^:]+):(\d+)$/.exec(value); + if (!m || m[1] !== this.generation) continue; + const seq = Number(m[2]); + // Same generation, but a seq this session never produced cannot identify an event. + if (Number.isSafeInteger(seq) && seq <= (this.lastIds.get(sessionId) ?? 0)) { + after = Math.max(after, seq); + } + } + return after; + } subscribe( sessionId: string, @@ -52,18 +84,10 @@ export class SseHub { const subscriber = { send, close: options.close ?? (() => undefined), closed: false }; this.subs.get(sessionId)!.add(subscriber); - 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 afterId = this.parseAfter(sessionId, options.after); const buf = this.buffers.get(sessionId) ?? []; for (const item of buf) { - if (item.id > afterId) send(item.event, item.data, item.id); + if (item.id > afterId) send(item.event, item.data, this.wireId(item.id)); } // A pending gate normally already exists in the buffer. If ring-buffer eviction removed it, @@ -74,18 +98,18 @@ export class SseHub { 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); + send(item.event, item.data, this.wireId(item.id)); } return () => this.subs.get(sessionId)?.delete(subscriber); } - publish(sessionId: string, event: string, data: object): number { + publish(sessionId: string, event: string, data: object): string { const item = this.buffer(sessionId, event, data); for (const subscriber of this.subs.get(sessionId) ?? []) { - subscriber.send(event, data, item.id); + subscriber.send(event, data, this.wireId(item.id)); } - return item.id; + return this.wireId(item.id); } private buffer(sessionId: string, event: string, data: object): BufferedEvent { diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts index 35343eb8..71ef0641 100644 --- a/backend/test/sse-hub.test.ts +++ b/backend/test/sse-hub.test.ts @@ -1,6 +1,12 @@ import { test, expect } from "vitest"; import { SseHub } from "../src/sse/sse-hub.js"; +// Wire ids are ":": seq is the per-session monotonic counter, +// generation identifies the hub instance (process). Tests derive both from +// published ids instead of hardcoding the random generation. +const seqOf = (id: string) => Number(id.split(":")[1]); +const genOf = (id: string) => id.split(":")[0]; + 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" }); @@ -10,14 +16,15 @@ test("publish assigns monotonic ids and a subscriber replays only ids newer than hub.subscribe( "s1", (event, data, id) => sent.push({ event, data, id }), - { afterId: firstId }, + { after: [firstId] }, ); const thirdId = hub.publish("s1", "text_delta", { text: "three" }); - expect([firstId, secondId, thirdId]).toEqual([1, 2, 3]); + expect([firstId, secondId, thirdId].map(seqOf)).toEqual([1, 2, 3]); + expect(new Set([firstId, secondId, thirdId].map(genOf)).size).toBe(1); expect(sent).toEqual([ - { event: "text_delta", data: { text: "two" }, id: 2 }, - { event: "text_delta", data: { text: "three" }, id: 3 }, + { event: "text_delta", data: { text: "two" }, id: secondId }, + { event: "text_delta", data: { text: "three" }, id: thirdId }, ]); }); @@ -38,33 +45,36 @@ test("clear closes every stale subscriber once and replays post-clear events fro (event, data, id) => secondStale.push({ event, data, id }), { close: () => { closeCalls[1] += 1; offSecond(); } }, ); - expect(hub.publish("s1", "info", { text: "before" })).toBe(1); + const beforeId = hub.publish("s1", "info", { text: "before" }); + expect(seqOf(beforeId)).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", { + const afterId = hub.publish("s1", "info", { text: "after" }); + expect(seqOf(afterId)).toBe(2); + const gateId = hub.publish("s1", "ui_request", { type: "ui_request", ui_request: { id: "gate-1", widget: "select" }, - })).toBe(3); + }); + expect(seqOf(gateId)).toBe(3); const resumed: any[] = []; hub.subscribe( "s1", (event, data, id) => resumed.push({ event, data, id }), - { afterId: 1 }, + { after: [beforeId] }, ); - const before = [{ event: "info", data: { text: "before" }, id: 1 }]; + const before = [{ event: "info", data: { text: "before" }, id: beforeId }]; expect(firstStale).toEqual(before); expect(secondStale).toEqual(before); expect(resumed).toEqual([ - { event: "info", data: { text: "after" }, id: 2 }, + { event: "info", data: { text: "after" }, id: afterId }, { event: "ui_request", data: { type: "ui_request", ui_request: { id: "gate-1", widget: "select" } }, - id: 3, + id: gateId, }, ]); }); @@ -86,95 +96,118 @@ test("forget closes every subscriber once and resets the per-session id sequence (event, data, id) => secondStale.push({ event, data, id }), { close: () => { closeCalls[1] += 1; offSecond(); } }, ); - expect(hub.publish("s1", "info", { text: "before" })).toBe(1); + const beforeId = hub.publish("s1", "info", { text: "before" }); + expect(seqOf(beforeId)).toBe(1); hub.forget("s1"); expect(closeCalls).toEqual([1, 1]); - expect(hub.publish("s1", "info", { text: "after" })).toBe(1); + const afterId = hub.publish("s1", "info", { text: "after" }); + expect(seqOf(afterId)).toBe(1); const fresh: any[] = []; hub.subscribe("s1", (event, data, id) => fresh.push({ event, data, id })); - const before = [{ event: "info", data: { text: "before" }, id: 1 }]; + const before = [{ event: "info", data: { text: "before" }, id: beforeId }]; expect(firstStale).toEqual(before); expect(secondStale).toEqual(before); - expect(fresh).toEqual([{ event: "info", data: { text: "after" }, id: 1 }]); + expect(fresh).toEqual([{ event: "info", data: { text: "after" }, id: afterId }]); }); 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 gateId = 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" } }, + { after: [], pending: { ...descriptor, title: "Runtime copy" } }, ); expect(replayed).toEqual([{ event: "ui_request", data: { type: "ui_request", ui_request: descriptor }, - id: 1, + id: gateId, }]); const alreadySeen: any[] = []; hub.subscribe( "s1", (event, data, id) => alreadySeen.push({ event, data, id }), - { afterId: 1, pending: descriptor }, + { after: [gateId], 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 readyId = 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" } }, + { after: [readyId], 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, + id: `${genOf(readyId)}:2`, }]); const reconnect: any[] = []; hub.subscribe( "s1", (event, data, id) => reconnect.push({ event, data, id }), - { afterId: 1, pending: { id: "gate-1", widget: "select" } }, + { after: [readyId], pending: { id: "gate-1", widget: "select" } }, ); 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 }); +test("a cursor from a previous hub generation is stale and replays from the beginning", () => { + // Simulates a backend restart: the browser auto-reconnects with the Last-Event-ID it + // got from the OLD process. The new hub must ignore it even when its own seqs have + // already reached (or passed) that number — same ids, different content. + const oldHub = new SseHub(); + oldHub.publish("s1", "info", { text: "old-1" }); + const oldCursor = oldHub.publish("s1", "info", { text: "old-2" }); + + const newHub = new SseHub(); + const freshFirst = newHub.publish("s1", "info", { text: "new-1" }); + const freshSecond = newHub.publish("s1", "info", { text: "new-2" }); const replayed: any[] = []; - hub.subscribe( + newHub.subscribe( "s1", (event, data, id) => replayed.push({ event, data, id }), - { afterId: 900, pending: { ...descriptor } }, + { after: [oldCursor] }, ); 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, - }, + { event: "info", data: { text: "new-1" }, id: freshFirst }, + { event: "info", data: { text: "new-2" }, id: freshSecond }, ]); }); + +test("legacy bare-number and too-new same-generation cursors replay from the beginning", () => { + const hub = new SseHub(); + const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" }; + const firstId = hub.publish("s1", "info", { type: "info", text: "fresh generation" }); + const gateId = hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor }); + const expected = [ + { event: "info", data: { type: "info", text: "fresh generation" }, id: firstId }, + { event: "ui_request", data: { type: "ui_request", ui_request: descriptor }, id: gateId }, + ]; + + const legacy: any[] = []; + hub.subscribe("s1", (event, data, id) => legacy.push({ event, data, id }), { after: ["2"] }); + expect(legacy).toEqual(expected); + + const tooNew: any[] = []; + hub.subscribe( + "s1", + (event, data, id) => tooNew.push({ event, data, id }), + { after: [`${genOf(firstId)}:900`], pending: { ...descriptor } }, + ); + expect(tooNew).toEqual(expected); +}); diff --git a/backend/test/sse-route.test.ts b/backend/test/sse-route.test.ts index 60c3cc98..f8ab7f8c 100644 --- a/backend/test/sse-route.test.ts +++ b/backend/test/sse-route.test.ts @@ -3,6 +3,8 @@ import { buildApp } from "../src/app.js"; import { loadConfig } from "../src/config.js"; import { SseHub } from "../src/sse/sse-hub.js"; +const seqOf = (id: string) => Number(id.split(":")[1]); + async function readUntil( reader: ReadableStreamDefaultReader, predicate: (text: string) => boolean, @@ -24,11 +26,15 @@ async function readUntil( return text; } -async function captureReplay(options: { query?: string; lastEventId?: string }): Promise { +async function captureReplay( + options: { query?: (ids: string[]) => string; lastEventId?: (ids: string[]) => string }, +): Promise<{ body: string; ids: string[] }> { 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 ids = [ + 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, @@ -38,17 +44,19 @@ async function captureReplay(options: { query?: string; lastEventId?: string }): const controller = new AbortController(); try { + const lastEventId = options.lastEventId?.(ids); + const query = options.query ? `?lastEventId=${encodeURIComponent(options.query(ids))}` : ""; const response = await fetch( - `http://127.0.0.1:${port}/sessions/s1/events${options.query ?? ""}`, + `http://127.0.0.1:${port}/sessions/s1/events${query}`, { - headers: options.lastEventId ? { "Last-Event-ID": options.lastEventId } : undefined, + headers: lastEventId ? { "Last-Event-ID": lastEventId } : undefined, signal: controller.signal, }, ); const reader = response.body!.getReader(); const body = await readUntil(reader, (text) => text.includes('"three"')); await reader.cancel(); - return body; + return { body, ids }; } finally { controller.abort(); await app.close(); @@ -67,36 +75,45 @@ async function expectStreamEnd(reader: ReadableStreamDefaultReader): } test("SSE emits ids and honors the native Last-Event-ID replay cursor", async () => { - const body = await captureReplay({ lastEventId: "1" }); + const { body, ids } = await captureReplay({ lastEventId: (published) => published[0] }); expect(body).not.toContain('"one"'); - expect(body).toContain("id: 2\nevent: info\n"); + expect(body).toContain(`id: ${ids[1]}\nevent: info\n`); expect(body).toContain('data: {"type":"info","text":"two"}'); - expect(body).toContain("id: 3\nevent: info\n"); + expect(body).toContain(`id: ${ids[2]}\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" }); + const { body, ids } = await captureReplay({ query: (published) => published[1] }); expect(body).not.toContain('"one"'); expect(body).not.toContain('"two"'); - expect(body).toContain("id: 3\nevent: info\n"); + expect(body).toContain(`id: ${ids[2]}\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" }); + const { body } = await captureReplay({ + query: (published) => published[1], + lastEventId: (published) => published[0], + }); expect(body).not.toContain('"two"'); - expect(body).toContain("id: 3\nevent: info\n"); + expect(body).toContain('"three"'); }); -test("SSE resets a cursor from an older process and emits a buffered pending gate once", async () => { +test("SSE resets a cursor from an older backend generation and emits a buffered pending gate once", async () => { + // The cursor comes from a PREVIOUS process: same session, ids restarted. It must be + // treated as stale (replay from the beginning), not honored against the new ids. + const oldGeneration = new SseHub(); + oldGeneration.publish("s1", "info", { type: "info", text: "old" }); + const staleCursor = oldGeneration.publish("s1", "info", { type: "info", text: "older" }); + 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 infoId = hub.publish("s1", "info", { type: "info", text: "fresh generation" }); + const gateId = hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor }); const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), { hub, thtRunner: {} as any, @@ -110,17 +127,17 @@ test("SSE resets a cursor from an older process and emits a buffered pending gat try { const response = await fetch(`http://127.0.0.1:${port}/sessions/s1/events`, { - headers: { "Last-Event-ID": "900" }, + headers: { "Last-Event-ID": staleCursor }, 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(`id: ${infoId}\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(`id: ${gateId}\nevent: ui_request\n`); expect(body).toContain('"ui_request":{"id":"gate-fresh","widget":"select","title":"Choose"}'); } finally { controller.abort(); @@ -145,32 +162,35 @@ test("clear ends every live SSE response and a cursor reconnect replays post-cle ))); const readers = responses.map((response) => response.body!.getReader()); - expect(hub.publish("s1", "info", { type: "info", text: "before" })).toBe(1); + const beforeId = hub.publish("s1", "info", { type: "info", text: "before" }); + expect(seqOf(beforeId)).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); + expect(initial.every((body) => body.includes(`id: ${beforeId}\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", { + const afterId = hub.publish("s1", "info", { type: "info", text: "after" }); + expect(seqOf(afterId)).toBe(2); + const gateId = hub.publish("s1", "ui_request", { type: "ui_request", ui_request: { id: "gate-1", widget: "select" }, - })).toBe(3); + }); + expect(seqOf(gateId)).toBe(3); const reconnected = await fetch( `http://127.0.0.1:${port}/sessions/s1/events`, { - headers: { "Last-Event-ID": "1" }, + headers: { "Last-Event-ID": beforeId }, 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(`id: ${afterId}\nevent: info\n`); expect(replay).toContain('data: {"type":"info","text":"after"}'); - expect(replay).toContain("id: 3\nevent: ui_request\n"); + expect(replay).toContain(`id: ${gateId}\nevent: ui_request\n`); expect(replay).toContain('"ui_request":{"id":"gate-1","widget":"select"}'); await reconnectReader.cancel(); } finally {