Using next-ws for better code structure
This commit is contained in:
-60
@@ -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');
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -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<WebSocket>;
|
|
||||||
};
|
|
||||||
|
|
||||||
const runs: Record<string, LiveRun> = {};
|
|
||||||
|
|
||||||
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");
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -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"]
|
|
||||||
}
|
|
||||||
@@ -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 });
|
||||||
|
}
|
||||||
@@ -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);
|
||||||
|
});
|
||||||
|
}
|
||||||
@@ -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");
|
|
||||||
});
|
|
||||||
}
|
|
||||||
@@ -4,6 +4,7 @@ import "server-only";
|
|||||||
import { LiveRun, LiveRunMove } from "./types";
|
import { LiveRun, LiveRunMove } from "./types";
|
||||||
import { SupabaseClient } from "@supabase/supabase-js";
|
import { SupabaseClient } from "@supabase/supabase-js";
|
||||||
import { redis } from "@/lib/redis";
|
import { redis } from "@/lib/redis";
|
||||||
|
import Redis from "ioredis";
|
||||||
|
|
||||||
/** -------------------- Live Runs in Redis -------------------- */
|
/** -------------------- Live Runs in Redis -------------------- */
|
||||||
|
|
||||||
@@ -43,6 +44,22 @@ export async function removeLiveRun(id: string) {
|
|||||||
await redis.del(`liveRun:${id}`, `liveRunMoves:${id}`);
|
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 -------------------- */
|
/** -------------------- Submit Run -------------------- */
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
Reference in New Issue
Block a user