import { test, expect } from "vitest"; import { buildApp } from "../src/app.js"; import { loadConfig } from "../src/config.js"; import { SseHub } from "../src/sse/sse-hub.js"; async function readUntil( reader: ReadableStreamDefaultReader, predicate: (text: string) => boolean, ): Promise { const decoder = new TextDecoder(); let text = ""; const deadline = Date.now() + 1_000; while (Date.now() < deadline) { const result = await Promise.race([ reader.read(), new Promise<{ timeout: true }>((resolve) => setTimeout(() => resolve({ timeout: true }), Math.max(0, deadline - Date.now())), ), ]); if ("timeout" in result || result.done) break; text += decoder.decode(result.value); if (predicate(text)) break; } return text; } async function captureReplay(options: { query?: string; lastEventId?: string }): Promise { const hub = new SseHub(); hub.publish("s1", "info", { type: "info", text: "one" }); hub.publish("s1", "info", { type: "info", text: "two" }); hub.publish("s1", "info", { type: "info", text: "three" }); const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), { hub, thtRunner: {} as any, }); await app.listen({ port: 0, host: "127.0.0.1" }); const port = (app.server.address() as { port: number }).port; const controller = new AbortController(); try { const response = await fetch( `http://127.0.0.1:${port}/sessions/s1/events${options.query ?? ""}`, { headers: options.lastEventId ? { "Last-Event-ID": options.lastEventId } : undefined, signal: controller.signal, }, ); const reader = response.body!.getReader(); const body = await readUntil(reader, (text) => text.includes('"three"')); await reader.cancel(); return body; } finally { controller.abort(); await app.close(); } } async function expectStreamEnd(reader: ReadableStreamDefaultReader): Promise { const result = await Promise.race([ reader.read(), new Promise<{ timeout: true }>((resolve) => setTimeout(() => resolve({ timeout: true }), 1_000), ), ]); expect(result).not.toEqual({ timeout: true }); expect("done" in result && result.done).toBe(true); } test("SSE emits ids and honors the native Last-Event-ID replay cursor", async () => { const body = await captureReplay({ lastEventId: "1" }); expect(body).not.toContain('"one"'); expect(body).toContain("id: 2\nevent: info\n"); expect(body).toContain('data: {"type":"info","text":"two"}'); expect(body).toContain("id: 3\nevent: info\n"); expect(body).toContain('data: {"type":"info","text":"three"}'); }); test("SSE honors the manual lastEventId query cursor", async () => { const body = await captureReplay({ query: "?lastEventId=2" }); expect(body).not.toContain('"one"'); expect(body).not.toContain('"two"'); expect(body).toContain("id: 3\nevent: info\n"); expect(body).toContain('data: {"type":"info","text":"three"}'); }); test("SSE uses the newer valid cursor when header and query are both present", async () => { const body = await captureReplay({ query: "?lastEventId=2", lastEventId: "1" }); expect(body).not.toContain('"two"'); expect(body).toContain("id: 3\nevent: info\n"); }); test("SSE resets a cursor from an older process and emits a buffered pending gate once", async () => { 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 app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), { hub, thtRunner: {} as any, mgr: { get: () => ({ bridge: { pendingWidget: () => ({ ...descriptor }) } }), } as any, }); await app.listen({ port: 0, host: "127.0.0.1" }); const port = (app.server.address() as { port: number }).port; const controller = new AbortController(); try { const response = await fetch(`http://127.0.0.1:${port}/sessions/s1/events`, { headers: { "Last-Event-ID": "900" }, signal: controller.signal, }); const reader = response.body!.getReader(); const body = await readUntil(reader, (text) => text.includes('"gate-fresh"')); await reader.cancel(); expect(body).toContain("id: 1\nevent: info\n"); expect(body).toContain('data: {"type":"info","text":"fresh generation"}'); expect(body.match(/event: ui_request/g)).toHaveLength(1); expect(body).toContain("id: 2\nevent: ui_request\n"); expect(body).toContain('"ui_request":{"id":"gate-fresh","widget":"select","title":"Choose"}'); } finally { controller.abort(); await app.close(); } }); test("clear ends every live SSE response and a cursor reconnect replays post-clear events", async () => { const hub = new SseHub(); const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), { hub, thtRunner: {} as any, }); await app.listen({ port: 0, host: "127.0.0.1" }); const port = (app.server.address() as { port: number }).port; const controllers = [new AbortController(), new AbortController(), new AbortController()]; try { const responses = await Promise.all(controllers.slice(0, 2).map((controller) => fetch( `http://127.0.0.1:${port}/sessions/s1/events`, { signal: controller.signal }, ))); const readers = responses.map((response) => response.body!.getReader()); expect(hub.publish("s1", "info", { type: "info", text: "before" })).toBe(1); const initial = await Promise.all(readers.map((reader) => readUntil(reader, (text) => text.includes('"before"')))); expect(initial.every((body) => body.includes("id: 1\nevent: info\n"))).toBe(true); hub.clear("s1"); await Promise.all(readers.map(expectStreamEnd)); expect(hub.publish("s1", "info", { type: "info", text: "after" })).toBe(2); expect(hub.publish("s1", "ui_request", { type: "ui_request", ui_request: { id: "gate-1", widget: "select" }, })).toBe(3); const reconnected = await fetch( `http://127.0.0.1:${port}/sessions/s1/events`, { headers: { "Last-Event-ID": "1" }, signal: controllers[2].signal, }, ); const reconnectReader = reconnected.body!.getReader(); const replay = await readUntil(reconnectReader, (text) => text.includes('"gate-1"')); expect(replay).toContain("id: 2\nevent: info\n"); expect(replay).toContain('data: {"type":"info","text":"after"}'); expect(replay).toContain("id: 3\nevent: ui_request\n"); expect(replay).toContain('"ui_request":{"id":"gate-1","widget":"select"}'); await reconnectReader.cancel(); } finally { for (const controller of controllers) controller.abort(); await app.close(); } });