51 lines
1.2 KiB
TypeScript
51 lines
1.2 KiB
TypeScript
import { NextRequest } from "next/server";
|
|
import { type SSEEventClient, subscribe, unsubscribe } from "@/lib/realtime/bus";
|
|
|
|
export const dynamic = "force-dynamic";
|
|
export const runtime = "nodejs";
|
|
|
|
export async function GET(req: NextRequest) {
|
|
const encoder = new TextEncoder();
|
|
|
|
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)),
|
|
};
|
|
subscribe(client);
|
|
|
|
controller.enqueue(encoder.encode(": connected\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",
|
|
},
|
|
});
|
|
}
|