diff --git a/backend/src/sse/sse-hub.ts b/backend/src/sse/sse-hub.ts new file mode 100644 index 00000000..0ccd9813 --- /dev/null +++ b/backend/src/sse/sse-hub.ts @@ -0,0 +1,13 @@ +type Send = (event: string, data: object) => void; +export class SseHub { + private subs = new Map>(); + subscribe(sessionId: string, send: Send, pending?: object | null): () => void { + if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set()); + this.subs.get(sessionId)!.add(send); + if (pending) send("ui_request", { ui_request: pending }); + return () => this.subs.get(sessionId)?.delete(send); + } + publish(sessionId: string, event: string, data: object): void { + for (const s of this.subs.get(sessionId) ?? []) s(event, data); + } +} diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts new file mode 100644 index 00000000..5c1501d8 --- /dev/null +++ b/backend/test/sse-hub.test.ts @@ -0,0 +1,16 @@ +import { test, expect } from "vitest"; +import { SseHub } from "../src/sse/sse-hub.js"; + +test("re-emette il widget pendente alla sottoscrizione", () => { + const hub = new SseHub(); const sent: any[] = []; + hub.subscribe("s1", (ev, data) => sent.push({ ev, data }), { id: "u1", widget: "select" }); + expect(sent[0]).toEqual({ ev: "ui_request", data: { ui_request: { id: "u1", widget: "select" } } }); +}); + +test("publish raggiunge i subscriber e unsubscribe li stacca", () => { + const hub = new SseHub(); const sent: any[] = []; + const off = hub.subscribe("s1", (ev, data) => sent.push({ ev, data })); + hub.publish("s1", "text_delta", { text: "x" }); + off(); hub.publish("s1", "text_delta", { text: "y" }); + expect(sent).toEqual([{ ev: "text_delta", data: { text: "x" } }]); +});