feat(phase3c): SSE /api/events/stream + live /events page
This commit is contained in:
73
apps/web/src/app/api/events/stream/route.ts
Normal file
73
apps/web/src/app/api/events/stream/route.ts
Normal file
@@ -0,0 +1,73 @@
|
||||
import { headers } from "next/headers";
|
||||
import { auth } from "@/lib/auth";
|
||||
import { prisma } from "@/lib/db";
|
||||
|
||||
export const dynamic = "force-dynamic";
|
||||
export const runtime = "nodejs";
|
||||
|
||||
export async function GET() {
|
||||
const session = await auth.api.getSession({ headers: await headers() });
|
||||
if (!session) return new Response("unauthorized", { status: 401 });
|
||||
|
||||
const encoder = new TextEncoder();
|
||||
let cursor = new Date(Date.now() - 60 * 1000);
|
||||
let alive = true;
|
||||
|
||||
const stream = new ReadableStream({
|
||||
async start(controller) {
|
||||
const send = (event: string, data: unknown) => {
|
||||
controller.enqueue(encoder.encode(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`));
|
||||
};
|
||||
send("hello", { ts: Date.now() });
|
||||
|
||||
const tick = async () => {
|
||||
if (!alive) return;
|
||||
try {
|
||||
const rows = await prisma.event.findMany({
|
||||
where: { receivedAt: { gt: cursor } },
|
||||
orderBy: { receivedAt: "asc" },
|
||||
take: 50,
|
||||
});
|
||||
for (const r of rows) {
|
||||
send("event", {
|
||||
id: r.id,
|
||||
streamId: r.streamId,
|
||||
projectKey: r.projectKey,
|
||||
eventType: r.eventType,
|
||||
version: r.version,
|
||||
occurredAt: r.occurredAt.toISOString(),
|
||||
receivedAt: r.receivedAt.toISOString(),
|
||||
payload: r.payload,
|
||||
});
|
||||
cursor = r.receivedAt;
|
||||
}
|
||||
} catch (e) {
|
||||
send("error", { message: e instanceof Error ? e.message : "tick failed" });
|
||||
}
|
||||
};
|
||||
|
||||
await tick();
|
||||
const interval = setInterval(tick, 2000);
|
||||
const heartbeat = setInterval(() => send("ping", { ts: Date.now() }), 25_000);
|
||||
|
||||
const close = () => {
|
||||
alive = false;
|
||||
clearInterval(interval);
|
||||
clearInterval(heartbeat);
|
||||
try { controller.close(); } catch {}
|
||||
};
|
||||
// close after 5 min — client will reconnect via EventSource
|
||||
setTimeout(close, 5 * 60 * 1000);
|
||||
},
|
||||
cancel() { alive = false; },
|
||||
});
|
||||
|
||||
return new Response(stream, {
|
||||
headers: {
|
||||
"Content-Type": "text/event-stream; charset=utf-8",
|
||||
"Cache-Control": "no-cache, no-transform",
|
||||
"Connection": "keep-alive",
|
||||
"X-Accel-Buffering": "no",
|
||||
},
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user