WebSocket Realtime Secure Engine
1. System Architecture & Prerequisites
- Node.js >= 18 LTS (runtime), npm >= 9, TypeScript >= 5.5,
tsx@^4.16.2 as the dev runner.
- Runtime deps:
socket.io@^4.7.5, socket.io-client@^4.7.5, jsonwebtoken@^9.0.2, redis@^4.7.0 (optional adapter), @socket.io/redis-adapter@^8.3.0 (optional horizontal scale), dotenv@^16.4.5.
- Dev deps:
@types/node@^20.14.9, @types/jsonwebtoken@^9.0.6, typescript@^5.5.2.
- Topologies: single-instance (in-memory sequence + room routing) and horizontally scaled (Redis adapter broadcasts
message/ping events across nodes; sequence log stays per-node unless moved to Redis Streams).
2. Input/Output Data Contracts
Handshake authentication
- Auth input:
socket.handshake.auth.token — a JWT access token (claims sub, role, email). Optional fallback header: x-access-token.
- Failure output:
next(new Error('unauthorized: <reason>')); the client receives a connect_error event and must not be connected.
Client → server events
room:join → { channel: string } → server emits room:joined { channel, resumeSeq } or room:join:error { code, message }.
room:leave → { channel: string }.
room:list → server emits room:list:result { rooms: string[] } filtered by role.
chat:send → { channel: string, body: string(1..2000) } → broadcast message { seq, channel, userId, payload: { from, body }, ts } or chat:send:error.
resync:since → { channel: string, since: number } → server emits resync:events { channel, events: StoredMessage[], currentSeq }.
delivered → { seq: number } — advances the per-user delivered-sequence watermark.
pong → any payload — marks the socket alive for the heartbeat monitor.
Server → client events
ping { at: number } every 25s (app-level liveness frame; the client must answer pong or it is disconnected).
message, resync:events, room:joined, room:list:result, and error events as defined above.
Output artifact paths
skills/websocket-realtime-secure-engine/package.json, tsconfig.json, .env.example
src/auth.ts, src/security.ts, src/server.ts, src/client.ts
3. Production Reference Implementation
# .env.example
PORT=8080
CORS_ORIGIN=http://localhost:3000
JWT_ACCESS_SECRET=change-me-at-least-32-chars-long
REDIS_URL=redis://127.0.0.1:6379
REDIS_ADAPTER=false
CHANNEL=channel:public
ACCESS_TOKEN=
// package.json
{
"name": "websocket-realtime-secure-engine",
"version": "1.0.0",
"private": true,
"type": "module",
"scripts": {
"server": "tsx src/server.ts",
"client": "tsx src/client.ts"
},
"dependencies": {
"@socket.io/redis-adapter": "^8.3.0",
"dotenv": "^16.4.5",
"jsonwebtoken": "^9.0.2",
"redis": "^4.7.0",
"socket.io": "^4.7.5",
"socket.io-client": "^4.7.5"
},
"devDependencies": {
"@types/jsonwebtoken": "^9.0.6",
"@types/node": "^20.14.9",
"tsx": "^4.16.2",
"typescript": "^5.5.2"
}
}
// tsconfig.json
{
"compilerOptions": {
"target": "ES2022",
"module": "ESNext",
"moduleResolution": "Bundler",
"strict": true,
"esModuleInterop": true,
"skipLibCheck": true,
"forceConsistentCasingInFileNames": true,
"types": ["node"]
},
"include": ["src/**/*.ts"]
}
// src/auth.ts
import jwt from 'jsonwebtoken';
export type Role = 'admin' | 'user' | 'guest';
export interface AuthUser {
userId: string;
role: Role;
email: string;
}
export interface AuthSuccess {
ok: true;
user: AuthUser;
}
export interface AuthFailure {
ok: false;
error: string;
}
export type AuthResult = AuthSuccess | AuthFailure;
export function verifyJwt(token: string): AuthResult {
const secret = process.env.JWT_ACCESS_SECRET;
if (!secret) {
return { ok: false, error: 'server JWT_ACCESS_SECRET not configured' };
}
try {
const payload = jwt.verify(token, secret) as {
sub?: string;
email?: string;
role?: Role;
type?: string;
};
if (typeof payload.sub !== 'string') {
return { ok: false, error: 'missing subject' };
}
return {
ok: true,
user: {
userId: payload.sub,
role: payload.role ?? 'user',
email: payload.email ?? ''
}
};
} catch (err) {
return { ok: false, error: err instanceof Error ? err.message : 'invalid token' };
}
}
// src/security.ts
export type Role = 'admin' | 'user' | 'guest';
export interface RoomPolicy {
channel: string;
description: string;
allowedRoles: Role[];
}
export const ROOMS: Record<string, RoomPolicy> = {
'channel:public': {
channel: 'channel:public',
description: 'Everyone',
allowedRoles: ['admin', 'user', 'guest']
},
'channel:premium': {
channel: 'channel:premium',
description: 'Admin and paying users',
allowedRoles: ['admin', 'user']
},
'channel:admin': {
channel: 'channel:admin',
description: 'Staff only',
allowedRoles: ['admin']
}
};
export class RealtimeError extends Error {
public readonly statusCode: number;
public readonly code: string;
constructor(statusCode: number, code: string, message: string) {
super(message);
this.statusCode = statusCode;
this.code = code;
}
}
export function canJoin(channel: string, role: Role): boolean {
const room = ROOMS[channel];
if (!room) return false;
return room.allowedRoles.includes(role);
}
export function assertCanJoin(channel: string, role: Role): void {
if (!canJoin(channel, role)) {
throw new RealtimeError(403, 'FORBIDDEN_CHANNEL', `Role '${role}' cannot join '${channel}'`);
}
}
// src/server.ts
import http from 'node:http';
import 'dotenv/config';
import type { Server as SocketServer, Socket } from 'socket.io';
import { Server } from 'socket.io';
import { createClient } from 'redis';
import { createAdapter } from '@socket.io/redis-adapter';
import { verifyJwt, type AuthUser } from './auth';
import { assertCanJoin, ROOMS, RealtimeError } from './security';
const PORT = Number(process.env.PORT ?? 8080);
const PING_INTERVAL_MS = 25_000;
const PING_TIMEOUT_MS = 10_000;
const MESSAGE_LOG_CAP = 10_000;
interface StoredMessage {
seq: number;
channel: string;
userId: string;
payload: unknown;
ts: number;
}
const messageLog = new Map<string, StoredMessage[]>();
const seqByUser = new Map<string, number>();
const sockets = new Map<string, Socket>();
const alive = new Map<string, boolean>();
let globalSeq = 0;
const httpServer = http.createServer((_req, res) => {
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ success: true, data: { service: 'realtime-engine' } }));
});
const io: SocketServer = new Server(httpServer, {
cors: { origin: process.env.CORS_ORIGIN ?? '*', credentials: true },
transports: ['websocket'],
pingInterval: PING_INTERVAL_MS,
pingTimeout: PING_TIMEOUT_MS
});
io.use((socket, next) => {
const token = (
socket.handshake.auth?.token ??
socket.handshake.headers['x-access-token']
) as string | undefined;
if (!token) {
return next(new Error('unauthorized: missing token'));
}
const result = verifyJwt(token);
if (!result.ok) {
return next(new Error(`unauthorized: ${result.error}`));
}
socket.data.user = result.user;
next();
});
function currentSeq(userId: string): number {
return seqByUser.get(userId) ?? 0;
}
function storeMessage(channel: string, userId: string, payload: unknown): StoredMessage {
globalSeq += 1;
const msg: StoredMessage = { seq: globalSeq, channel, userId, payload, ts: Date.now() };
const log = messageLog.get(channel) ?? [];
log.push(msg);
if (log.length > MESSAGE_LOG_CAP) log.shift();
messageLog.set(channel, log);
return msg;
}
function publish(channel: string, userId: string, payload: unknown): void {
const msg = storeMessage(channel, userId, payload);
io.to(channel).emit('message', msg);
}
function handleRoomJoin(socket: Socket, user: AuthUser, channel: string): void {
try {
assertCanJoin(channel, user.role);
} catch (err) {
const e = err instanceof RealtimeError ? err : new RealtimeError(500, 'INTERNAL', 'Join failed');
socket.emit('room:join:error', { code: e.code, message: e.message });
return;
}
void socket.join(channel);
socket.emit('room:joined', { channel, resumeSeq: currentSeq(user.userId) });
}
function handleResync(socket: Socket, user: AuthUser, channel: string, since: number): void {
const log = messageLog.get(channel) ?? [];
const events = log.filter((m) => m.seq > since && m.seq <= currentSeq(user.userId));
socket.emit('resync:events', { channel, events, currentSeq: currentSeq(user.userId) });
}
const heartbeat = setInterval(() => {
const now = Date.now();
for (const [id, socket] of sockets) {
if (alive.get(id) === false) {
console.log(`[heartbeat] socket=${id} missed pong, disconnecting`);
socket.disconnect(true);
sockets.delete(id);
alive.delete(id);
continue;
}
alive.set(id, false);
socket.emit('ping', { at: now });
}
}, PING_INTERVAL_MS);
io.on('connection', (socket) => {
const user = socket.data.user as AuthUser;
sockets.set(socket.id, socket);
alive.set(socket.id, true);
console.log(`[conn] socket=${socket.id} userId=${user.userId} role=${user.role}`);
socket.on('pong', () => {
alive.set(socket.id, true);
});
socket.on('room:join', (payload: { channel?: string }) => {
handleRoomJoin(socket, user, payload?.channel ?? 'channel:public');
});
socket.on('room:leave', (payload: { channel?: string }) => {
if (payload?.channel) {
void socket.leave(payload.channel);
}
});
socket.on('room:list', () => {
const rooms = Object.values(ROOMS)
.filter((room) => room.allowedRoles.includes(user.role))
.map((room) => room.channel);
socket.emit('room:list:result', { rooms });
});
socket.on('chat:send', (payload: { channel: string; body: string }) => {
try {
assertCanJoin(payload.channel, user.role);
if (typeof payload.body !== 'string' || payload.body.length === 0 || payload.body.length > 2000) {
throw new RealtimeError(400, 'INVALID_BODY', 'Message body must be 1-2000 characters');
}
publish(payload.channel, user.userId, { from: user.userId, body: payload.body });
} catch (err) {
const e = err instanceof RealtimeError ? err : new RealtimeError(500, 'INTERNAL', 'Publish failed');
socket.emit('chat:send:error', { code: e.code, message: e.message });
}
});
socket.on('resync:since', (payload: { channel: string; since?: number }) => {
const since = Number.isSafeInteger(payload?.since) ? (payload.since as number) : 0;
handleResync(socket, user, payload?.channel ?? 'channel:public', since);
});
socket.on('delivered', (payload: { seq?: number }) => {
const seq = payload?.seq;
if (typeof seq === 'number' && seq > currentSeq(user.userId)) {
seqByUser.set(user.userId, seq);
}
});
socket.on('disconnect', (reason) => {
sockets.delete(socket.id);
alive.delete(socket.id);
console.log(`[disc] socket=${socket.id} reason=${reason}`);
});
});
async function attachRedisAdapter(): Promise<void> {
if (process.env.REDIS_ADAPTER !== 'true') {
return;
}
const pub = createClient({ url: process.env.REDIS_URL ?? 'redis://127.0.0.1:6379' });
const sub = pub.duplicate();
await Promise.all([pub.connect(), sub.connect()]);
io.adapter(createAdapter(pub, sub));
console.log('[redis-adapter] attached — broadcasts fan out across nodes');
}
function shutdown(): void {
clearInterval(heartbeat);
io.close();
httpServer.close(() => process.exit(0));
setTimeout(() => process.exit(1), 5_000).unref();
}
process.on('SIGTERM', shutdown);
process.on('SIGINT', shutdown);
async function bootstrap(): Promise<void> {
await attachRedisAdapter();
httpServer.listen(PORT, () => {
console.log(`[realtime] engine on ws://localhost:${PORT}`);
});
}
void bootstrap();
// src/client.ts
import 'dotenv/config';
import { io, type Socket } from 'socket.io-client';
const url = process.env.WS_URL ?? 'http://localhost:8080';
const channel = process.env.CHANNEL ?? 'channel:public';
const token = process.env.ACCESS_TOKEN;
function start(): void {
if (!token) {
console.error('[client] ACCESS_TOKEN env var is required');
process.exit(1);
}
let lastSeenSeq = 0;
const socket: Socket = io(url, {
auth: { token },
transports: ['websocket'],
reconnection: true,
reconnectionAttempts: Infinity,
reconnectionDelay: 1_000,
reconnectionDelayMax: 8_000,
randomizationFactor: 0.5,
timeout: 20_000
});
socket.on('connect', () => {
console.log(`[client] connected socket=${socket.id}`);
socket.emit('room:join', { channel });
if (lastSeenSeq > 0) {
socket.emit('resync:since', { channel, since: lastSeenSeq });
}
});
socket.on('room:joined', (payload: { channel: string; resumeSeq: number }) => {
lastSeenSeq = payload.resumeSeq;
console.log(`[client] joined ${payload.channel} resumeSeq=${lastSeenSeq}`);
socket.emit('chat:send', { channel, body: `hello from ${socket.id}` });
});
socket.on('room:join:error', (err: { code: string; message: string }) => {
console.error(`[client] join rejected code=${err.code} message=${err.message}`);
});
socket.on('message', (msg: { seq: number; channel: string; payload: { from: string; body: string }; ts: number }) => {
if (msg.seq > lastSeenSeq) lastSeenSeq = msg.seq;
console.log(`[client] seq=${msg.seq} from=${msg.payload.from}: ${msg.payload.body}`);
});
socket.on('resync:events', (payload: {
events: Array<{ seq: number; payload: { from: string; body: string } }>;
currentSeq: number;
}) => {
for (const evt of payload.events) {
if (evt.seq > lastSeenSeq) lastSeenSeq = evt.seq;
console.log(`[client] replay seq=${evt.seq} from=${evt.payload.from}: ${evt.payload.body}`);
}
if (payload.currentSeq > lastSeenSeq) lastSeenSeq = payload.currentSeq;
socket.emit('delivered', { seq: lastSeenSeq });
});
socket.on('chat:send:error', (err: { code: string; message: string }) => {
console.error(`[client] send failed code=${err.code} message=${err.message}`);
});
socket.on('ping', () => {
socket.emit('pong', { at: Date.now() });
});
socket.on('connect_error', (err) => {
console.error(`[client] connect_error ${err.message}`);
});
socket.on('disconnect', (reason) => {
console.warn(`[client] disconnected reason=${reason}`);
});
}
start();
4. Execution Protocol & Step-by-Step Workflow
- Scaffold: create the project directory and copy all files from section 3 (package.json, tsconfig.json,
.env.example, src/*).
- Configure: copy
.env.example to .env; set JWT_ACCESS_SECRET (min 32 chars) and generate at least one test token with jsonwebtoken or jose: node -e "console.log(require('jsonwebtoken').sign({sub:'user-1',role:'admin',email:'admin@example.com'},'<secret>',{expiresIn:'15m'}))".
- Install: run
npm install.
- Run server: execute
npm run server; a plain-HTTP health probe is served at http://localhost:8080 returning {"service":"realtime-engine"}.
- Run client A: execute
ACCESS_TOKEN=<token> npm run client with ROLE-appropriate token; observe connected, joined, the broadcast message, and periodic ping/pong liveness in the logs.
- Run client B in a second terminal to confirm fan-out of
chat:send to every member of the room on that node.
- Verify authorization: join
channel:admin with a role: 'user' token and confirm room:join:error FORBIDDEN_CHANNEL (403 semantics) is emitted and the client is not admitted.
- Verify heartbeat: stop client A abruptly; within ~35s (
25s ping interval + 10s timeout) the server logs missed pong, disconnecting and cleans up the socket map.
- Verify resync: notify client A is offline, have client B send several messages, then restart client A — it reconnects, re-joins, requests
resync:since from its lastSeenSeq, replays the missed events, and thresholds currentSeq.
- Horizontal scale (optional): set
REDIS_ADAPTER=true, run two server instances, and confirm a chat:send on node 1 reaches a client connected to node 2.
# Example authorization matrix exercised by the guards above
# channel | admin | user | guest
# public | X | X | X
# premium | X | X | -
# admin | X | - | -
5. Edge Cases & Error Handling
- Unauthenticated or expired-token handshakes are rejected in the
io.use middleware with an unauthorized: <reason> connect_error; expired JWTs (via exp claim) fail signature/expiry checks so idle connections cannot linger past the token TTL without re-authenticating.
- The heartbeat monitor uses a check-and-set
alive flag per socket: each tick flips the flag to false and emits ping, and any socket still false at the next tick is force-disconnected and evicted from both maps; the built-in engine pingTimeout (10s) additionally terminates stalled transports, so a dead peer is recycled in ~35s total.
- Channel authorization is enforced at every entry point (
room:join AND chat:send) — joining once does not grant later send rights, and unknown channel names fail canJoin so a typo never falls back to a public broadcast.
- Message size/type guards (
INVALID_BODY) reject non-string or oversized payloads before publishing so the per-user sequence log cannot be poisoned by malformed frames.
- Reconnect reconciliation is duplicate-safe on the client: it only processes
resync:events with seq > lastSeenSeq, idempotently replays, and acknowledges with delivered so the per-user watermark advances monotonically; the ring-buffer cap (MESSAGE_LOG_CAP) drops the oldest events rather than the newest to protect low-latency delivery.
forgotten clients that never ack delivered are bounded by the ring-buffer cap on the server; a Redis Streams-backed journal (consumer groups + XACK) is the documented replacement when exactly-once delivery across restarts is required.
- Redis adapter startup failure is fatal (the process exits) so a half-fan-out production split is impossible; when
REDIS_ADAPTER is unset/false the single-node in-memory path needs no Redis at all, keeping local dev dependency-free.
- Graceful shutdown clears the heartbeat interval, closes the socket.io engine, and drains the HTTP server before exit; a hard-exit timer (unref'd) prevents hangs on stubborn connections.
1---2name: websocket-realtime-secure-engine3description: Scaffolds high-throughput, bidirectional real-time communication layers over WebSockets or Server-Sent Events (SSE) with robust reconnection and authentication controls.4---56# WebSocket Realtime Secure Engine78## 1. System Architecture & Prerequisites910- Node.js >= 18 LTS (runtime), npm >= 9, TypeScript >= 5.5, `tsx@^4.16.2` as the dev runner.11- Runtime deps: `socket.io@^4.7.5`, `socket.io-client@^4.7.5`, `jsonwebtoken@^9.0.2`, `redis@^4.7.0` (optional adapter), `@socket.io/redis-adapter@^8.3.0` (optional horizontal scale), `dotenv@^16.4.5`.12- Dev deps: `@types/node@^20.14.9`, `@types/jsonwebtoken@^9.0.6`, `typescript@^5.5.2`.13- Topologies: single-instance (in-memory sequence + room routing) and horizontally scaled (Redis adapter broadcasts `message`/`ping` events across nodes; sequence log stays per-node unless moved to Redis Streams).1415## 2. Input/Output Data Contracts1617### Handshake authentication1819- Auth input: `socket.handshake.auth.token` — a JWT access token (claims `sub`, `role`, `email`). Optional fallback header: `x-access-token`.20- Failure output: `next(new Error('unauthorized: <reason>'))`; the client receives a `connect_error` event and must not be connected.2122### Client → server events2324- `room:join` → `{ channel: string }` → server emits `room:joined` `{ channel, resumeSeq }` or `room:join:error` `{ code, message }`.25- `room:leave` → `{ channel: string }`.26- `room:list` → server emits `room:list:result` `{ rooms: string[] }` filtered by role.27- `chat:send` → `{ channel: string, body: string(1..2000) }` → broadcast `message` `{ seq, channel, userId, payload: { from, body }, ts }` or `chat:send:error`.28- `resync:since` → `{ channel: string, since: number }` → server emits `resync:events` `{ channel, events: StoredMessage[], currentSeq }`.29- `delivered` → `{ seq: number }` — advances the per-user delivered-sequence watermark.30- `pong` → any payload — marks the socket alive for the heartbeat monitor.3132### Server → client events3334- `ping` `{ at: number }` every 25s (app-level liveness frame; the client must answer `pong` or it is disconnected).35- `message`, `resync:events`, `room:joined`, `room:list:result`, and error events as defined above.3637### Output artifact paths3839- `skills/websocket-realtime-secure-engine/package.json`, `tsconfig.json`, `.env.example`40- `src/auth.ts`, `src/security.ts`, `src/server.ts`, `src/client.ts`4142## 3. Production Reference Implementation4344```text45# .env.example46PORT=808047CORS_ORIGIN=http://localhost:300048JWT_ACCESS_SECRET=change-me-at-least-32-chars-long49REDIS_URL=redis://127.0.0.1:637950REDIS_ADAPTER=false51CHANNEL=channel:public52ACCESS_TOKEN=53```5455```json56// package.json57{58 "name": "websocket-realtime-secure-engine",59 "version": "1.0.0",60 "private": true,61 "type": "module",62 "scripts": {63 "server": "tsx src/server.ts",64 "client": "tsx src/client.ts"65 },66 "dependencies": {67 "@socket.io/redis-adapter": "^8.3.0",68 "dotenv": "^16.4.5",69 "jsonwebtoken": "^9.0.2",70 "redis": "^4.7.0",71 "socket.io": "^4.7.5",72 "socket.io-client": "^4.7.5"73 },74 "devDependencies": {75 "@types/jsonwebtoken": "^9.0.6",76 "@types/node": "^20.14.9",77 "tsx": "^4.16.2",78 "typescript": "^5.5.2"79 }80}81```8283```json84// tsconfig.json85{86 "compilerOptions": {87 "target": "ES2022",88 "module": "ESNext",89 "moduleResolution": "Bundler",90 "strict": true,91 "esModuleInterop": true,92 "skipLibCheck": true,93 "forceConsistentCasingInFileNames": true,94 "types": ["node"]95 },96 "include": ["src/**/*.ts"]97}98```99100```typescript101// src/auth.ts102import jwt from 'jsonwebtoken';103104export type Role = 'admin' | 'user' | 'guest';105106export interface AuthUser {107 userId: string;108 role: Role;109 email: string;110}111112export interface AuthSuccess {113 ok: true;114 user: AuthUser;115}116117export interface AuthFailure {118 ok: false;119 error: string;120}121122export type AuthResult = AuthSuccess | AuthFailure;123124export function verifyJwt(token: string): AuthResult {125 const secret = process.env.JWT_ACCESS_SECRET;126 if (!secret) {127 return { ok: false, error: 'server JWT_ACCESS_SECRET not configured' };128 }129 try {130 const payload = jwt.verify(token, secret) as {131 sub?: string;132 email?: string;133 role?: Role;134 type?: string;135 };136 if (typeof payload.sub !== 'string') {137 return { ok: false, error: 'missing subject' };138 }139 return {140 ok: true,141 user: {142 userId: payload.sub,143 role: payload.role ?? 'user',144 email: payload.email ?? ''145 }146 };147 } catch (err) {148 return { ok: false, error: err instanceof Error ? err.message : 'invalid token' };149 }150}151```152153```typescript154// src/security.ts155export type Role = 'admin' | 'user' | 'guest';156157export interface RoomPolicy {158 channel: string;159 description: string;160 allowedRoles: Role[];161}162163export const ROOMS: Record<string, RoomPolicy> = {164 'channel:public': {165 channel: 'channel:public',166 description: 'Everyone',167 allowedRoles: ['admin', 'user', 'guest']168 },169 'channel:premium': {170 channel: 'channel:premium',171 description: 'Admin and paying users',172 allowedRoles: ['admin', 'user']173 },174 'channel:admin': {175 channel: 'channel:admin',176 description: 'Staff only',177 allowedRoles: ['admin']178 }179};180181export class RealtimeError extends Error {182 public readonly statusCode: number;183 public readonly code: string;184185 constructor(statusCode: number, code: string, message: string) {186 super(message);187 this.statusCode = statusCode;188 this.code = code;189 }190}191192export function canJoin(channel: string, role: Role): boolean {193 const room = ROOMS[channel];194 if (!room) return false;195 return room.allowedRoles.includes(role);196}197198export function assertCanJoin(channel: string, role: Role): void {199 if (!canJoin(channel, role)) {200 throw new RealtimeError(403, 'FORBIDDEN_CHANNEL', `Role '${role}' cannot join '${channel}'`);201 }202}203```204205```typescript206// src/server.ts207import http from 'node:http';208import 'dotenv/config';209import type { Server as SocketServer, Socket } from 'socket.io';210import { Server } from 'socket.io';211import { createClient } from 'redis';212import { createAdapter } from '@socket.io/redis-adapter';213import { verifyJwt, type AuthUser } from './auth';214import { assertCanJoin, ROOMS, RealtimeError } from './security';215216const PORT = Number(process.env.PORT ?? 8080);217const PING_INTERVAL_MS = 25_000;218const PING_TIMEOUT_MS = 10_000;219const MESSAGE_LOG_CAP = 10_000;220221interface StoredMessage {222 seq: number;223 channel: string;224 userId: string;225 payload: unknown;226 ts: number;227}228229const messageLog = new Map<string, StoredMessage[]>();230const seqByUser = new Map<string, number>();231const sockets = new Map<string, Socket>();232const alive = new Map<string, boolean>();233let globalSeq = 0;234235const httpServer = http.createServer((_req, res) => {236 res.writeHead(200, { 'Content-Type': 'application/json' });237 res.end(JSON.stringify({ success: true, data: { service: 'realtime-engine' } }));238});239240const io: SocketServer = new Server(httpServer, {241 cors: { origin: process.env.CORS_ORIGIN ?? '*', credentials: true },242 transports: ['websocket'],243 pingInterval: PING_INTERVAL_MS,244 pingTimeout: PING_TIMEOUT_MS245});246247io.use((socket, next) => {248 const token = (249 socket.handshake.auth?.token ??250 socket.handshake.headers['x-access-token']251 ) as string | undefined;252 if (!token) {253 return next(new Error('unauthorized: missing token'));254 }255 const result = verifyJwt(token);256 if (!result.ok) {257 return next(new Error(`unauthorized: ${result.error}`));258 }259 socket.data.user = result.user;260 next();261});262263function currentSeq(userId: string): number {264 return seqByUser.get(userId) ?? 0;265}266267function storeMessage(channel: string, userId: string, payload: unknown): StoredMessage {268 globalSeq += 1;269 const msg: StoredMessage = { seq: globalSeq, channel, userId, payload, ts: Date.now() };270 const log = messageLog.get(channel) ?? [];271 log.push(msg);272 if (log.length > MESSAGE_LOG_CAP) log.shift();273 messageLog.set(channel, log);274 return msg;275}276277function publish(channel: string, userId: string, payload: unknown): void {278 const msg = storeMessage(channel, userId, payload);279 io.to(channel).emit('message', msg);280}281282function handleRoomJoin(socket: Socket, user: AuthUser, channel: string): void {283 try {284 assertCanJoin(channel, user.role);285 } catch (err) {286 const e = err instanceof RealtimeError ? err : new RealtimeError(500, 'INTERNAL', 'Join failed');287 socket.emit('room:join:error', { code: e.code, message: e.message });288 return;289 }290 void socket.join(channel);291 socket.emit('room:joined', { channel, resumeSeq: currentSeq(user.userId) });292}293294function handleResync(socket: Socket, user: AuthUser, channel: string, since: number): void {295 const log = messageLog.get(channel) ?? [];296 const events = log.filter((m) => m.seq > since && m.seq <= currentSeq(user.userId));297 socket.emit('resync:events', { channel, events, currentSeq: currentSeq(user.userId) });298}299300const heartbeat = setInterval(() => {301 const now = Date.now();302 for (const [id, socket] of sockets) {303 if (alive.get(id) === false) {304 console.log(`[heartbeat] socket=${id} missed pong, disconnecting`);305 socket.disconnect(true);306 sockets.delete(id);307 alive.delete(id);308 continue;309 }310 alive.set(id, false);311 socket.emit('ping', { at: now });312 }313}, PING_INTERVAL_MS);314315io.on('connection', (socket) => {316 const user = socket.data.user as AuthUser;317 sockets.set(socket.id, socket);318 alive.set(socket.id, true);319 console.log(`[conn] socket=${socket.id} userId=${user.userId} role=${user.role}`);320321 socket.on('pong', () => {322 alive.set(socket.id, true);323 });324325 socket.on('room:join', (payload: { channel?: string }) => {326 handleRoomJoin(socket, user, payload?.channel ?? 'channel:public');327 });328329 socket.on('room:leave', (payload: { channel?: string }) => {330 if (payload?.channel) {331 void socket.leave(payload.channel);332 }333 });334335 socket.on('room:list', () => {336 const rooms = Object.values(ROOMS)337 .filter((room) => room.allowedRoles.includes(user.role))338 .map((room) => room.channel);339 socket.emit('room:list:result', { rooms });340 });341342 socket.on('chat:send', (payload: { channel: string; body: string }) => {343 try {344 assertCanJoin(payload.channel, user.role);345 if (typeof payload.body !== 'string' || payload.body.length === 0 || payload.body.length > 2000) {346 throw new RealtimeError(400, 'INVALID_BODY', 'Message body must be 1-2000 characters');347 }348 publish(payload.channel, user.userId, { from: user.userId, body: payload.body });349 } catch (err) {350 const e = err instanceof RealtimeError ? err : new RealtimeError(500, 'INTERNAL', 'Publish failed');351 socket.emit('chat:send:error', { code: e.code, message: e.message });352 }353 });354355 socket.on('resync:since', (payload: { channel: string; since?: number }) => {356 const since = Number.isSafeInteger(payload?.since) ? (payload.since as number) : 0;357 handleResync(socket, user, payload?.channel ?? 'channel:public', since);358 });359360 socket.on('delivered', (payload: { seq?: number }) => {361 const seq = payload?.seq;362 if (typeof seq === 'number' && seq > currentSeq(user.userId)) {363 seqByUser.set(user.userId, seq);364 }365 });366367 socket.on('disconnect', (reason) => {368 sockets.delete(socket.id);369 alive.delete(socket.id);370 console.log(`[disc] socket=${socket.id} reason=${reason}`);371 });372});373374async function attachRedisAdapter(): Promise<void> {375 if (process.env.REDIS_ADAPTER !== 'true') {376 return;377 }378 const pub = createClient({ url: process.env.REDIS_URL ?? 'redis://127.0.0.1:6379' });379 const sub = pub.duplicate();380 await Promise.all([pub.connect(), sub.connect()]);381 io.adapter(createAdapter(pub, sub));382 console.log('[redis-adapter] attached — broadcasts fan out across nodes');383}384385function shutdown(): void {386 clearInterval(heartbeat);387 io.close();388 httpServer.close(() => process.exit(0));389 setTimeout(() => process.exit(1), 5_000).unref();390}391392process.on('SIGTERM', shutdown);393process.on('SIGINT', shutdown);394395async function bootstrap(): Promise<void> {396 await attachRedisAdapter();397 httpServer.listen(PORT, () => {398 console.log(`[realtime] engine on ws://localhost:${PORT}`);399 });400}401402void bootstrap();403```404405```typescript406// src/client.ts407import 'dotenv/config';408import { io, type Socket } from 'socket.io-client';409410const url = process.env.WS_URL ?? 'http://localhost:8080';411const channel = process.env.CHANNEL ?? 'channel:public';412const token = process.env.ACCESS_TOKEN;413414function start(): void {415 if (!token) {416 console.error('[client] ACCESS_TOKEN env var is required');417 process.exit(1);418 }419420 let lastSeenSeq = 0;421422 const socket: Socket = io(url, {423 auth: { token },424 transports: ['websocket'],425 reconnection: true,426 reconnectionAttempts: Infinity,427 reconnectionDelay: 1_000,428 reconnectionDelayMax: 8_000,429 randomizationFactor: 0.5,430 timeout: 20_000431 });432433 socket.on('connect', () => {434 console.log(`[client] connected socket=${socket.id}`);435 socket.emit('room:join', { channel });436 if (lastSeenSeq > 0) {437 socket.emit('resync:since', { channel, since: lastSeenSeq });438 }439 });440441 socket.on('room:joined', (payload: { channel: string; resumeSeq: number }) => {442 lastSeenSeq = payload.resumeSeq;443 console.log(`[client] joined ${payload.channel} resumeSeq=${lastSeenSeq}`);444 socket.emit('chat:send', { channel, body: `hello from ${socket.id}` });445 });446447 socket.on('room:join:error', (err: { code: string; message: string }) => {448 console.error(`[client] join rejected code=${err.code} message=${err.message}`);449 });450451 socket.on('message', (msg: { seq: number; channel: string; payload: { from: string; body: string }; ts: number }) => {452 if (msg.seq > lastSeenSeq) lastSeenSeq = msg.seq;453 console.log(`[client] seq=${msg.seq} from=${msg.payload.from}: ${msg.payload.body}`);454 });455456 socket.on('resync:events', (payload: {457 events: Array<{ seq: number; payload: { from: string; body: string } }>;458 currentSeq: number;459 }) => {460 for (const evt of payload.events) {461 if (evt.seq > lastSeenSeq) lastSeenSeq = evt.seq;462 console.log(`[client] replay seq=${evt.seq} from=${evt.payload.from}: ${evt.payload.body}`);463 }464 if (payload.currentSeq > lastSeenSeq) lastSeenSeq = payload.currentSeq;465 socket.emit('delivered', { seq: lastSeenSeq });466 });467468 socket.on('chat:send:error', (err: { code: string; message: string }) => {469 console.error(`[client] send failed code=${err.code} message=${err.message}`);470 });471472 socket.on('ping', () => {473 socket.emit('pong', { at: Date.now() });474 });475476 socket.on('connect_error', (err) => {477 console.error(`[client] connect_error ${err.message}`);478 });479480 socket.on('disconnect', (reason) => {481 console.warn(`[client] disconnected reason=${reason}`);482 });483}484485start();486```487488## 4. Execution Protocol & Step-by-Step Workflow4894901. Scaffold: create the project directory and copy all files from section 3 (package.json, tsconfig.json, `.env.example`, `src/*`).4912. Configure: copy `.env.example` to `.env`; set `JWT_ACCESS_SECRET` (min 32 chars) and generate at least one test token with `jsonwebtoken` or `jose`: `node -e "console.log(require('jsonwebtoken').sign({sub:'user-1',role:'admin',email:'admin@example.com'},'<secret>',{expiresIn:'15m'}))"`.4923. Install: run `npm install`.4934. Run server: execute `npm run server`; a plain-HTTP health probe is served at `http://localhost:8080` returning `{"service":"realtime-engine"}`.4945. Run client A: execute `ACCESS_TOKEN=<token> npm run client` with `ROLE`-appropriate token; observe `connected`, `joined`, the broadcast `message`, and periodic `ping`/`pong` liveness in the logs.4956. Run client B in a second terminal to confirm fan-out of `chat:send` to every member of the room on that node.4967. Verify authorization: join `channel:admin` with a `role: 'user'` token and confirm `room:join:error` `FORBIDDEN_CHANNEL` (403 semantics) is emitted and the client is not admitted.4978. Verify heartbeat: stop client A abruptly; within ~35s (`25s` ping interval + `10s` timeout) the server logs `missed pong, disconnecting` and cleans up the socket map.4989. Verify resync: notify client A is offline, have client B send several messages, then restart client A — it reconnects, re-joins, requests `resync:since` from its `lastSeenSeq`, replays the missed events, and thresholds `currentSeq`.49910. Horizontal scale (optional): set `REDIS_ADAPTER=true`, run two server instances, and confirm a `chat:send` on node 1 reaches a client connected to node 2.500501```text502# Example authorization matrix exercised by the guards above503# channel | admin | user | guest504# public | X | X | X505# premium | X | X | -506# admin | X | - | -507```508509## 5. Edge Cases & Error Handling510511- Unauthenticated or expired-token handshakes are rejected in the `io.use` middleware with an `unauthorized: <reason>` `connect_error`; expired JWTs (via `exp` claim) fail signature/expiry checks so idle connections cannot linger past the token TTL without re-authenticating.512- The heartbeat monitor uses a check-and-set `alive` flag per socket: each tick flips the flag to `false` and emits `ping`, and any socket still `false` at the next tick is force-disconnected and evicted from both maps; the built-in engine `pingTimeout` (10s) additionally terminates stalled transports, so a dead peer is recycled in ~35s total.513- Channel authorization is enforced at every entry point (`room:join` AND `chat:send`) — joining once does not grant later send rights, and unknown channel names fail `canJoin` so a typo never falls back to a public broadcast.514- Message size/type guards (`INVALID_BODY`) reject non-string or oversized payloads before publishing so the per-user sequence log cannot be poisoned by malformed frames.515- Reconnect reconciliation is duplicate-safe on the client: it only processes `resync:events` with `seq > lastSeenSeq`, idempotently replays, and acknowledges with `delivered` so the per-user watermark advances monotonically; the ring-buffer cap (`MESSAGE_LOG_CAP`) drops the oldest events rather than the newest to protect low-latency delivery.516- `forgotten` clients that never ack `delivered` are bounded by the ring-buffer cap on the server; a Redis Streams-backed journal (consumer groups + `XACK`) is the documented replacement when exactly-once delivery across restarts is required.517- Redis adapter startup failure is fatal (the process exits) so a half-fan-out production split is impossible; when `REDIS_ADAPTER` is unset/false the single-node in-memory path needs no Redis at all, keeping local dev dependency-free.518- Graceful shutdown clears the heartbeat interval, closes the socket.io engine, and drains the HTTP server before exit; a hard-exit timer (unref'd) prevents hangs on stubborn connections.