feat(events): stream Redis invalidations over SSE

This commit is contained in:
2026-08-28 20:03:11 +00:00 Unverified
parent b59a038351
commit 455e4240ac
16 changed files with 453 additions and 1 deletions
+44
View File
@@ -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();
});
});
+58
View File
@@ -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<Response> {
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" },
},
);
}
+30
View File
@@ -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();
});
});
+58
View File
@@ -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<typeof setTimeout> | 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;
},
};
}
+47
View File
@@ -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);
});
});
+94
View File
@@ -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<Uint8Array>;
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<Uint8Array> | undefined;
let heartbeat: ReturnType<typeof setInterval> | undefined;
let lifetime: ReturnType<typeof setTimeout> | 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<Uint8Array>({
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);
},
};
}