10 · Project — A Real-Time Chat App¶
This project combines Level 3 into a working multi-room chat: an Express app that issues
short-lived tokens and serves a small browser client, and a WebSocket server that
authenticates connections during the upgrade, validates every message, keeps bounded
per-room history, rate-limits each connection, detects dead peers, and shuts down
cleanly. It's tested end to end with Node's built-in WebSocket client.
Features and the lesson each comes from¶
| Feature | Lesson |
|---|---|
| Upgrade-time auth with a signed, expiring token and Origin check | 07 WebSockets, 03 Buffers (HMAC bytes) |
| Messages validated with a Zod discriminated union | L2-04 |
| Rooms, presence, typing, bounded history | 07 |
Slow-consumer protection via bufferedAmount |
02 Streams & backpressure |
| Per-connection token-bucket rate limit | 09 Security |
| Heartbeats with ping/pong | 07 |
| Graceful shutdown with close code 1001 | preview of L4-09 |
Layout¶
chat/
src/protocol.js # message schemas
src/rooms.js # membership, history, broadcast
src/tokens.js # signed expiring tokens
src/server.js # createChatServer(): Express + ws
src/main.js # entry point
public/index.html # minimal browser client
test/chat.test.js
Dependencies: express, ws, zod.
The protocol¶
import { z } from 'zod';
const room = z.string().regex(/^[a-z0-9-]{1,32}$/, 'room names are lowercase letters, digits, dashes');
// Every message a client may send, discriminated by `type`
export const ClientMessage = z.discriminatedUnion('type', [
z.object({ type: z.literal('join'), room }),
z.object({ type: z.literal('leave'), room }),
z.object({ type: z.literal('say'), room, text: z.string().trim().min(1).max(1000) }),
z.object({ type: z.literal('typing'), room }),
]);
A single schema describes every message a client may send. Anything else — wrong type, missing room, a 5 KB "text", a room name with spaces — is rejected before it reaches the room logic.
Room state¶
// In-memory room state for ONE process. Scaling out requires pub/sub (see lesson text).
export class Rooms {
#members = new Map(); // room -> Set<socket>
#history = new Map(); // room -> recent messages (bounded)
constructor({ historySize = 50 } = {}) { this.historySize = historySize; }
join(room, socket) {
if (!this.#members.has(room)) this.#members.set(room, new Set());
this.#members.get(room).add(socket);
return this.#history.get(room) ?? [];
}
leave(room, socket) {
const set = this.#members.get(room);
if (!set) return;
set.delete(socket);
if (set.size === 0) this.#members.delete(room);
}
leaveAll(socket) {
for (const room of [...this.#members.keys()]) this.leave(room, socket);
}
isMember(room, socket) { return this.#members.get(room)?.has(socket) ?? false; }
record(room, message) {
const list = this.#history.get(room) ?? [];
list.push(message);
if (list.length > this.historySize) list.shift();
this.#history.set(room, list);
}
broadcast(room, payload, { except } = {}) {
const data = JSON.stringify(payload); // serialize once
for (const socket of this.#members.get(room) ?? []) {
if (socket === except || socket.readyState !== socket.OPEN) continue;
if (socket.bufferedAmount > 1_000_000) { socket.terminate(); continue; } // slow consumer
socket.send(data);
}
}
count(room) { return this.#members.get(room)?.size ?? 0; }
}
broadcast serializes once, skips sockets that are closing, and terminates any client
with more than ~1 MB queued — a slow consumer must not be allowed to grow server memory
without bound.
Tokens¶
import { createHmac, timingSafeEqual } from 'node:crypto';
// Tiny signed, expiring token: base64url(json).base64url(hmac). Enough for a demo;
// in the Level 2 API you would reuse the JWT helpers instead.
export function issueToken(name, secret, ttlSeconds = 60) {
const body = Buffer.from(JSON.stringify({ name, exp: Date.now() + ttlSeconds * 1000 })).toString('base64url');
const sig = createHmac('sha256', secret).update(body).digest('base64url');
return `${body}.${sig}`;
}
export function verifyToken(token, secret) {
const [body, sig] = String(token).split('.');
if (!body || !sig) return null;
const expected = createHmac('sha256', secret).update(body).digest();
const given = Buffer.from(sig, 'base64url');
if (given.length !== expected.length || !timingSafeEqual(given, expected)) return null;
const payload = JSON.parse(Buffer.from(body, 'base64url').toString());
return payload.exp > Date.now() ? payload : null;
}
The server¶
import express from 'express';
import { createServer } from 'node:http';
import { WebSocketServer } from 'ws';
import { z } from 'zod';
import { ClientMessage } from './protocol.js';
import { Rooms } from './rooms.js';
import { issueToken, verifyToken } from './tokens.js';
export function createChatServer({ secret, allowedOrigins = [], heartbeatMs = 30_000 }) {
const app = express();
app.use(express.json({ limit: '10kb' }));
app.use(express.static(new URL('../public', import.meta.url).pathname));
// "Login": pick a display name, get a short-lived token for the WebSocket upgrade
const Login = z.object({ name: z.string().trim().regex(/^[\w-]{2,20}$/) });
app.post('/session', (req, res) => {
const parsed = Login.safeParse(req.body);
if (!parsed.success) return res.status(400).json({ error: 'name must be 2-20 letters, digits, _ or -' });
res.json({ token: issueToken(parsed.data.name, secret) });
});
const server = createServer(app);
const wss = new WebSocketServer({ noServer: true, maxPayload: 16 * 1024 });
const rooms = new Rooms();
server.on('upgrade', (req, socket, head) => {
const url = new URL(req.url, 'http://localhost');
const origin = req.headers.origin;
const user = url.pathname === '/ws' ? verifyToken(url.searchParams.get('token'), secret) : null;
if (!user || (origin && allowedOrigins.length && !allowedOrigins.includes(origin))) {
socket.write('HTTP/1.1 401 Unauthorized\r\nConnection: close\r\n\r\n');
return socket.destroy();
}
wss.handleUpgrade(req, socket, head, (ws) => wss.emit('connection', ws, user));
});
wss.on('connection', (ws, user) => {
ws.user = user;
ws.isAlive = true;
ws.tokens = 10; // simple per-socket token bucket
ws.on('pong', () => { ws.isAlive = true; });
ws.on('error', () => {}); // 'close' follows; nothing else to do
ws.on('close', () => rooms.leaveAll(ws));
const send = (payload) => ws.send(JSON.stringify(payload));
ws.on('message', (data, isBinary) => {
if (isBinary) return ws.close(1003, 'text frames only');
if (ws.tokens <= 0) return send({ type: 'error', error: 'slow down' });
ws.tokens--;
let msg;
try { msg = ClientMessage.parse(JSON.parse(data.toString())); }
catch { return send({ type: 'error', error: 'invalid message' }); }
switch (msg.type) {
case 'join': {
const history = rooms.join(msg.room, ws);
send({ type: 'joined', room: msg.room, history, online: rooms.count(msg.room) });
rooms.broadcast(msg.room, { type: 'presence', room: msg.room, user: user.name, event: 'join' }, { except: ws });
break;
}
case 'leave':
rooms.leave(msg.room, ws);
rooms.broadcast(msg.room, { type: 'presence', room: msg.room, user: user.name, event: 'leave' });
break;
case 'say': {
if (!rooms.isMember(msg.room, ws)) return send({ type: 'error', error: `join ${msg.room} first` });
const message = { type: 'message', room: msg.room, user: user.name, text: msg.text, at: new Date().toISOString() };
rooms.record(msg.room, message);
rooms.broadcast(msg.room, message);
break;
}
case 'typing':
if (rooms.isMember(msg.room, ws)) rooms.broadcast(msg.room, { type: 'typing', room: msg.room, user: user.name }, { except: ws });
break;
}
});
});
// Refill rate-limit tokens and detect dead connections
const refill = setInterval(() => { for (const ws of wss.clients) ws.tokens = Math.min(10, ws.tokens + 5); }, 1000);
const heartbeat = setInterval(() => {
for (const ws of wss.clients) {
if (!ws.isAlive) { ws.terminate(); continue; }
ws.isAlive = false;
ws.ping();
}
}, heartbeatMs);
function close() {
clearInterval(refill);
clearInterval(heartbeat);
for (const ws of wss.clients) ws.close(1001, 'server shutting down');
return new Promise((resolve) => server.close(resolve));
}
return { server, close };
}
import { createChatServer } from './server.js';
const secret = process.env.CHAT_SECRET;
if (!secret || secret.length < 32) throw new Error('CHAT_SECRET (32+ chars) is required');
const { server, close } = createChatServer({
secret,
allowedOrigins: (process.env.ALLOWED_ORIGINS ?? 'http://localhost:3400').split(','),
});
server.listen(Number(process.env.PORT ?? 3400), () => console.log('chat on http://localhost:3400'));
process.on('SIGTERM', async () => { await close(); process.exit(0); });
process.on('SIGINT', async () => { await close(); process.exit(0); });
Why noServer: true and a manual 'upgrade' handler? It lets us reject bad
connections with a proper HTTP 401 before they become WebSockets, instead of accepting
them and closing immediately.
A minimal browser client¶
<!doctype html>
<html lang="en">
<meta charset="utf-8">
<title>Node chat</title>
<style>
body { font: 16px system-ui, sans-serif; max-width: 40rem; margin: 2rem auto; padding: 0 1rem; }
#log { border: 1px solid #ccc; height: 20rem; overflow-y: auto; padding: .5rem; }
.sys { color: #777; font-style: italic; }
</style>
<form id="login"><input id="name" placeholder="your name" required> <button>Join #general</button></form>
<div id="log" hidden></div>
<form id="say" hidden><input id="text" autocomplete="off" required> <button>Send</button></form>
<script type="module">
const $ = (id) => document.getElementById(id);
const line = (text, cls) => {
const p = document.createElement('p');
p.textContent = text; // textContent, never innerHTML: no XSS
if (cls) p.className = cls;
$('log').append(p);
$('log').scrollTop = $('log').scrollHeight;
};
$('login').onsubmit = async (e) => {
e.preventDefault();
const res = await fetch('/session', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ name: $('name').value }) });
const { token, error } = await res.json();
if (error) return alert(error);
const ws = new WebSocket(`${location.origin.replace('http', 'ws')}/ws?token=${encodeURIComponent(token)}`);
ws.onopen = () => { ws.send(JSON.stringify({ type: 'join', room: 'general' })); $('login').hidden = true; $('log').hidden = $('say').hidden = false; };
ws.onmessage = (e) => {
const m = JSON.parse(e.data);
if (m.type === 'joined') { m.history.forEach((h) => line(`${h.user}: ${h.text}`)); line(`joined #${m.room} (${m.online} online)`, 'sys'); }
if (m.type === 'message') line(`${m.user}: ${m.text}`);
if (m.type === 'presence') line(`${m.user} ${m.event === 'join' ? 'joined' : 'left'}`, 'sys');
if (m.type === 'error') line(`error: ${m.error}`, 'sys');
};
ws.onclose = () => line('disconnected', 'sys');
$('say').onsubmit = (e) => { e.preventDefault(); ws.send(JSON.stringify({ type: 'say', room: 'general', text: $('text').value })); $('text').value = ''; };
};
</script>
</html>
Run it with CHAT_SECRET=<32+ random chars> node src/main.js, open
http://localhost:3400 in two browser windows, and chat. Starting it and checking the
basics from a terminal:
$ curl -s -o /dev/null -w "%{http_code} %{content_type}\n" localhost:3400/
200 text/html; charset=utf-8
$ curl -s -X POST localhost:3400/session -H 'content-type: application/json' -d '{"name":"x"}'
{"error":"name must be 2-20 letters, digits, _ or -"}
Tests¶
import { test, before, after } from 'node:test';
import assert from 'node:assert/strict';
import { once } from 'node:events';
import { createChatServer } from '../src/server.js';
const secret = 'test-secret-with-at-least-32-characters!';
let chat, base;
before(async () => {
chat = createChatServer({ secret });
chat.server.listen(0);
await once(chat.server, 'listening');
base = `localhost:${chat.server.address().port}`;
});
after(() => chat.close());
async function login(name) {
const res = await fetch(`http://${base}/session`, {
method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ name }),
});
return (await res.json()).token;
}
// Connect and collect messages; next(type) resolves with the next message of that type
async function connect(name) {
const ws = new WebSocket(`ws://${base}/ws?token=${await login(name)}`);
const inbox = [];
const waiters = [];
ws.addEventListener('message', (e) => {
const m = JSON.parse(e.data);
const i = waiters.findIndex((w) => w.type === m.type);
if (i >= 0) waiters.splice(i, 1)[0].resolve(m); else inbox.push(m);
});
await once(ws, 'open');
ws.json = (obj) => ws.send(JSON.stringify(obj));
ws.next = (type) => {
const i = inbox.findIndex((m) => m.type === type);
if (i >= 0) return Promise.resolve(inbox.splice(i, 1)[0]);
return new Promise((resolve) => waiters.push({ type, resolve }));
};
return ws;
}
test('rejects upgrades without a valid token', async () => {
const ws = new WebSocket(`ws://${base}/ws?token=forged.token`);
const [event] = await once(ws, 'error');
assert.equal(event.type, 'error');
});
test('messages reach everyone in the room, with history for late joiners', async () => {
const alice = await connect('alice');
const bob = await connect('bob');
alice.json({ type: 'join', room: 'general' });
await alice.next('joined');
bob.json({ type: 'join', room: 'general' });
const joined = await bob.next('joined');
assert.equal(joined.online, 2);
assert.equal((await alice.next('presence')).user, 'bob');
alice.json({ type: 'say', room: 'general', text: ' hello bob ' });
const [a, b] = await Promise.all([alice.next('message'), bob.next('message')]);
assert.equal(a.text, 'hello bob');
assert.equal(b.user, 'alice');
const carol = await connect('carol');
carol.json({ type: 'join', room: 'general' });
const late = await carol.next('joined');
assert.deepEqual(late.history.map((m) => m.text), ['hello bob']);
for (const ws of [alice, bob, carol]) ws.close();
});
test('invalid messages and non-members get errors, not crashes', async () => {
const eve = await connect('eve');
eve.send('not json');
assert.equal((await eve.next('error')).error, 'invalid message');
eve.json({ type: 'say', room: 'general', text: 'sneaky' });
assert.equal((await eve.next('error')).error, 'join general first');
eve.close();
});
test('per-connection rate limit', async () => {
const spammer = await connect('spammer');
spammer.json({ type: 'join', room: 'spam' });
await spammer.next('joined');
for (let i = 0; i < 15; i++) spammer.json({ type: 'say', room: 'spam', text: `msg ${i}` });
assert.equal((await spammer.next('error')).error, 'slow down');
spammer.close();
});
$ node --test 'test/*.test.js'
✔ rejects upgrades without a valid token
✔ messages reach everyone in the room, with history for late joiners
✔ invalid messages and non-members get errors, not crashes
✔ per-connection rate limit
ℹ tests 4
ℹ pass 4
ℹ fail 0
The helper connect() buffers incoming messages and exposes next(type), which turns
the event stream into awaitable steps — the same idea as events.once, filtered by
message type. Tests listen on port 0 so the OS picks a free port.
Scaling out: what changes with several processes¶
Everything in Rooms lives in one process's memory. To run several instances behind a
load balancer:
- Fan-out via pub/sub. When a process receives a
say, it publishes{ room, message }to a Redis channel (e.g.chat:general). Every process subscribes and calls its localrooms.broadcastfor messages it receives. A subscriber needs its own dedicated Redis connection. - Shared history. Store recent messages in a Redis list per room (
LPUSH+LTRIMto bound it) or in Postgres, rather than in process memory. - Presence becomes eventually consistent: keep a Redis set per room with entries that expire unless refreshed by heartbeats.
- Load balancer settings. WebSockets are long-lived; configure idle timeouts above
your heartbeat interval and make sure the balancer passes
Upgradeheaders.
This needs a Redis server, so it's left as an exercise rather than shown with output.
How It Actually Works¶
Follow one chat message:
- Alice's browser sends a masked text frame over her TCP connection.
- libuv's poll phase reports her socket readable;
wsunmasks and assembles the frame, checks it againstmaxPayload, and emits'message'with a Buffer. - Our handler spends a rate-limit token, parses JSON, validates with Zod, and checks room membership — all synchronous, microseconds of work.
rooms.broadcastserializes once and callssendfor each member. Eachsendframes the data and writes it to that member's socket. Writes that the kernel can't accept yet are queued in user space — visible asbufferedAmount.- The kernel transmits; each browser's
onmessagefires.
No step waits on I/O, so a single process handles a lot of chat traffic; the costs grow
with (messages per second × members per room). The heartbeat timer, the token refill
timer, and the listening socket are the handles that keep the process alive; close()
clears the timers, sends close frames with code 1001 ("going away"), and closes the server
so the process can exit.
Common mistakes in real-time projects¶
- Rendering messages with
innerHTML(stored XSS through chat). The client usestextContent. - Trusting client-sent usernames; here the name comes from the verified token.
- Unbounded history, per-room state that's never cleaned up, or sockets never removed
from rooms on
close. - No rate limiting — one client can flood a room.
- Forgetting that timers keep the process alive, so tests hang.
Exercise¶
- Add direct messages:
{ type: "dm", to: "bob", text }, delivered to all of Bob's connections (he may have several tabs). - Persist history in SQLite or Postgres (reusing Level 2 patterns) so it survives restarts; load the last 50 on join.
- Implement the Redis pub/sub fan-out described above, run two instances on different ports, and prove a message sent to one reaches a client on the other.
- Add a test for the heartbeat: create the server with
heartbeatMs: 50, connect a raw TCP client that completes the upgrade but never answers pings, and assert the server terminates it.