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"; const seqOf = (id: string) => Number(id.split(":")[1]); 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?: (ids: string[]) => string; lastEventId?: (ids: string[]) => string }, ): Promise<{ body: string; ids: string[] }> { const hub = new SseHub(); const ids = [ 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 lastEventId = options.lastEventId?.(ids); const query = options.query ? `?lastEventId=${encodeURIComponent(options.query(ids))}` : ""; const response = await fetch( `http://127.0.0.1:${port}/sessions/s1/events${query}`, { headers: lastEventId ? { "Last-Event-ID": lastEventId } : undefined, signal: controller.signal, }, ); const reader = response.body!.getReader(); const body = await readUntil(reader, (text) => text.includes('"three"')); await reader.cancel(); return { body, ids }; } 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, ids } = await captureReplay({ lastEventId: (published) => published[0] }); expect(body).not.toContain('"one"'); expect(body).toContain(`id: ${ids[1]}\nevent: info\n`); expect(body).toContain('data: {"type":"info","text":"two"}'); expect(body).toContain(`id: ${ids[2]}\nevent: info\n`); expect(body).toContain('data: {"type":"info","text":"three"}'); }); test("SSE honors the manual lastEventId query cursor", async () => { const { body, ids } = await captureReplay({ query: (published) => published[1] }); expect(body).not.toContain('"one"'); expect(body).not.toContain('"two"'); expect(body).toContain(`id: ${ids[2]}\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: (published) => published[1], lastEventId: (published) => published[0], }); expect(body).not.toContain('"two"'); expect(body).toContain('"three"'); }); test("SSE resets a cursor from an older backend generation and emits a buffered pending gate once", async () => { // The cursor comes from a PREVIOUS process: same session, ids restarted. It must be // treated as stale (replay from the beginning), not honored against the new ids. const oldGeneration = new SseHub(); oldGeneration.publish("s1", "info", { type: "info", text: "old" }); const staleCursor = oldGeneration.publish("s1", "info", { type: "info", text: "older" }); const hub = new SseHub(); const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" }; const infoId = hub.publish("s1", "info", { type: "info", text: "fresh generation" }); const gateId = 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": staleCursor }, 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: ${infoId}\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: ${gateId}\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()); const beforeId = hub.publish("s1", "info", { type: "info", text: "before" }); expect(seqOf(beforeId)).toBe(1); const initial = await Promise.all(readers.map((reader) => readUntil(reader, (text) => text.includes('"before"')))); expect(initial.every((body) => body.includes(`id: ${beforeId}\nevent: info\n`))).toBe(true); hub.clear("s1"); await Promise.all(readers.map(expectStreamEnd)); const afterId = hub.publish("s1", "info", { type: "info", text: "after" }); expect(seqOf(afterId)).toBe(2); const gateId = hub.publish("s1", "ui_request", { type: "ui_request", ui_request: { id: "gate-1", widget: "select" }, }); expect(seqOf(gateId)).toBe(3); const reconnected = await fetch( `http://127.0.0.1:${port}/sessions/s1/events`, { headers: { "Last-Event-ID": beforeId }, signal: controllers[2].signal, }, ); const reconnectReader = reconnected.body!.getReader(); const replay = await readUntil(reconnectReader, (text) => text.includes('"gate-1"')); expect(replay).toContain(`id: ${afterId}\nevent: info\n`); expect(replay).toContain('data: {"type":"info","text":"after"}'); expect(replay).toContain(`id: ${gateId}\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(); } });