Files
ThothII/backend/test/sse-hub.test.ts
T

181 lines
5.7 KiB
TypeScript

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,
},
]);
});