204 lines
5.8 KiB
JavaScript
204 lines
5.8 KiB
JavaScript
// Minimal RFC 6455 WebSocket server built only on Node's built-in modules.
|
|
// The project rule is "no external JS libraries other than jQuery", so the
|
|
// transport is implemented here rather than pulled from npm. It supports the
|
|
// subset the game needs: text and binary messages, fragmentation, ping/pong,
|
|
// close, and large frames (state snapshots are hundreds of KiB).
|
|
|
|
import { createHash } from "node:crypto";
|
|
import { EventEmitter } from "node:events";
|
|
|
|
const GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11";
|
|
|
|
const OP_CONTINUATION = 0x0;
|
|
const OP_TEXT = 0x1;
|
|
const OP_BINARY = 0x2;
|
|
const OP_CLOSE = 0x8;
|
|
const OP_PING = 0x9;
|
|
const OP_PONG = 0xa;
|
|
|
|
export class WebSocketServer extends EventEmitter {
|
|
constructor(httpServer, options = {}) {
|
|
super();
|
|
this.path = options.path || "/ws";
|
|
this._nextId = 1;
|
|
this._connections = new Map();
|
|
httpServer.on("upgrade", (request, socket, head) => {
|
|
const url = request.url || "";
|
|
const path = url.split("?")[0];
|
|
if (path !== this.path) {
|
|
socket.destroy();
|
|
return;
|
|
}
|
|
this._accept(request, socket, head);
|
|
});
|
|
}
|
|
|
|
connectionCount() {
|
|
return this._connections.size;
|
|
}
|
|
|
|
_accept(request, socket, head) {
|
|
const wsKey = request.headers["sec-websocket-key"];
|
|
if (!wsKey || (request.headers.upgrade || "").toLowerCase() !== "websocket") {
|
|
socket.destroy();
|
|
return;
|
|
}
|
|
const accept = createHash("sha1").update(wsKey + GUID).digest("base64");
|
|
socket.write(
|
|
"HTTP/1.1 101 Switching Protocols\r\n" +
|
|
"Upgrade: websocket\r\n" +
|
|
"Connection: Upgrade\r\n" +
|
|
`Sec-WebSocket-Accept: ${accept}\r\n` +
|
|
"\r\n"
|
|
);
|
|
socket.setNoDelay(true);
|
|
|
|
const id = this._nextId++;
|
|
const connection = new Connection(id, socket, head);
|
|
this._connections.set(id, connection);
|
|
this.emit("connection", connection);
|
|
connection.on("close", () => this._connections.delete(id));
|
|
}
|
|
}
|
|
|
|
class Connection extends EventEmitter {
|
|
constructor(id, socket, head) {
|
|
super();
|
|
this.id = id;
|
|
this.socket = socket;
|
|
this._buffer = Buffer.alloc(0);
|
|
this._fragments = [];
|
|
this._fragmentOpcode = null;
|
|
this._closed = false;
|
|
|
|
if (head && head.length > 0) this._buffer = Buffer.concat([this._buffer, head]);
|
|
socket.on("data", (chunk) => {
|
|
this._buffer = Buffer.concat([this._buffer, chunk]);
|
|
this._parse();
|
|
});
|
|
socket.on("error", () => this._close());
|
|
socket.on("close", () => this._close());
|
|
if (this._buffer.length > 0) this._parse();
|
|
}
|
|
|
|
send(data) {
|
|
if (this._closed) return;
|
|
const payload = Buffer.isBuffer(data) ? data : Buffer.from(String(data), "utf8");
|
|
const opcode = Buffer.isBuffer(data) ? OP_BINARY : OP_TEXT;
|
|
try {
|
|
this.socket.write(encodeFrame(opcode, payload));
|
|
} catch {
|
|
this._close();
|
|
}
|
|
}
|
|
|
|
close(code = 1000, reason = "") {
|
|
if (this._closed) return;
|
|
const reasonBuffer = Buffer.from(reason, "utf8");
|
|
const payload = Buffer.alloc(2 + reasonBuffer.length);
|
|
payload.writeUInt16BE(code, 0);
|
|
reasonBuffer.copy(payload, 2);
|
|
try {
|
|
// Flush the close frame before shutting the socket down; destroying
|
|
// immediately can discard it.
|
|
this.socket.write(encodeFrame(OP_CLOSE, payload));
|
|
this.socket.end();
|
|
} catch {
|
|
this._close();
|
|
}
|
|
}
|
|
|
|
_close() {
|
|
if (this._closed) return;
|
|
this._closed = true;
|
|
try {
|
|
this.socket.destroy();
|
|
} catch {
|
|
// ignore
|
|
}
|
|
this.emit("close");
|
|
}
|
|
|
|
_parse() {
|
|
while (true) {
|
|
const buffer = this._buffer;
|
|
if (buffer.length < 2) return;
|
|
const b0 = buffer[0];
|
|
const b1 = buffer[1];
|
|
const fin = (b0 & 0x80) !== 0;
|
|
const opcode = b0 & 0x0f;
|
|
const masked = (b1 & 0x80) !== 0;
|
|
let length = b1 & 0x7f;
|
|
let offset = 2;
|
|
if (length === 126) {
|
|
if (buffer.length < 4) return;
|
|
length = buffer.readUInt16BE(2);
|
|
offset = 4;
|
|
} else if (length === 127) {
|
|
if (buffer.length < 10) return;
|
|
length = Number(buffer.readBigUInt64BE(2));
|
|
offset = 10;
|
|
}
|
|
let maskKey = null;
|
|
if (masked) {
|
|
if (buffer.length < offset + 4) return;
|
|
maskKey = buffer.subarray(offset, offset + 4);
|
|
offset += 4;
|
|
}
|
|
if (buffer.length < offset + length) return;
|
|
let payload = buffer.subarray(offset, offset + length);
|
|
if (masked) {
|
|
const out = Buffer.allocUnsafe(length);
|
|
for (let i = 0; i < length; i++) out[i] = payload[i] ^ maskKey[i & 3];
|
|
payload = out;
|
|
}
|
|
this._buffer = buffer.subarray(offset + length);
|
|
this._handleFrame(fin, opcode, payload);
|
|
}
|
|
}
|
|
|
|
_handleFrame(fin, opcode, payload) {
|
|
if (opcode === OP_CLOSE) {
|
|
this.close(1000);
|
|
return;
|
|
}
|
|
if (opcode === OP_PING) {
|
|
this.socket.write(encodeFrame(OP_PONG, payload));
|
|
return;
|
|
}
|
|
if (opcode === OP_PONG) return;
|
|
|
|
if (opcode === OP_CONTINUATION) {
|
|
this._fragments.push(payload);
|
|
} else {
|
|
this._fragmentOpcode = opcode;
|
|
this._fragments = [payload];
|
|
}
|
|
if (!fin) return;
|
|
const full = Buffer.concat(this._fragments);
|
|
const isBinary = this._fragmentOpcode === OP_BINARY;
|
|
this._fragments = [];
|
|
this._fragmentOpcode = null;
|
|
this.emit("message", isBinary ? full : full.toString("utf8"));
|
|
}
|
|
}
|
|
|
|
function encodeFrame(opcode, payload) {
|
|
const length = payload.length;
|
|
let header;
|
|
if (length < 126) {
|
|
header = Buffer.alloc(2);
|
|
header[1] = length;
|
|
} else if (length < 65536) {
|
|
header = Buffer.alloc(4);
|
|
header[1] = 126;
|
|
header.writeUInt16BE(length, 2);
|
|
} else {
|
|
header = Buffer.alloc(10);
|
|
header[1] = 127;
|
|
header.writeBigUInt64BE(BigInt(length), 2);
|
|
}
|
|
header[0] = 0x80 | opcode;
|
|
return Buffer.concat([header, payload]);
|
|
}
|