You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

313 lines
7.9 KiB

import { createHash, randomUUID } from 'node:crypto';
import { createServer } from 'node:http';
const HOST = process.env.CONN_CHAT_HOST ?? '127.0.0.1';
const PORT = Number(process.env.CONN_CHAT_PORT ?? 8787);
const WEBSOCKET_GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11';
const HEADER_UPGRADE = 'upgrade';
const HEADER_WEBSOCKET_KEY = 'sec-websocket-key';
const UPGRADE_WEBSOCKET = 'websocket';
const OPCODE_TEXT = 0x1;
const OPCODE_CLOSE = 0x8;
const OPCODE_PING = 0x9;
const OPCODE_PONG = 0xa;
const CLOSE_CODE_NORMAL = 1000;
const CLOSE_CODE_PROTOCOL_ERROR = 1002;
const FRAME_TYPE_AUTH = 'conn.auth';
const FRAME_TYPE_JOIN = 'conn.join';
const FRAME_TYPE_LEAVE = 'conn.leave';
const FRAME_TYPE_PING = 'conn.ping';
const FRAME_TYPE_PONG = 'conn.pong';
const SERVER_REPLY_EVENT = 'server.reply';
const CHAT_EVENT_MESSAGE = 'chat.message';
const CHAT_EVENT_SYSTEM = 'chat.system';
const CHAT_EVENT_TYPING = 'chat.typing';
const CHAT_EVENT_WHOAMI = 'chat.whoami';
const DEMO_EVENT_DISCONNECT = 'demo.disconnect';
const BOT_NAME = 'Ava Bot';
const DEFAULT_ROOM = 'room:lobby';
const DEFAULT_USER = 'Anon';
const BOT_REPLY_DELAY_MS = 450;
const FORCED_CLOSE_DELAY_MS = 80;
const clients = new Map();
function makeAcceptKey(key) {
return createHash('sha1').update(`${key}${WEBSOCKET_GUID}`).digest('base64');
}
function encodeFrame(opcode, payload = Buffer.alloc(0)) {
const data = Buffer.isBuffer(payload) ? payload : Buffer.from(String(payload));
const len = data.length;
let header;
if (len < 126) {
header = Buffer.alloc(2);
header[1] = len;
} else if (len <= 0xffff) {
header = Buffer.alloc(4);
header[1] = 126;
header.writeUInt16BE(len, 2);
} else {
header = Buffer.alloc(10);
header[1] = 127;
header.writeBigUInt64BE(BigInt(len), 2);
}
header[0] = 0x80 | opcode;
return Buffer.concat([header, data]);
}
function sendText(socket, text) {
if (socket.destroyed) return;
socket.write(encodeFrame(OPCODE_TEXT, text));
}
function sendJson(socket, frame) {
sendText(socket, JSON.stringify(frame));
}
function closeSocket(socket, code = CLOSE_CODE_NORMAL, reason = '') {
if (socket.destroyed) return;
const body = Buffer.alloc(2 + Buffer.byteLength(reason));
body.writeUInt16BE(code, 0);
body.write(reason, 2);
socket.write(encodeFrame(OPCODE_CLOSE, body));
socket.end();
}
function createMessage(user, text, kind) {
return {
id: randomUUID(),
user,
text,
at: Date.now(),
kind
};
}
function reply(socket, frame, payload, error) {
if (typeof frame.id !== 'string') return;
sendJson(socket, {
type: SERVER_REPLY_EVENT,
replyTo: frame.id,
payload,
error
});
}
function broadcast(room, frame) {
for (const [socket, client] of clients) {
if (client.room !== room) continue;
sendJson(socket, frame);
}
}
function broadcastSystem(room, text) {
broadcast(room, {
topic: room,
type: CHAT_EVENT_SYSTEM,
payload: createMessage('Sistema', text, 'system')
});
}
function parseFrames(client) {
const out = [];
let offset = 0;
const buffer = client.buffer;
while (offset + 2 <= buffer.length) {
const first = buffer[offset];
const second = buffer[offset + 1];
const opcode = first & 0x0f;
const masked = (second & 0x80) !== 0;
let len = second & 0x7f;
let cursor = offset + 2;
if (len === 126) {
if (cursor + 2 > buffer.length) break;
len = buffer.readUInt16BE(cursor);
cursor += 2;
} else if (len === 127) {
if (cursor + 8 > buffer.length) break;
const longLen = buffer.readBigUInt64BE(cursor);
if (longLen > BigInt(Number.MAX_SAFE_INTEGER)) {
throw new Error('websocket frame too large');
}
len = Number(longLen);
cursor += 8;
}
const maskOffset = cursor;
if (masked) cursor += 4;
if (cursor + len > buffer.length) break;
const payload = Buffer.from(buffer.subarray(cursor, cursor + len));
if (masked) {
const mask = buffer.subarray(maskOffset, maskOffset + 4);
for (let i = 0; i < payload.length; i += 1) payload[i] ^= mask[i % 4];
}
out.push({ opcode, payload });
offset = cursor + len;
}
client.buffer = buffer.subarray(offset);
return out;
}
function handleFrame(socket, rawFrame) {
const client = clients.get(socket);
if (client === undefined) return;
if (rawFrame.opcode === OPCODE_CLOSE) {
closeSocket(socket);
return;
}
if (rawFrame.opcode === OPCODE_PING) {
socket.write(encodeFrame(OPCODE_PONG, rawFrame.payload));
return;
}
if (rawFrame.opcode !== OPCODE_TEXT) return;
let frame;
try {
frame = JSON.parse(rawFrame.payload.toString('utf8'));
} catch {
closeSocket(socket, CLOSE_CODE_PROTOCOL_ERROR, 'invalid json');
return;
}
if (frame.type === FRAME_TYPE_AUTH) {
client.authed = true;
client.user = frame.payload?.user ?? client.user;
reply(socket, frame, { accepted: true, server: 'conn-chat-server' });
return;
}
if (frame.type === FRAME_TYPE_JOIN) {
const room = typeof frame.topic === 'string' ? frame.topic : DEFAULT_ROOM;
client.room = room;
client.user = frame.payload?.user ?? client.user;
broadcastSystem(room, `${client.user} ha entrado en ${room}.`);
return;
}
if (frame.type === FRAME_TYPE_LEAVE) {
const room = client.room;
client.room = null;
if (room !== null) broadcastSystem(room, `${client.user} ha salido de ${room}.`);
return;
}
if (frame.type === FRAME_TYPE_PING) {
sendJson(socket, { type: FRAME_TYPE_PONG, payload: { at: Date.now() } });
return;
}
if (frame.type === CHAT_EVENT_WHOAMI) {
reply(socket, frame, {
user: client.user,
room: client.room,
connection: client.id,
generation: client.generation
});
return;
}
if (frame.type === DEMO_EVENT_DISCONNECT) {
setTimeout(() => {
closeSocket(socket, CLOSE_CODE_NORMAL, 'demo disconnect');
}, FORCED_CLOSE_DELAY_MS);
return;
}
if (frame.type === CHAT_EVENT_MESSAGE) {
const room = typeof frame.topic === 'string' ? frame.topic : client.room;
if (room === null) return;
broadcast(room, frame);
broadcast(room, {
topic: room,
type: CHAT_EVENT_TYPING,
payload: { user: BOT_NAME, active: true }
});
setTimeout(() => {
broadcast(room, {
topic: room,
type: CHAT_EVENT_TYPING,
payload: { user: BOT_NAME, active: false }
});
broadcast(room, {
topic: room,
type: CHAT_EVENT_MESSAGE,
payload: createMessage(
BOT_NAME,
`Recibido por ${client.user}: "${frame.payload?.text ?? ''}"`,
'bot'
)
});
}, BOT_REPLY_DELAY_MS);
}
}
const server = createServer((_, response) => {
response.writeHead(200, { 'content-type': 'text/plain; charset=utf-8' });
response.end('conn chat websocket server\n');
});
server.on('upgrade', (request, socket) => {
const upgrade = request.headers[HEADER_UPGRADE];
const key = request.headers[HEADER_WEBSOCKET_KEY];
if (upgrade !== UPGRADE_WEBSOCKET || typeof key !== 'string') {
socket.destroy();
return;
}
socket.write(
[
'HTTP/1.1 101 Switching Protocols',
'Upgrade: websocket',
'Connection: Upgrade',
`Sec-WebSocket-Accept: ${makeAcceptKey(key)}`,
'',
''
].join('\r\n')
);
const client = {
id: randomUUID(),
user: DEFAULT_USER,
room: null,
authed: false,
generation: Date.now(),
buffer: Buffer.alloc(0)
};
clients.set(socket, client);
console.log(`[conn-chat] connected ${client.id}`);
socket.on('data', (chunk) => {
const current = clients.get(socket);
if (current === undefined) return;
current.buffer = Buffer.concat([current.buffer, chunk]);
try {
for (const frame of parseFrames(current)) handleFrame(socket, frame);
} catch (error) {
console.error('[conn-chat] protocol error', error);
closeSocket(socket, CLOSE_CODE_PROTOCOL_ERROR, 'protocol error');
}
});
socket.on('close', () => {
const current = clients.get(socket);
clients.delete(socket);
if (current?.room) broadcastSystem(current.room, `${current.user} se ha desconectado.`);
console.log(`[conn-chat] disconnected ${current?.id ?? 'unknown'}`);
});
socket.on('error', (error) => {
console.error('[conn-chat] socket error', error);
});
});
server.listen(PORT, HOST, () => {
console.log(`[conn-chat] ws://${HOST}:${PORT}`);
});

Powered by TurnKey Linux.