diff --git a/lib/events/invalidation.test.ts b/lib/events/invalidation.test.ts
new file mode 100644
index 0000000..da05cea
--- /dev/null
+++ b/lib/events/invalidation.test.ts
@@ -0,0 +1,44 @@
+import { describe, expect, it } from "vitest";
+
+import {
+ parseInvalidationEvent,
+ toInvalidationEvent,
+} from "./invalidation";
+
+describe("invalidation event envelopes", () => {
+ it("strips every field except type, opaque id, and version", () => {
+ const event = toInvalidationEvent({
+ topic: "directory",
+ aggregateId: "4afeea7b-6f24-43f3-b747-2db98335e01e",
+ eventType: "directory.updated",
+ payload: {
+ id: "4afeea7b-6f24-43f3-b747-2db98335e01e",
+ version: 3,
+ title: "private draft title",
+ blocks: [{ content: "private" }],
+ },
+ });
+ expect(event).toEqual({
+ type: "directory.updated",
+ id: "4afeea7b-6f24-43f3-b747-2db98335e01e",
+ version: 3,
+ });
+ });
+
+ it("rejects unknown, malformed, and topic-mismatched messages", () => {
+ expect(parseInvalidationEvent("not json")).toBeNull();
+ expect(
+ parseInvalidationEvent(
+ JSON.stringify({ type: "page.content", id: "id", version: 1 }),
+ ),
+ ).toBeNull();
+ expect(
+ toInvalidationEvent({
+ topic: "admin",
+ aggregateId: "id",
+ eventType: "page.updated",
+ payload: { id: "id", version: 1 },
+ }),
+ ).toBeNull();
+ });
+});
diff --git a/lib/events/redis-stream.ts b/lib/events/redis-stream.ts
new file mode 100644
index 0000000..0ad0788
--- /dev/null
+++ b/lib/events/redis-stream.ts
@@ -0,0 +1,58 @@
+import {
+ expectedTopic,
+ parseInvalidationEvent,
+} from "@/lib/events/invalidation";
+import {
+ createEventStream,
+ eventStreamHeaders,
+ type EventStreamChannel,
+} from "@/lib/events/sse";
+import {
+ connectRedisClient,
+ createRedisClient,
+ redisEventChannel,
+} from "@/lib/redis/client";
+
+export async function createRedisEventResponse(
+ topic: string,
+ signal: AbortSignal,
+): Promise {
+ const redis = createRedisClient(`buzz-sheet-sse:${topic}`);
+ const channel = redisEventChannel(topic);
+ try {
+ await connectRedisClient(redis);
+ if (signal.aborted) throw new Error("Request closed before subscription.");
+ await redis.subscribe(channel);
+ } catch (cause) {
+ redis.disconnect();
+ throw cause;
+ }
+
+ const onMessage = (receivedChannel: string, value: string) => {
+ if (receivedChannel !== channel) return;
+ const event = parseInvalidationEvent(value);
+ if (event && expectedTopic(event) === topic) eventStream.send(event);
+ };
+ const eventStream: EventStreamChannel = createEventStream({
+ signal,
+ onClose: () => {
+ redis.off("message", onMessage);
+ redis.disconnect();
+ },
+ });
+ redis.on("message", onMessage);
+
+ return new Response(eventStream.stream, {
+ headers: eventStreamHeaders(),
+ });
+}
+
+export function eventsUnavailableResponse(): Response {
+ return Response.json(
+ { error: "events-unavailable" },
+ {
+ status: 503,
+ headers: { "Cache-Control": "no-store" },
+ },
+ );
+}
diff --git a/lib/events/refresh.test.ts b/lib/events/refresh.test.ts
new file mode 100644
index 0000000..fbc3540
--- /dev/null
+++ b/lib/events/refresh.test.ts
@@ -0,0 +1,30 @@
+import { afterEach, describe, expect, it, vi } from "vitest";
+
+import { createRefreshScheduler } from "./refresh";
+
+afterEach(() => {
+ vi.useRealTimers();
+});
+
+describe("authoritative refresh scheduler", () => {
+ it("refetches on connection and coalesces duplicate notifications", async () => {
+ vi.useFakeTimers();
+ const refresh = vi.fn();
+ const scheduler = createRefreshScheduler(refresh, { delayMs: 25 });
+ const event = { type: "page.updated" as const, id: "page", version: 5 };
+
+ scheduler.connected();
+ await vi.advanceTimersByTimeAsync(25);
+ expect(refresh).toHaveBeenCalledTimes(1);
+
+ scheduler.notified(event);
+ scheduler.notified(event);
+ await vi.advanceTimersByTimeAsync(25);
+ expect(refresh).toHaveBeenCalledTimes(2);
+
+ scheduler.connected();
+ await vi.advanceTimersByTimeAsync(25);
+ expect(refresh).toHaveBeenCalledTimes(3);
+ scheduler.dispose();
+ });
+});
diff --git a/lib/events/refresh.ts b/lib/events/refresh.ts
new file mode 100644
index 0000000..19a5f0c
--- /dev/null
+++ b/lib/events/refresh.ts
@@ -0,0 +1,58 @@
+import type { InvalidationEvent } from "./invalidation";
+
+export interface RefreshScheduler {
+ connected(): void;
+ notified(event: InvalidationEvent): void;
+ dispose(): void;
+}
+
+export interface RefreshSchedulerOptions {
+ delayMs?: number;
+ duplicateWindowMs?: number;
+ now?: () => number;
+}
+
+export function createRefreshScheduler(
+ refresh: () => void,
+ options: RefreshSchedulerOptions = {},
+): RefreshScheduler {
+ const delayMs = options.delayMs ?? 75;
+ const duplicateWindowMs = options.duplicateWindowMs ?? 2_000;
+ const now = options.now ?? Date.now;
+ let pending: ReturnType | undefined;
+ let lastEventKey: string | undefined;
+ let lastEventAt = 0;
+
+ const schedule = (event?: InvalidationEvent) => {
+ if (event) {
+ const key = `${event.type}:${event.id}:${event.version}`;
+ const timestamp = now();
+ if (
+ key === lastEventKey &&
+ timestamp - lastEventAt < duplicateWindowMs
+ ) {
+ return;
+ }
+ lastEventKey = key;
+ lastEventAt = timestamp;
+ }
+ if (pending) return;
+ pending = setTimeout(() => {
+ pending = undefined;
+ refresh();
+ }, delayMs);
+ };
+
+ return {
+ connected() {
+ schedule();
+ },
+ notified(event) {
+ schedule(event);
+ },
+ dispose() {
+ if (pending) clearTimeout(pending);
+ pending = undefined;
+ },
+ };
+}
diff --git a/lib/events/sse.test.ts b/lib/events/sse.test.ts
new file mode 100644
index 0000000..6a4099b
--- /dev/null
+++ b/lib/events/sse.test.ts
@@ -0,0 +1,47 @@
+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',
+ );
+ await vi.advanceTimersByTimeAsync(50);
+ expect(decoder.decode((await reader.read()).value)).toMatch(/^: heartbeat /u);
+ await vi.advanceTimersByTimeAsync(49);
+ expect((await reader.read()).done).toBe(true);
+ expect(onClose).toHaveBeenCalledTimes(1);
+ });
+});
diff --git a/lib/events/sse.ts b/lib/events/sse.ts
new file mode 100644
index 0000000..b23491c
--- /dev/null
+++ b/lib/events/sse.ts
@@ -0,0 +1,94 @@
+import type { InvalidationEvent } from "./invalidation";
+
+export const SSE_HEARTBEAT_MS = 90_000;
+export const SSE_MAX_LIFETIME_MS = 30 * 60 * 1_000;
+
+const encoder = new TextEncoder();
+
+export interface EventStreamOptions {
+ signal?: AbortSignal;
+ heartbeatMs?: number;
+ maximumLifetimeMs?: number;
+ onClose?: () => void;
+}
+
+export interface EventStreamChannel {
+ stream: ReadableStream;
+ send(event: InvalidationEvent): boolean;
+ close(): void;
+}
+
+export function eventStreamHeaders(): Headers {
+ return new Headers({
+ "Cache-Control": "no-cache, no-transform",
+ Connection: "keep-alive",
+ "Content-Type": "text/event-stream; charset=utf-8",
+ "X-Accel-Buffering": "no",
+ });
+}
+
+export function createEventStream(
+ options: EventStreamOptions = {},
+): EventStreamChannel {
+ const heartbeatMs = options.heartbeatMs ?? SSE_HEARTBEAT_MS;
+ const maximumLifetimeMs =
+ options.maximumLifetimeMs ?? SSE_MAX_LIFETIME_MS;
+ let controller: ReadableStreamDefaultController | undefined;
+ let heartbeat: ReturnType | undefined;
+ let lifetime: ReturnType | undefined;
+ let closed = false;
+
+ const enqueue = (value: string): boolean => {
+ if (closed || !controller) return false;
+ try {
+ controller.enqueue(encoder.encode(value));
+ return true;
+ } catch {
+ close(false);
+ return false;
+ }
+ };
+
+ const abort = () => close(true);
+ const close = (closeController: boolean) => {
+ if (closed) return;
+ closed = true;
+ if (heartbeat) clearInterval(heartbeat);
+ if (lifetime) clearTimeout(lifetime);
+ options.signal?.removeEventListener("abort", abort);
+ if (closeController) {
+ try {
+ controller?.close();
+ } catch {
+ // The browser may have cancelled the stream before cleanup ran.
+ }
+ }
+ options.onClose?.();
+ };
+
+ const stream = new ReadableStream({
+ start(streamController) {
+ controller = streamController;
+ enqueue(": connected\n\n");
+ heartbeat = setInterval(() => {
+ enqueue(`: heartbeat ${Date.now()}\n\n`);
+ }, heartbeatMs);
+ lifetime = setTimeout(() => close(true), maximumLifetimeMs);
+ if (options.signal?.aborted) close(true);
+ else options.signal?.addEventListener("abort", abort, { once: true });
+ },
+ cancel() {
+ close(false);
+ },
+ });
+
+ return {
+ stream,
+ send(event) {
+ return enqueue(`event: invalidation\ndata: ${JSON.stringify(event)}\n\n`);
+ },
+ close() {
+ close(true);
+ },
+ };
+}
diff --git a/tests/security-boundaries.test.ts b/tests/security-boundaries.test.ts
index 6ec89df..916bd05 100644
--- a/tests/security-boundaries.test.ts
+++ b/tests/security-boundaries.test.ts
@@ -5,6 +5,7 @@ const protectedEntryPoints = [
"app/api/admin/media/presign/route.ts",
"app/api/admin/media/[id]/complete/route.ts",
"app/api/admin/media/[id]/route.ts",
+ "app/api/events/admin/route.ts",
] as const;
describe("administrative security boundaries", () => {