diff --git a/src/app/api/run/route.ts b/src/app/api/run/route.ts index 4d964ee..3fb1bc7 100644 --- a/src/app/api/run/route.ts +++ b/src/app/api/run/route.ts @@ -7,13 +7,6 @@ import { ClientMessage } from "./types"; import { NextRequest } from "next/server"; import { redis } from "@/lib/redis"; -export function GET(req: NextRequest) { - const headers = new Headers(); - headers.set("Connection", "Upgrade"); - headers.set("Upgrade", "websocket"); - return new Response("Upgrade Required", { status: 426, headers }); -} - export function UPGRADE(client: WebSocket, server: WebSocketServer) { client.once("message", async (raw) => { const { @@ -244,3 +237,10 @@ export function UPGRADE(client: WebSocket, server: WebSocketServer) { // optional cleanup }); } + +export function GET(req: NextRequest) { + const headers = new Headers(); + headers.set("Connection", "Upgrade"); + headers.set("Upgrade", "websocket"); + return new Response("Upgrade Required", { status: 426, headers }); +} diff --git a/src/app/api/view-run/route.ts b/src/app/api/view-run/route.ts index 8686012..7d921b3 100644 --- a/src/app/api/view-run/route.ts +++ b/src/app/api/view-run/route.ts @@ -1,67 +1,63 @@ +// ws/liveRun.ts import { WebSocketServer, WebSocket } from "ws"; -import Redis from "ioredis"; -import { NextRequest, NextResponse } from "next/server"; +import { NextRequest } from "next/server"; import { LiveRunEvent } from "@/modules/live-run/types"; import { updateLiveRunViews } from "@/modules/live-run/core"; -import { redis } from "@/lib/redis"; +import { redis, redisSubscriber } from "@/lib/redis"; -export function GET(req: NextRequest) { - const headers = new Headers(); - headers.set("Connection", "Upgrade"); - headers.set("Upgrade", "websocket"); - return new Response("Upgrade Required", { status: 426, headers }); -} +// Map of runId → Set of WebSockets +const runSubscribers: Map> = new Map(); -export function UPGRADE(client: WebSocket, server: WebSocketServer) { - const subscriber = new Redis(process.env.REDIS_URL!); +// Handle incoming Pub/Sub messages +redisSubscriber.on("message", (channel, message) => { + const runId = channel.split(":")[1]; + const sockets = runSubscribers.get(runId); + if (!sockets) return; + + const data: LiveRunEvent = JSON.parse(message); + + sockets.forEach((ws) => { + if (ws.readyState === WebSocket.OPEN) { + ws.send(JSON.stringify({ ...data, text: null })); + } + }); +}); + +export async function UPGRADE(client: WebSocket, server: WebSocketServer) { let runId: string; + // Receive runId from client client.once("message", async (message) => { runId = message.toString(); + // Track view count await updateLiveRunViews(runId, 1); - let lastEvent: LiveRunEvent | null = null; + // Send initial text + const text = await redis.get(`liveRunText:${runId}`); + client.send( + JSON.stringify({ file: null, language: null, moves: [], text }), + ); - const sendText = async () => { - const text = await redis.get(`liveRunText:${runId}`); - - if (client.readyState === WebSocket.OPEN) { - client.send( - JSON.stringify({ - file: lastEvent?.file, - language: lastEvent?.language, - moves: [], - text, - }), - ); - } - }; - - sendText(); - - client.on("message", sendText); - - subscriber.subscribe(`liveRunEvent:${runId}`, (err, count) => { - if (err) { - console.error("Failed to subscribe", err); - client.close(1011, "Redis subscription failed"); - } - }); - - subscriber.on("message", (channel, message) => { - const data: LiveRunEvent = JSON.parse(message); - - lastEvent = data; - if (client.readyState === WebSocket.OPEN) { - client.send(JSON.stringify({ ...data, text: null })); - } - }); + // Add WebSocket to subscribers + if (!runSubscribers.has(runId)) { + runSubscribers.set(runId, new Set()); + await redisSubscriber.subscribe(`liveRunEvent:${runId}`); + } + runSubscribers.get(runId)!.add(client); }); + // Handle client disconnect client.once("close", async () => { - await subscriber.unsubscribe(`liveRunEvent:${runId}`); - subscriber.quit(); + const set = runSubscribers.get(runId); + if (set) { + set.delete(client); + if (set.size === 0) { + await redisSubscriber.unsubscribe(`liveRunEvent:${runId}`); + runSubscribers.delete(runId); + } + } + await updateLiveRunViews(runId, -1); }); @@ -69,3 +65,10 @@ export function UPGRADE(client: WebSocket, server: WebSocketServer) { console.error("WebSocket error:", err); }); } + +export function GET(req: NextRequest) { + const headers = new Headers(); + headers.set("Connection", "Upgrade"); + headers.set("Upgrade", "websocket"); + return new Response("Upgrade Required", { status: 426, headers }); +} diff --git a/src/lib/redis.ts b/src/lib/redis.ts index f433e79..33321b3 100644 --- a/src/lib/redis.ts +++ b/src/lib/redis.ts @@ -3,6 +3,12 @@ import Redis from "ioredis"; export const redis = new Redis(process.env.REDIS_URL!); +export const redisSubscriber = new Redis(process.env.REDIS_URL!); + redis.on("error", (err) => { console.error("Redis error:", err); }); + +redisSubscriber.on("error", (err) => { + console.error("Redis subscriber error:", err); +});