fix(backend): generation-aware SSE event ids — stale cursors can no longer eat events
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>:<seq>" (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 <noreply@anthropic.com>
This commit is contained in:
@@ -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 "<generation>:<seq>" 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.
|
||||
|
||||
+42
-18
@@ -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: `<generation>:<seq>`, 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<string, Set<Subscriber>>();
|
||||
private buffers = new Map<string, BufferedEvent[]>();
|
||||
private lastIds = new Map<string, number>();
|
||||
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 {
|
||||
|
||||
@@ -1,6 +1,12 @@
|
||||
import { test, expect } from "vitest";
|
||||
import { SseHub } from "../src/sse/sse-hub.js";
|
||||
|
||||
// Wire ids are "<generation>:<seq>": 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);
|
||||
});
|
||||
|
||||
@@ -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<Uint8Array>,
|
||||
predicate: (text: string) => boolean,
|
||||
@@ -24,11 +26,15 @@ async function readUntil(
|
||||
return text;
|
||||
}
|
||||
|
||||
async function captureReplay(options: { query?: string; lastEventId?: string }): Promise<string> {
|
||||
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<Uint8Array>):
|
||||
}
|
||||
|
||||
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 {
|
||||
|
||||
Reference in New Issue
Block a user