Files
EpicNext-Cms/src/app/api/notifications/stream/route.ts
T
openhands 83f4585167
Local Build and Deploy / deploy (push) Successful in 1m5s
Feature: Plugin UI, event hooks, DB notifications, web push, webhook retry, games pages, rate limiting, tests
- Plugin admin UI: enable/disable per plugin with config forms
- Event hooks system: HookManager with on/off/emit for plugin events
- Database notifications: persistent via WebsiteNotification table + Redis pub/sub
- Notification preferences: per-user settings with API endpoints
- Web Push API: browser push notification subscription
- Webhook retry: 3 attempts with exponential backoff (1s/3s/9s) + delivery logs
- Webhook test button + paginated delivery log viewer
- Games: rank/challenge/reward admin pages
- Games: public frontend leaderboard
- Rate limiting: search API (30/min), notification stream dedup
- Unit tests: 28 new tests (registry, notifications, webhooks)
- Prisma migration 0017: website_notifications, website_notification_preferences, website_webhook_logs
2026-07-18 14:04:41 +02:00

108 lines
2.4 KiB
TypeScript

import { auth } from "@/lib/auth";
import { redis } from "@/lib/redis";
export const dynamic = "force-dynamic";
export const runtime = "nodejs";
export async function GET(request: Request) {
const session = await auth();
if (!session?.user?.id) {
return new Response("Unauthorized", { status: 401 });
}
const userId = Number(session.user.id);
if (!userId) {
return new Response("Unauthorized", { status: 401 });
}
// Track active connections per user — close duplicate connections
if (redis) {
try {
const connKey = `stream:conn:${userId}`;
const existing = await redis.get(connKey);
if (existing) {
// Signal the previous connection to close via pub/sub
await redis.publish(
`stream:close:${userId}`,
JSON.stringify({ reason: "duplicate" }),
);
}
await redis.set(connKey, "active", "EX", 60);
} catch {}
}
const stream = new ReadableStream({
start(controller) {
const encoder = new TextEncoder();
let cleanup: (() => void) | null = null;
let closed = false;
const sendEvent = (data: string) => {
if (!closed) {
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
}
};
sendEvent(JSON.stringify({ type: "connected" }));
if (redis) {
const subscriber = redis.duplicate();
// Listen for close signals (duplicate connection)
subscriber.subscribe(`stream:close:${userId}`, () => {
sendEvent(
JSON.stringify({
type: "closed",
message: "Duplicate connection detected",
}),
);
closed = true;
controller.close();
});
subscriber.subscribe(`user:${userId}`, (err) => {
if (err) {
sendEvent(
JSON.stringify({
type: "error",
message: "Subscription failed",
}),
);
}
});
subscriber.on("message", (_channel, message) => {
sendEvent(message);
});
cleanup = () => {
subscriber.unsubscribe();
subscriber.quit();
};
}
const keepAlive = setInterval(() => {
sendEvent(JSON.stringify({ type: "ping" }));
}, 30000);
request.signal.addEventListener("abort", () => {
clearInterval(keepAlive);
cleanup?.();
if (redis) {
try {
redis.del(`stream:conn:${userId}`);
} catch {}
}
});
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
},
});
}