import { test, expect } from "vitest"; import { SseHub } from "../src/sse/sse-hub.js"; 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("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(); } }, ); expect(hub.publish("s1", "info", { text: "before" })).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", { type: "ui_request", ui_request: { id: "gate-1", widget: "select" }, })).toBe(3); const resumed: any[] = []; hub.subscribe( "s1", (event, data, id) => resumed.push({ event, data, id }), { afterId: 1 }, ); const before = [{ event: "info", data: { text: "before" }, id: 1 }]; expect(firstStale).toEqual(before); expect(secondStale).toEqual(before); expect(resumed).toEqual([ { event: "info", data: { text: "after" }, id: 2 }, { event: "ui_request", data: { type: "ui_request", ui_request: { id: "gate-1", widget: "select" } }, id: 3, }, ]); }); 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(); } }, ); expect(hub.publish("s1", "info", { text: "before" })).toBe(1); hub.forget("s1"); expect(closeCalls).toEqual([1, 1]); expect(hub.publish("s1", "info", { text: "after" })).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 }]; expect(firstStale).toEqual(before); expect(secondStale).toEqual(before); expect(fresh).toEqual([{ event: "info", data: { text: "after" }, id: 1 }]); }); 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); }); 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, }, ]); });