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.

67 lines
1.9 KiB

Nexo realtime — WebSocket distributor for chat fan-out Replaces the chat poll-every-3-seconds compromise with a real WebSocket pipeline. Server side: `realtime.mjs`: - In-memory hub mapping `userId → Set<WebSocket>`. Single-process, zero external deps; fine for the local demo. A clustered deployment would swap in Redis pub/sub behind the same surface. - `attach(userId, socket)` registers a socket and self-detaches on close/error. - `broadcast(userIds, event)` JSON-encodes once and dispatches to every subscriber of every listed user, swallowing per-socket errors. `server.mjs`: - Adds a `WebSocketServer({ noServer: true })` that listens on the HTTP server's `'upgrade'` event for `/api/realtime` paths. - Authentication mirrors the HTTP path: extract `dating_session` from the upgrade request's `Cookie` header, pass through `currentSession` and reject with 401 if it doesn't resolve. - Rejects upgrades from origins outside the CORS allowlist (the cookie-based auth is the second line of defence; origin gating is the first). - On accept, attaches the socket to the hub and sends a `hello` envelope so the client can confirm authentication round-trip. `routes.mjs:messagesPost`: - After persisting a new message, calls `broadcast(relationIds(match.users), { type: 'dating.message.created', message: view })`. Both members (sender and counterpart) get the push, so multi-device sessions stay in sync. Adds `ws` (8.20) + `@types/ws` to dependencies. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
5 months ago
/**
* Realtime distributor for the dating server. Holds an in-memory map
* of `userId -> Set<WebSocket>` so handlers can broadcast events to
* specific users without each handler caring about transport details.
*
* Single-process only: this hub does NOT persist subscriptions or
* coordinate across instances. Sufficient for the local demo; a
* production deployment would replace this with Redis pub/sub or
* similar.
*/
const sockets = new Map(); // userId -> Set<WebSocket>
export function attach(userId, socket) {
if (!userId || !socket) return () => {};
let set = sockets.get(userId);
if (!set) {
set = new Set();
sockets.set(userId, set);
}
set.add(socket);
const detach = () => {
const current = sockets.get(userId);
if (!current) return;
current.delete(socket);
if (current.size === 0) sockets.delete(userId);
};
socket.on('close', detach);
socket.on('error', detach);
return detach;
}
/**
* Send `event` to every socket subscribed by any of `userIds`. Each
* `event` is a plain object that gets JSON-encoded once for all
* recipients.
*
* Failures on individual sockets are swallowed: the close/error
* handlers attached in `attach` clean them up later, and a single bad
* socket must never block the rest of the broadcast.
*/
export function broadcast(userIds, event) {
if (!Array.isArray(userIds) || userIds.length === 0) return;
const payload = JSON.stringify(event);
const seen = new Set();
for (const userId of userIds) {
if (!userId || seen.has(userId)) continue;
seen.add(userId);
const set = sockets.get(userId);
if (!set || set.size === 0) continue;
for (const socket of set) {
try {
if (socket.readyState === socket.OPEN) socket.send(payload);
} catch {
// best-effort
}
}
}
}
export function snapshot() {
const result = {};
for (const [userId, set] of sockets) result[userId] = set.size;
return result;
}

Powered by TurnKey Linux.