1
0
Fork 0
polaris-task-force/src/app/api/realtime/route.ts

62 lines
1.9 KiB
TypeScript

import { NextRequest } from "next/server";
import config from "@payload-config";
import { getPayload } from "payload";
import { type SSEEventClient, subscribe, unsubscribe } from "@/lib/realtime/bus";
import { getPresence } from "@/lib/realtime/presence";
export const dynamic = "force-dynamic";
export const runtime = "nodejs";
export async function GET(req: NextRequest) {
const payload = await getPayload({ config: await config });
const { user } = await payload.auth({ headers: req.headers, canSetHeaders: false });
const encoder = new TextEncoder();
const channel = req.nextUrl.searchParams.get("channel");
let client: SSEEventClient | null = null;
let heartbeat: ReturnType<typeof setInterval> | null = null;
const cleanup = () => {
if (heartbeat) clearInterval(heartbeat);
if (client) unsubscribe(client);
client = null;
heartbeat = null;
};
const stream = new ReadableStream({
start(controller) {
client = {
send: (data) => controller.enqueue(encoder.encode(data)),
userId: user ? Number(user.id) : undefined,
events: channel ? new Set(["presence"]) : undefined,
};
subscribe(client);
controller.enqueue(encoder.encode(": connected\n\n"));
if (channel) {
controller.enqueue(encoder.encode(`event: presence\ndata: ${JSON.stringify({ channel, action: "snapshot", participants: getPresence(channel) })}\n\n`));
}
heartbeat = setInterval(() => {
try {
controller.enqueue(encoder.encode(": keepalive\n\n"));
} catch {
cleanup();
}
}, 25_000);
req.signal.addEventListener("abort", cleanup);
},
cancel() {
cleanup();
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
},
});
}