From 6c9b703ed1f00dd15aa9cc37fe8b24691c76b2d4 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Mon, 9 Feb 2026 13:41:03 +0100 Subject: [PATCH] Using next-ws for better code structure --- server.mjs | 60 ---------------------- server.ts | 86 -------------------------------- server.tsconfig.json | 16 ------ src/app/api/run/move/route.ts | 35 +++++++++++++ src/app/api/view-run/route.ts | 47 +++++++++++++++++ src/app/api/view-run/ws/route.ts | 14 ------ src/modules/live-run/core.ts | 17 +++++++ 7 files changed, 99 insertions(+), 176 deletions(-) delete mode 100644 server.mjs delete mode 100644 server.ts delete mode 100644 server.tsconfig.json create mode 100644 src/app/api/run/move/route.ts create mode 100644 src/app/api/view-run/route.ts delete mode 100644 src/app/api/view-run/ws/route.ts diff --git a/server.mjs b/server.mjs deleted file mode 100644 index 988b8c6..0000000 --- a/server.mjs +++ /dev/null @@ -1,60 +0,0 @@ -import { createServer } from 'http'; -import next from 'next'; -import { WebSocketServer } from 'ws'; - -const dev = process.env.NODE_ENV !== 'production'; -const app = next({ dev }); -const handle = app.getRequestHandler(); - -app.prepare().then(() => { - const server = createServer((req, res) => { - handle(req, res); - }); - - const runs = {}; - - const wss = new WebSocketServer({ noServer: true }); - - // Handle WebSocket connections - wss.on('connection', (ws, req) => { - ws.on('message', (msg) => { - const runId = msg.toString(); - if (!runs[runId]) runs[runId] = { clients: [] } - runs[runId].clients.push(ws); - }); - }); - - // Handle POST requests - server.on('request', async (req, res) => { - console.log(req.method, req.url) - if (req.method === 'POST' && req.url === '/api/run/move') { - let body = ''; - req.on('data', chunk => body += chunk); - req.on('end', () => { - const move = JSON.parse(body); - - runs[move.runId]?.clients.forEach(client => { - if (client.readyState === client.OPEN) { - client.send(JSON.stringify(move.moves)); - } - }); - - res.writeHead(200, { 'Content-Type': 'application/json' }); - res.end(JSON.stringify({ status: 'ok' })); - }); - } - }); - - // Upgrade HTTP connections to WebSocket when needed - server.on('upgrade', (req, socket, head) => { - if (req.url === '/api/view-run') { - wss.handleUpgrade(req, socket, head, (ws) => { - wss.emit('connection', ws, req); - }); - } - }); - - server.listen(3000, () => { - console.log('Server running on http://localhost:3000'); - }); -}); diff --git a/server.ts b/server.ts deleted file mode 100644 index 56f1e16..0000000 --- a/server.ts +++ /dev/null @@ -1,86 +0,0 @@ -import { createServer, IncomingMessage } from "http"; -import next from "next"; -import { WebSocketServer, WebSocket } from "ws"; - -const dev = process.env.NODE_ENV !== "production"; -const app = next({ dev }); -const handle = app.getRequestHandler(); - -/** - * One run = many websocket clients - */ -type LiveRun = { - clients: Set; -}; - -const runs: Record = {}; - -app.prepare().then(() => { - const server = createServer((req, res) => { - handle(req, res); - }); - - const wss = new WebSocketServer({ noServer: true }); - - // WebSocket connections - wss.on("connection", (ws: WebSocket, req: IncomingMessage) => { - ws.on("message", (msg: Buffer) => { - const runId = msg.toString(); - - if (!runs[runId]) { - runs[runId] = { clients: new Set() }; - } - - runs[runId].clients.add(ws); - - ws.once("close", () => { - runs[runId]?.clients.delete(ws); - }); - }); - }); - - // Handle POST requests - server.on("request", (req, res) => { - if (req.method === "POST" && req.url === "/api/run/move") { - let body = ""; - - req.on("data", (chunk: Buffer) => { - body += chunk.toString(); - }); - - req.on("end", () => { - try { - const move: { - runId: string; - moves: unknown; - } = JSON.parse(body); - - runs[move.runId]?.clients.forEach((client) => { - if (client.readyState === WebSocket.OPEN) { - client.send(JSON.stringify(move.moves)); - } - }); - - res.writeHead(200, { "Content-Type": "application/json" }); - res.end(JSON.stringify({ status: "ok" })); - } catch (err) { - res.writeHead(400); - res.end("Invalid JSON"); - } - }); - } - }); - - // Upgrade HTTP → WebSocket - server.on("upgrade", (req, socket, head) => { - if (req.url === "/api/view-run") { - wss.handleUpgrade(req, socket, head, (ws) => { - wss.emit("connection", ws, req); - }); - } - }); - - server.listen(3000, () => { - console.log("Server running on http://localhost:3000"); - }); -}); diff --git a/server.tsconfig.json b/server.tsconfig.json deleted file mode 100644 index aaced07..0000000 --- a/server.tsconfig.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "extends": "./tsconfig.json", - "compilerOptions": { - "noEmit": false, - "outDir": "./.next", - "target": "ES2020", - "module": "NodeNext", - "moduleResolution": "NodeNext", - "lib": ["ES2020"], - "types": ["node"], - "skipLibCheck": true, - "incremental": false, - "composite": false - }, - "include": ["server.ts"] -} diff --git a/src/app/api/run/move/route.ts b/src/app/api/run/move/route.ts new file mode 100644 index 0000000..d691e2b --- /dev/null +++ b/src/app/api/run/move/route.ts @@ -0,0 +1,35 @@ +import { NextRequest, NextResponse } from "next/server"; +import { createClient } from "@/lib/supabase/server"; +import * as Core from "@/modules/live-run/core"; +import { getSession } from "@/modules/account/core"; + +export async function POST(req: NextRequest) { + const { runId, moves } = await req.json(); + + if (!runId || !moves) + return NextResponse.json( + { error: "body requires runId and moves" }, + { status: 400 }, + ); + + const supabase = await createClient(); + + const session = await getSession(supabase); + + if (!session) + return NextResponse.json({ error: "Invalid session" }, { status: 401 }); + + const liveRun = await Core.getLiveRun(runId); + + if (!liveRun) + return NextResponse.json( + { error: "Run is not live or does not exist" }, + { status: 401 }, + ); + + console.log(moves); + + const run = await Core.addLiveRunMoves(liveRun.id, moves); + + return NextResponse.json({ ok: true, run }); +} diff --git a/src/app/api/view-run/route.ts b/src/app/api/view-run/route.ts new file mode 100644 index 0000000..72e0e73 --- /dev/null +++ b/src/app/api/view-run/route.ts @@ -0,0 +1,47 @@ +import { WebSocketServer, WebSocket } from "ws"; +import Redis from "ioredis"; +import { LiveRunMove } from "@/modules/live-run/types"; +import { NextRequest, NextResponse } from "next/server"; + +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) { + console.log("A client connected"); + + const subscriber = new Redis(); + let runId: string; + + client.once("message", (message) => { + runId = message.toString(); + + subscriber.subscribe(`liveRunMoves:${runId}`, (err, count) => { + if (err) { + console.error("Failed to subscribe", err); + client.close(1011, "Redis subscription failed"); + } else { + console.log(`Subscribed to ${count} channels`); + } + }); + + subscriber.on("message", (channel, message) => { + if (client.readyState === WebSocket.OPEN) { + client.send(message); + } + }); + }); + + client.once("close", async () => { + if (runId) await subscriber.unsubscribe(`liveRun:${runId}`); + subscriber.quit(); + console.log("Client disconnected, Redis unsubscribed"); + }); + + client.on("error", (err) => { + console.error("WebSocket error:", err); + }); +} diff --git a/src/app/api/view-run/ws/route.ts b/src/app/api/view-run/ws/route.ts deleted file mode 100644 index 71fc7f3..0000000 --- a/src/app/api/view-run/ws/route.ts +++ /dev/null @@ -1,14 +0,0 @@ -import { WebSocketServer, WebSocket } from "ws"; - -export function UPGRADE(client: WebSocket, server: WebSocketServer) { - console.log("A client connected"); - - client.on("message", (message) => { - console.log("Received message:", message); - client.send(message); - }); - - client.once("close", () => { - console.log("A client disconnected"); - }); -} diff --git a/src/modules/live-run/core.ts b/src/modules/live-run/core.ts index a90e73d..40784fc 100644 --- a/src/modules/live-run/core.ts +++ b/src/modules/live-run/core.ts @@ -4,6 +4,7 @@ import "server-only"; import { LiveRun, LiveRunMove } from "./types"; import { SupabaseClient } from "@supabase/supabase-js"; import { redis } from "@/lib/redis"; +import Redis from "ioredis"; /** -------------------- Live Runs in Redis -------------------- */ @@ -43,6 +44,22 @@ export async function removeLiveRun(id: string) { await redis.del(`liveRun:${id}`, `liveRunMoves:${id}`); } +/** -------------------- Live Run Moves -------------------- */ + +/** + * Add moves to a live run (ephemeral) + */ +export async function addLiveRunMoves(runId: string, moves: LiveRunMove[]) { + const runExists = await redis.exists(`liveRun:${runId}`); + + if (!runExists) throw "Invalid live run"; + + const publisher = new Redis(); + + // Send a message + await publisher.publish(`liveRunMoves:${runId}`, JSON.stringify(moves)); +} + /** -------------------- Submit Run -------------------- */ /**