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" }); const secondId = hub.publish("s1", "text_delta", { text: "two" }); const sent: any[] = []; hub.subscribe( "s1", (event, data, id) => sent.push({ event, data, id }), { after: [firstId] }, ); const thirdId = hub.publish("s1", "text_delta", { text: "three" }); 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: secondId }, { event: "text_delta", data: { text: "three" }, id: thirdId }, ]); }); test("clear closes every stale subscriber once and replays post-clear events from the cursor", () => { const hub = new SseHub(); const firstStale: any[] = []; const secondStale: any[] = []; const closeCalls = [0, 0]; let offFirst: () => void = () => undefined; let offSecond: () => void = () => undefined; offFirst = hub.subscribe( "s1", (event, data, id) => firstStale.push({ event, data, id }), { close: () => { closeCalls[0] += 1; offFirst(); } }, ); offSecond = hub.subscribe( "s1", (event, data, id) => secondStale.push({ event, data, id }), { close: () => { closeCalls[1] += 1; offSecond(); } }, ); const beforeId = hub.publish("s1", "info", { text: "before" }); expect(seqOf(beforeId)).toBe(1); hub.clear("s1"); offFirst(); offSecond(); expect(closeCalls).toEqual([1, 1]); 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" }, }); expect(seqOf(gateId)).toBe(3); const resumed: any[] = []; hub.subscribe( "s1", (event, data, id) => resumed.push({ event, data, id }), { after: [beforeId] }, ); 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: afterId }, { event: "ui_request", data: { type: "ui_request", ui_request: { id: "gate-1", widget: "select" } }, id: gateId, }, ]); }); test("forget closes every subscriber once and resets the per-session id sequence", () => { const hub = new SseHub(); const firstStale: any[] = []; const secondStale: any[] = []; const closeCalls = [0, 0]; let offFirst: () => void = () => undefined; let offSecond: () => void = () => undefined; offFirst = hub.subscribe( "s1", (event, data, id) => firstStale.push({ event, data, id }), { close: () => { closeCalls[0] += 1; offFirst(); } }, ); offSecond = hub.subscribe( "s1", (event, data, id) => secondStale.push({ event, data, id }), { close: () => { closeCalls[1] += 1; offSecond(); } }, ); const beforeId = hub.publish("s1", "info", { text: "before" }); expect(seqOf(beforeId)).toBe(1); hub.forget("s1"); expect(closeCalls).toEqual([1, 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: beforeId }]; expect(firstStale).toEqual(before); expect(secondStale).toEqual(before); 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" }; 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 }), { after: [], pending: { ...descriptor, title: "Runtime copy" } }, ); expect(replayed).toEqual([{ event: "ui_request", data: { type: "ui_request", ui_request: descriptor }, id: gateId, }]); const alreadySeen: any[] = []; hub.subscribe( "s1", (event, data, id) => alreadySeen.push({ event, data, id }), { after: [gateId], pending: descriptor }, ); expect(alreadySeen).toEqual([]); }); test("an unbuffered pending gate receives one fresh buffered id", () => { const hub = new SseHub(); const readyId = hub.publish("s1", "info", { text: "ready" }); const first: any[] = []; hub.subscribe( "s1", (event, data, id) => first.push({ event, data, id }), { 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: `${genOf(readyId)}:2`, }]); const reconnect: any[] = []; hub.subscribe( "s1", (event, data, id) => reconnect.push({ event, data, id }), { after: [readyId], pending: { id: "gate-1", widget: "select" } }, ); expect(reconnect).toEqual(first); }); 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[] = []; newHub.subscribe( "s1", (event, data, id) => replayed.push({ event, data, id }), { after: [oldCursor] }, ); expect(replayed).toEqual([ { 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); });