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.

123 lines
4.3 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
import { createServer } from 'node:http';
import { Readable } from 'node:stream';
import { WebSocketServer } from 'ws';
import { config } from './env.mjs';
import { currentSession } from './domain.mjs';
import { applyCors, corsPreflight, handleError, parseCookies, writeResponse } from './http.mjs';
import { attach as attachSocket } from './realtime.mjs';
import { route } from './routes.mjs';
const server = createServer(async (nodeRequest, nodeResponse) => {
try {
const request = toWebRequest(nodeRequest);
if (request.method === 'OPTIONS') {
writeResponse(nodeResponse, applyCors(request, corsPreflight(request)));
return;
}
const response = await route(request);
logResponse(request, response);
writeResponse(nodeResponse, applyCors(request, response));
} catch (error) {
const request = toWebRequest(nodeRequest, false);
const response = handleError(error);
logResponse(request, response);
writeResponse(nodeResponse, applyCors(request, response));
}
});
// ── Realtime ─────────────────────────────────────────────────────────────
//
// The WebSocket endpoint at `/api/realtime` authenticates via the same
// `dating_session` cookie the HTTP API uses, then registers the socket
// with the in-memory hub. Endpoint handlers (currently `messagesPost`)
// call `broadcast([recipientId], event)` to push events out without
// the consumer needing to track connections themselves.
//
// `noServer: true` lets us reuse the existing HTTP server for upgrades
// and run the cookie-based auth check before accepting the WebSocket.
const wss = new WebSocketServer({ noServer: true });
server.on('upgrade', async (nodeRequest, socket, head) => {
const url = nodeRequest.url || '/';
if (!url.startsWith('/api/realtime')) {
socket.destroy();
return;
}
// Reject upgrades from origins not in the CORS allowlist; Browsers
// are happy to upgrade cross-origin and the server has to gate it
// explicitly.
const origin = nodeRequest.headers.origin;
if (origin && !config.allowedOrigins.includes(origin)) {
socket.write('HTTP/1.1 403 Forbidden\r\n\r\n');
socket.destroy();
return;
}
const session = await authenticateUpgrade(nodeRequest);
if (session === null) {
socket.write('HTTP/1.1 401 Unauthorized\r\n\r\n');
socket.destroy();
return;
}
wss.handleUpgrade(nodeRequest, socket, head, (ws) => {
attachSocket(session.user.id, ws);
try {
ws.send(JSON.stringify({ type: 'hello', userId: session.user.id }));
} catch {
// initial frame is best-effort; the socket will recover on the next event
}
});
});
async function authenticateUpgrade(nodeRequest) {
const cookies = parseCookies(nodeRequest.headers.cookie);
const token = cookies.get(config.cookieName);
if (!token) return null;
// Reuse the HTTP path's session helper. It expects a `Headers`
// object so we forge one carrying the bearer the helper looks for.
const fakeRequest = new Request('http://internal/realtime', {
headers: { authorization: `Bearer ${token}` }
});
try {
return await currentSession(fakeRequest);
} catch {
return null;
}
}
server.listen(config.port, config.host, () => {
console.log(`Dating server listening on http://${config.host}:${config.port}`);
console.log(`PocketBase: ${config.pocketBaseUrl}`);
});
function toWebRequest(nodeRequest, includeBody = true) {
const protocol = nodeRequest.headers['x-forwarded-proto'] || 'http';
const host = nodeRequest.headers.host || `${config.host}:${config.port}`;
const url = `${protocol}://${host}${nodeRequest.url || '/'}`;
const method = nodeRequest.method || 'GET';
const headers = new Headers();
for (const [name, value] of Object.entries(nodeRequest.headers)) {
if (Array.isArray(value)) {
for (const item of value) headers.append(name, item);
} else if (value != null) {
headers.set(name, String(value));
}
}
const hasBody = includeBody && !['GET', 'HEAD'].includes(method);
return new Request(url, {
method,
headers,
body: hasBody ? Readable.toWeb(nodeRequest) : undefined,
duplex: hasBody ? 'half' : undefined
});
}
function logResponse(request, response) {
if (response.status < 400) return;
const url = new URL(request.url);
const code = response.body?.error?.code || 'http_error';
console.warn(`${request.method} ${url.pathname} -> ${response.status} ${code}`);
}

Powered by TurnKey Linux.