import { describe, expect, test } from "bun:test"; import { BuildEventStreamHub } from "../../server/build-event-stream"; import type { BuildStatus } from "../../shared/build-protocol"; const base: BuildStatus = { version: 1, id: "a", state: "running", createdAt: "now", }; const tick = () => Bun.sleep(10); describe("build SSE hub", () => { test("reports an explicit gap when the resume cursor predates retained logs", async () => { const abort = new AbortController(); const hub = new BuildEventStreamHub( { getBuildStatus: async () => base, reconcileBuild: async () => base, getBuildEvents: async (_id, after) => [ { type: "status" as const, status: base }, ...[4, 5] .filter((sequence) => sequence > (after ?? 0)) .map((sequence) => ({ type: "log" as const, id: "a", sequence, message: `log ${sequence}\n`, })), ], }, { pollMs: 100, heartbeatMs: 1_000 }, ); const response = await hub.open( "a", new Request("https://test/builds/a/events", { headers: { "last-event-id": "1" }, signal: abort.signal, }), ); const reader = response.body!.getReader(); const decoder = new TextDecoder(); const first = decoder.decode((await reader.read()).value); expect(first).toContain("event: gap\n"); expect(first).toContain('"missing":2'); expect(first).toContain("id: 3\n"); expect(decoder.decode((await reader.read()).value)).toContain("id: 4\nevent: log"); abort.abort(); await reader.cancel(); }); test("coalesces reconciliation, resumes logs, dedupes status and drains terminal", async () => { let calls = 0; let status = base; const hub = new BuildEventStreamHub( { getBuildStatus: async () => status, reconcileBuild: async () => { calls++; return status; }, getBuildEvents: async (_id, after) => [ { type: "status" as const, status: base }, ...[1, 2] .filter((sequence) => sequence > (after ?? 0)) .map((sequence) => ({ type: "log" as const, id: "a", sequence, message: `log ${sequence}\n`, })), ], }, { pollMs: 20, heartbeatMs: 40 }, ); const a = new AbortController(); const b = new AbortController(); const request = (signal: AbortSignal, after: string) => new Request(`https://test/builds/a/events?after=${after}`, { signal }); const first = await hub.open("a", request(a.signal, "1")); const second = await hub.open("a", request(b.signal, "2")); expect(first.headers.get("content-type")).toContain("text/event-stream"); const reader1 = first.body!.getReader(); const reader2 = second.body!.getReader(); const text = new TextDecoder(); expect(text.decode((await reader1.read()).value)).toContain( "id: 2\nevent: log\n", ); expect(text.decode((await reader1.read()).value)).toContain( "event: status\n", ); expect(text.decode((await reader2.read()).value)).toContain( "event: status\n", ); await tick(); expect(calls).toBe(1); status = { ...base, state: "succeeded" }; const remaining = async (reader: typeof reader1) => { let output = ""; for (;;) { const { value, done } = await reader.read(); if (done) return output; output += text.decode(value); } }; expect(await remaining(reader1)).toContain('"state":"succeeded"'); expect(await remaining(reader2)).toContain('"state":"succeeded"'); a.abort(); b.abort(); }); test("validates cursor, stops aborted subscribers and emits heartbeat comments", async () => { const hub = new BuildEventStreamHub( { getBuildStatus: async () => base, reconcileBuild: async () => base, getBuildEvents: async () => [], }, { pollMs: 5, heartbeatMs: 5 }, ); await expect( hub.open("a", new Request("https://test/events?after=-1")), ).rejects.toThrow(RangeError); const abort = new AbortController(); const response = await hub.open( "a", new Request("https://test/events", { signal: abort.signal }), ); const reader = response.body!.getReader(); expect(new TextDecoder().decode((await reader.read()).value)).toContain( "event: status", ); expect(new TextDecoder().decode((await reader.read()).value)).toBe( ": heartbeat\n\n", ); abort.abort(); expect((await reader.read()).done).toBe(true); }); test("reads persisted progress when another replica owns reconciliation", async () => { let status = base; let sequence = 0; let reconciles = 0; const hub = new BuildEventStreamHub( { getBuildStatus: async () => status, reconcileBuild: async () => { reconciles++; throw new Error("build lease held by another replica"); }, getBuildEvents: async (_id, after) => sequence > (after ?? 0) ? [ { type: "log" as const, id: "a", sequence, message: "remote log\n", }, ] : [], }, { pollMs: 5, heartbeatMs: 50 }, ); const abort = new AbortController(); const response = await hub.open( "a", new Request("https://test/events", { signal: abort.signal }), ); const reader = response.body!.getReader(); const decoder = new TextDecoder(); expect(decoder.decode((await reader.read()).value)).toContain( '"state":"running"', ); sequence = 1; expect(decoder.decode((await reader.read()).value)).toContain( "id: 1\nevent: log", ); status = { ...base, state: "succeeded" }; let terminalFrame = ""; for (;;) { const { value, done } = await reader.read(); if (done) break; terminalFrame += decoder.decode(value); } expect(terminalFrame).toContain('"state":"succeeded"'); expect(reconciles).toBeGreaterThan(0); abort.abort(); }); test("heartbeats during slow reconciliation and tears down on abort", async () => { const blocked = new Promise(() => {}); const signals: AbortSignal[] = []; let observations = 0; const hub = new BuildEventStreamHub( { getBuildStatus: async () => base, reconcileBuild: async (_id, options) => { if (options?.signal) signals.push(options.signal); return blocked; }, getBuildEvents: async () => { observations++; return []; }, }, { pollMs: 5, heartbeatMs: 5 }, ); const abort = new AbortController(); const response = await hub.open( "a", new Request("https://test/events", { signal: abort.signal }), ); const reader = response.body!.getReader(); expect(new TextDecoder().decode((await reader.read()).value)).toContain( "event: status", ); expect(new TextDecoder().decode((await reader.read()).value)).toBe( ": heartbeat\n\n", ); await Bun.sleep(25); expect(observations).toBeGreaterThan(1); expect(signals).toHaveLength(1); abort.abort(); for (;;) { if ((await reader.read()).done) break; } expect(signals[0]?.aborted).toBe(true); const stoppedAt = observations; await Bun.sleep(20); expect(observations).toBe(stoppedAt); const next = new AbortController(); const reopened = await hub.open( "a", new Request("https://test/events", { signal: next.signal }), ); const nextReader = reopened.body!.getReader(); expect(new TextDecoder().decode((await nextReader.read()).value)).toContain( "event: status", ); expect(signals).toHaveLength(2); expect(signals[1]?.aborted).toBe(false); next.abort(); await nextReader.cancel(); }); test("bounds shared reconciliation writes while observing two subscribers independently", async () => { let reconciles = 0; let observations = 0; const hub = new BuildEventStreamHub( { getBuildStatus: async () => base, reconcileBuild: async () => { reconciles++; return base; }, getBuildEvents: async () => { observations++; return []; }, }, { pollMs: 10, reconcileMs: 80, heartbeatMs: 1_000 }, ); const a = new AbortController(); const b = new AbortController(); const first = await hub.open( "a", new Request("https://test/events", { signal: a.signal }), ); const second = await hub.open( "a", new Request("https://test/events", { signal: b.signal }), ); expect( new TextDecoder().decode((await first.body!.getReader().read()).value), ).toContain("event: status"); expect( new TextDecoder().decode((await second.body!.getReader().read()).value), ).toContain("event: status"); await Bun.sleep(45); expect(reconciles).toBe(1); // One lease/write attempt, not one per poll or subscriber. expect(observations).toBeGreaterThanOrEqual(3); for (let i = 0; i < 15 && reconciles < 2; i++) await tick(); expect(reconciles).toBe(2); a.abort(); b.abort(); }); test("closes a slow subscriber at the bounded queue limit for resumable logs", async () => { const hub = new BuildEventStreamHub( { getBuildStatus: async () => base, reconcileBuild: async () => base, getBuildEvents: async (_id, after) => Array.from({ length: 40 }, (_, index) => index + 1) .filter((sequence) => sequence > (after ?? 0)) .map((sequence) => ({ type: "log" as const, id: "a", sequence, message: `log ${sequence}\n`, })), }, { pollMs: 5, heartbeatMs: 50 }, ); const response = await hub.open("a", new Request("https://test/events")); await tick(); // Allow the producer to fill its queue before consumption. const reader = response.body!.getReader(); const frames: string[] = []; for (;;) { const { value, done } = await reader.read(); if (done) break; frames.push(new TextDecoder().decode(value)); } expect(frames).toHaveLength(16); expect(frames[0]).toContain("id: 1\nevent: log"); expect(frames[15]).toContain("id: 16\nevent: log"); const resumed = await hub.open( "a", new Request("https://test/events", { headers: { "last-event-id": "16" }, }), ); const resumedReader = resumed.body!.getReader(); expect( new TextDecoder().decode((await resumedReader.read()).value), ).toContain("id: 17\nevent: log"); await resumedReader.cancel(); }); });