Max redis clients fix
This commit is contained in:
@@ -7,13 +7,6 @@ import { ClientMessage } from "./types";
|
|||||||
import { NextRequest } from "next/server";
|
import { NextRequest } from "next/server";
|
||||||
import { redis } from "@/lib/redis";
|
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) {
|
export function UPGRADE(client: WebSocket, server: WebSocketServer) {
|
||||||
client.once("message", async (raw) => {
|
client.once("message", async (raw) => {
|
||||||
const {
|
const {
|
||||||
@@ -244,3 +237,10 @@ export function UPGRADE(client: WebSocket, server: WebSocketServer) {
|
|||||||
// optional cleanup
|
// 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 });
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,67 +1,63 @@
|
|||||||
|
// ws/liveRun.ts
|
||||||
import { WebSocketServer, WebSocket } from "ws";
|
import { WebSocketServer, WebSocket } from "ws";
|
||||||
import Redis from "ioredis";
|
import { NextRequest } from "next/server";
|
||||||
import { NextRequest, NextResponse } from "next/server";
|
|
||||||
import { LiveRunEvent } from "@/modules/live-run/types";
|
import { LiveRunEvent } from "@/modules/live-run/types";
|
||||||
import { updateLiveRunViews } from "@/modules/live-run/core";
|
import { updateLiveRunViews } from "@/modules/live-run/core";
|
||||||
import { redis } from "@/lib/redis";
|
import { redis, redisSubscriber } from "@/lib/redis";
|
||||||
|
|
||||||
export function GET(req: NextRequest) {
|
// Map of runId → Set of WebSockets
|
||||||
const headers = new Headers();
|
const runSubscribers: Map<string, Set<WebSocket>> = new Map();
|
||||||
headers.set("Connection", "Upgrade");
|
|
||||||
headers.set("Upgrade", "websocket");
|
|
||||||
return new Response("Upgrade Required", { status: 426, headers });
|
|
||||||
}
|
|
||||||
|
|
||||||
export function UPGRADE(client: WebSocket, server: WebSocketServer) {
|
// Handle incoming Pub/Sub messages
|
||||||
const subscriber = new Redis(process.env.REDIS_URL!);
|
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;
|
let runId: string;
|
||||||
|
|
||||||
|
// Receive runId from client
|
||||||
client.once("message", async (message) => {
|
client.once("message", async (message) => {
|
||||||
runId = message.toString();
|
runId = message.toString();
|
||||||
|
|
||||||
|
// Track view count
|
||||||
await updateLiveRunViews(runId, 1);
|
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 () => {
|
// Add WebSocket to subscribers
|
||||||
const text = await redis.get(`liveRunText:${runId}`);
|
if (!runSubscribers.has(runId)) {
|
||||||
|
runSubscribers.set(runId, new Set());
|
||||||
if (client.readyState === WebSocket.OPEN) {
|
await redisSubscriber.subscribe(`liveRunEvent:${runId}`);
|
||||||
client.send(
|
}
|
||||||
JSON.stringify({
|
runSubscribers.get(runId)!.add(client);
|
||||||
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 }));
|
|
||||||
}
|
|
||||||
});
|
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Handle client disconnect
|
||||||
client.once("close", async () => {
|
client.once("close", async () => {
|
||||||
await subscriber.unsubscribe(`liveRunEvent:${runId}`);
|
const set = runSubscribers.get(runId);
|
||||||
subscriber.quit();
|
if (set) {
|
||||||
|
set.delete(client);
|
||||||
|
if (set.size === 0) {
|
||||||
|
await redisSubscriber.unsubscribe(`liveRunEvent:${runId}`);
|
||||||
|
runSubscribers.delete(runId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
await updateLiveRunViews(runId, -1);
|
await updateLiveRunViews(runId, -1);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -69,3 +65,10 @@ export function UPGRADE(client: WebSocket, server: WebSocketServer) {
|
|||||||
console.error("WebSocket error:", err);
|
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 });
|
||||||
|
}
|
||||||
|
|||||||
@@ -3,6 +3,12 @@ import Redis from "ioredis";
|
|||||||
|
|
||||||
export const redis = new Redis(process.env.REDIS_URL!);
|
export const redis = new Redis(process.env.REDIS_URL!);
|
||||||
|
|
||||||
|
export const redisSubscriber = new Redis(process.env.REDIS_URL!);
|
||||||
|
|
||||||
redis.on("error", (err) => {
|
redis.on("error", (err) => {
|
||||||
console.error("Redis error:", err);
|
console.error("Redis error:", err);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
redisSubscriber.on("error", (err) => {
|
||||||
|
console.error("Redis subscriber error:", err);
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user