import { afterEach, describe, expect, it, vi } from "vitest"; import { createEventStream, eventStreamHeaders, SSE_HEARTBEAT_MS, SSE_MAX_LIFETIME_MS, } from "./sse"; const decoder = new TextDecoder(); afterEach(() => { vi.useRealTimers(); }); describe("SSE stream", () => { it("uses the required heartbeat, reconnect lifetime, and proxy headers", () => { expect(SSE_HEARTBEAT_MS).toBe(90_000); expect(SSE_MAX_LIFETIME_MS).toBe(30 * 60 * 1_000); const headers = eventStreamHeaders(); expect(headers.get("content-type")).toContain("text/event-stream"); expect(headers.get("cache-control")).toBe("no-cache, no-transform"); expect(headers.get("x-accel-buffering")).toBe("no"); }); it("heartbeats, sends only the supplied envelope, and closes on schedule", async () => { vi.useFakeTimers(); const onClose = vi.fn(); const channel = createEventStream({ heartbeatMs: 50, maximumLifetimeMs: 99, onClose, }); const reader = channel.stream.getReader(); expect(decoder.decode((await reader.read()).value)).toBe(": connected\n\n"); channel.send({ type: "page.updated", id: "opaque", version: 7 }); expect(decoder.decode((await reader.read()).value)).toBe( 'event: invalidation\ndata: {"type":"page.updated","id":"opaque","version":7}\n\n', ); channel.sendNamed("deployment", "release-2"); expect(decoder.decode((await reader.read()).value)).toBe( "event: deployment\ndata: release-2\n\n", ); vi.advanceTimersByTime(50); expect(decoder.decode((await reader.read()).value)).toMatch(/^: heartbeat /u); vi.advanceTimersByTime(49); expect((await reader.read()).done).toBe(true); expect(onClose).toHaveBeenCalledTimes(1); }); it("rejects event names that could inject another SSE field", () => { const channel = createEventStream(); expect(() => channel.sendNamed("deployment\ndata", "release-2")).toThrow( "SSE event names", ); channel.close(); }); });