Working ish architecture

This commit is contained in:
2026-02-09 11:26:12 +01:00
parent 04bb5d1b08
commit 25fd5f207e
7 changed files with 123 additions and 108 deletions
+2 -2
View File
@@ -7,11 +7,11 @@
"license": "Apache-2.0", "license": "Apache-2.0",
"type": "module", "type": "module",
"bin": { "bin": {
"next-websockets": "dist/index.js" "next-websockets": "dist/bin/index.js"
}, },
"scripts": { "scripts": {
"build": "tsc", "build": "tsc",
"dev": "tsc && node dist/index.js" "dev": "tsc && node dist/bin/index.js"
}, },
"dependencies": { "dependencies": {
"@types/cli-color": "^2.0.6", "@types/cli-color": "^2.0.6",
+76
View File
@@ -0,0 +1,76 @@
#!/usr/bin/env node
import clc from "cli-color";
import { createServer, IncomingMessage } from "http";
import next from "next";
import path from "path";
import { pathToFileURL } from "url";
import { WebSocketServer, WebSocket } from "ws";
import { NextWebsocketConfig } from "..";
import { Readable } from "stream";
const dev = process.env.NODE_ENV !== "production";
const app = next({ dev });
const handle = app.getRequestHandler();
const config = await loadOptionalConfig();
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) => {});
// Upgrade HTTP → WebSocket
server.on("upgrade", (req, socket, head) => {
if (!req.url) return;
config?.app.emit("upgrade", req.url, toRequest(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(clc.red(" ▲ Next.js (websockets)"));
console.log(` - Local: http://localhost:3000`);
console.log(` - Network: http://192.168.100.3:3000`);
console.log(` - Environments: .env`);
});
});
export async function loadOptionalConfig(
filename = "ws.config.ts",
): Promise<NextWebsocketConfig | undefined> {
const absPath = path.resolve(process.cwd(), filename);
try {
const mod = await import(pathToFileURL(absPath).href);
return mod.default ?? mod;
} catch {
return undefined;
}
}
function toRequest(req: IncomingMessage): Request {
const host = req.headers.host ?? "localhost";
const url = new URL(req.url ?? "/", `http://${host}`);
const body =
req.method !== "GET" && req.method !== "HEAD"
? new ReadableStream(Readable.toWeb(req))
: undefined;
return new Request(url.toString(), {
method: req.method,
headers: req.headers as HeadersInit,
body,
});
}
View File
+35 -81
View File
@@ -1,92 +1,46 @@
#!/usr/bin/env node import { Duplex } from "stream";
import clc from "cli-color"; export interface NextWebsocketConfig {
import { createServer, IncomingMessage } from "http"; app: NextWebsocketApp;
import next from "next"; }
import { WebSocketServer, WebSocket } from "ws";
const dev = process.env.NODE_ENV !== "production"; export interface WsRouteHandler {
const app = next({ dev }); upgrade?: (req: Request, socket: Duplex, head: Buffer<ArrayBuffer>) => void;
const handle = app.getRequestHandler(); post?: (req: Request) => void;
}
/** export class NextWebsocketApp {
* One run = many websocket clients handlers: Record<string, WsRouteHandler> = {};
*/
type LiveRun = {
clients: Set<WebSocket>;
};
const runs: Record<string, LiveRun> = {}; on<K extends keyof WsRouteHandler>(
method: K,
app.prepare().then(() => { route: string,
const server = createServer((req, res) => { handler: WsRouteHandler[K],
handle(req, res); ) {
}); if (!this.handlers[route]) this.handlers[route] = {};
this.handlers[route][method] = handler;
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); upgrade(
route: string,
ws.once("close", () => { handler: (req: Request, socket: Duplex, head: Buffer<ArrayBuffer>) => void,
runs[runId]?.clients.delete(ws); ) {
}); this.on("upgrade", route, handler);
});
});
// 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" }); post(route: string, handler: () => void) {
res.end(JSON.stringify({ status: "ok" })); this.on("post", route, handler);
} catch (err) {
res.writeHead(400);
res.end("Invalid JSON");
} }
});
}
});
// Upgrade HTTP → WebSocket emit<K extends keyof WsRouteHandler>(
server.on("upgrade", (req, socket, head) => { method: K,
if (req.url === "/api/view-run") { route: string,
wss.handleUpgrade(req, socket, head, (ws) => { ...args: Parameters<NonNullable<WsRouteHandler[K]>>
wss.emit("connection", ws, req); ) {
}); const handler = this.handlers[route]?.[method] as
} | ((...args: Parameters<NonNullable<WsRouteHandler[K]>>) => unknown)
}); | undefined;
server.listen(3000, () => { handler?.(...args);
console.log(clc.red(" ▲ Next.js (websockets)")); }
console.log(` - Local: http://localhost:3000`); }
console.log(` - Network: http://192.168.100.3:3000`);
console.log(` - Environments: .env`);
});
});
View File
-22
View File
@@ -1,22 +0,0 @@
import fs from "fs";
import path from "path";
export function listWsFiles(dir: string) {
const files: string[] = [];
function walk(currentDir: string) {
const entries = fs.readdirSync(currentDir, { withFileTypes: true });
for (const entry of entries) {
const fullPath = path.join(currentDir, entry.name);
if (entry.isDirectory()) {
walk(fullPath);
} else if (entry.isFile() && entry.name === "ws.ts") {
files.push(fullPath);
}
}
}
walk(dir);
return files;
}
+7
View File
@@ -0,0 +1,7 @@
import { NextWebsocketApp, NextWebsocketConfig } from "./src/index";
const app = new NextWebsocketApp();
export default {
app,
} satisfies NextWebsocketConfig;