feat(backend): SSE hub with pending-widget re-emit on (re)subscribe

This commit is contained in:
2026-06-27 21:10:48 +02:00
parent f4c6126270
commit 39eb35116e
2 changed files with 29 additions and 0 deletions
+13
View File
@@ -0,0 +1,13 @@
type Send = (event: string, data: object) => void;
export class SseHub {
private subs = new Map<string, Set<Send>>();
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);
}
}
+16
View File
@@ -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" } }]);
});