WebSockets
A full-duplex communication protocol enabling persistent, bidirectional connections between clients and servers over a single TCP connection.
WebSockets Cheatsheet
A full-duplex communication protocol enabling persistent, bidirectional connections between clients and servers over a single TCP connection.
Overview
WebSockets provide a standardised way to establish long-lived connections for real-time data exchange. Unlike HTTP's request-response model, WebSockets allow both parties to send data independently, making them ideal for applications requiring low-latency updates such as chat systems, live dashboards, and collaborative tools.
sequenceDiagram
participant Client
participant Server
Note over Client,Server: HTTP Upgrade Handshake
Client->>Server: GET /ws HTTP/1.1<br/>Upgrade: websocket<br/>Connection: Upgrade
Server-->>Client: HTTP/1.1 101 Switching Protocols<br/>Upgrade: websocket
Note over Client,Server: Full-Duplex Communication
Client->>Server: Message Frame
Server->>Client: Message Frame
Server->>Client: Message Frame
Client->>Server: Message Frame
Note over Client,Server: Connection Close
Client->>Server: Close Frame
Server-->>Client: Close Frame
Connection Establishment and Handshake
Key Concepts
- Upgrade Request: HTTP GET request with specific headers to initiate WebSocket connection
- Sec-WebSocket-Key: Base64-encoded random value sent by client for handshake validation
- Sec-WebSocket-Accept: Server's response hash proving it understands WebSocket protocol
- Protocol Negotiation: Optional subprotocol selection via
Sec-WebSocket-Protocolheader - Extensions: Optional features like compression via
Sec-WebSocket-Extensions
Common Patterns
# Client Request
GET /chat HTTP/1.1
Host: server.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Sec-WebSocket-Version: 13
Sec-WebSocket-Protocol: chat, superchat
Origin: http://example.com
# Server Response
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
Sec-WebSocket-Protocol: chat
// Browser client connection
const ws = new WebSocket('wss://server.example.com/chat');
// Connection with subprotocol
const ws = new WebSocket('wss://server.example.com/chat', ['chat', 'superchat']);
// Event handlers
ws.onopen = (event) => {
console.log('Connection established');
ws.send('Hello Server');
};
ws.onmessage = (event) => {
console.log('Received:', event.data);
};
ws.onerror = (event) => {
console.error('WebSocket error:', event);
};
ws.onclose = (event) => {
console.log(`Connection closed: ${event.code} ${event.reason}`);
};
Examples
Node.js Server with ws library:
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
wss.on('connection', (ws, req) => {
// Access client IP
const ip = req.socket.remoteAddress;
console.log(`Client connected: ${ip}`);
// Access URL parameters
const url = new URL(req.url, 'ws://localhost');
const userId = url.searchParams.get('userId');
ws.on('message', (message) => {
console.log(`Received: ${message}`);
// Echo back
ws.send(`Echo: ${message}`);
});
ws.on('close', (code, reason) => {
console.log(`Client disconnected: ${code} ${reason}`);
});
});
Message Framing and Formats
Key Concepts
- Frame Types: Text (0x1), Binary (0x2), Close (0x8), Ping (0x9), Pong (0xA)
- Masking: Client-to-server frames must be masked; server-to-client must not
- Fragmentation: Large messages split into multiple frames with FIN bit
- Payload Length: 7-bit for <126 bytes, 16-bit for <64KB, 64-bit for larger
graph LR
subgraph "WebSocket Frame Structure"
A[FIN<br/>1 bit] --> B[RSV<br/>3 bits]
B --> C[Opcode<br/>4 bits]
C --> D[Mask<br/>1 bit]
D --> E[Payload Length<br/>7+ bits]
E --> F[Masking Key<br/>0/32 bits]
F --> G[Payload Data]
end
style A fill:#e3f2fd
style C fill:#fff3e0
style G fill:#e8f5e9
Common Patterns
Sending different data types:
// Text message
ws.send('Hello World');
// JSON message
ws.send(JSON.stringify({
type: 'chat',
payload: { user: 'Alice', message: 'Hello!' }
}));
// Binary data (ArrayBuffer)
const buffer = new ArrayBuffer(8);
const view = new DataView(buffer);
view.setInt32(0, 12345);
ws.send(buffer);
// Blob data
const blob = new Blob(['Hello'], { type: 'text/plain' });
ws.send(blob);
Message protocol design:
// Structured message format
const MessageTypes = {
CHAT: 'chat',
NOTIFICATION: 'notification',
PRESENCE: 'presence',
ERROR: 'error'
};
function createMessage(type, payload) {
return JSON.stringify({
type,
payload,
timestamp: Date.now(),
id: crypto.randomUUID()
});
}
// Sending
ws.send(createMessage(MessageTypes.CHAT, {
room: 'general',
text: 'Hello everyone!'
}));
// Receiving and parsing
ws.onmessage = (event) => {
try {
const message = JSON.parse(event.data);
switch (message.type) {
case MessageTypes.CHAT:
handleChat(message.payload);
break;
case MessageTypes.NOTIFICATION:
handleNotification(message.payload);
break;
default:
console.warn('Unknown message type:', message.type);
}
} catch (error) {
console.error('Failed to parse message:', error);
}
};
Examples
Binary protocol with Protocol Buffers:
// Using protobufjs
const protobuf = require('protobufjs');
async function setupProtobuf() {
const root = await protobuf.load('messages.proto');
const ChatMessage = root.lookupType('ChatMessage');
// Encode message
const payload = { userId: 1, text: 'Hello', timestamp: Date.now() };
const errMsg = ChatMessage.verify(payload);
if (errMsg) throw Error(errMsg);
const message = ChatMessage.create(payload);
const buffer = ChatMessage.encode(message).finish();
ws.send(buffer);
// Decode message
ws.binaryType = 'arraybuffer';
ws.onmessage = (event) => {
const decoded = ChatMessage.decode(new Uint8Array(event.data));
console.log(decoded);
};
}
Heartbeat and Keep-Alive Mechanisms
Key Concepts
- Ping/Pong Frames: Built-in WebSocket control frames for connection health checks
- Application-Level Heartbeat: Custom messages when native ping/pong unavailable
- Idle Timeout: Disconnect clients that fail to respond within threshold
- Reconnection Strategy: Exponential backoff for automatic reconnection
sequenceDiagram
participant Client
participant Server
loop Every 30 seconds
Server->>Client: Ping Frame
Client-->>Server: Pong Frame
end
Note over Client,Server: Client stops responding
Server->>Client: Ping Frame
Note over Server: Timeout after 10s
Server->>Client: Close Frame (1001)
Common Patterns
Server-side heartbeat (Node.js):
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
// Track connection state
function heartbeat() {
this.isAlive = true;
}
wss.on('connection', (ws) => {
ws.isAlive = true;
ws.on('pong', heartbeat);
ws.on('message', (message) => {
// Handle message
});
});
// Check connections every 30 seconds
const interval = setInterval(() => {
wss.clients.forEach((ws) => {
if (ws.isAlive === false) {
console.log('Terminating dead connection');
return ws.terminate();
}
ws.isAlive = false;
ws.ping(); // Send ping frame
});
}, 30000);
wss.on('close', () => {
clearInterval(interval);
});
Client-side reconnection with exponential backoff:
class ReconnectingWebSocket {
constructor(url, protocols = []) {
this.url = url;
this.protocols = protocols;
this.reconnectAttempts = 0;
this.maxReconnectAttempts = 10;
this.baseDelay = 1000;
this.maxDelay = 30000;
this.connect();
}
connect() {
this.ws = new WebSocket(this.url, this.protocols);
this.ws.onopen = () => {
console.log('Connected');
this.reconnectAttempts = 0;
if (this.onopen) this.onopen();
};
this.ws.onclose = (event) => {
if (event.code !== 1000) { // Abnormal closure
this.scheduleReconnect();
}
if (this.onclose) this.onclose(event);
};
this.ws.onerror = (error) => {
if (this.onerror) this.onerror(error);
};
this.ws.onmessage = (event) => {
if (this.onmessage) this.onmessage(event);
};
}
scheduleReconnect() {
if (this.reconnectAttempts >= this.maxReconnectAttempts) {
console.error('Max reconnection attempts reached');
return;
}
// Exponential backoff with jitter
const delay = Math.min(
this.baseDelay * Math.pow(2, this.reconnectAttempts) +
Math.random() * 1000,
this.maxDelay
);
console.log(`Reconnecting in ${delay}ms...`);
this.reconnectAttempts++;
setTimeout(() => this.connect(), delay);
}
send(data) {
if (this.ws.readyState === WebSocket.OPEN) {
this.ws.send(data);
}
}
close() {
this.ws.close(1000, 'Client closing');
}
}
Examples
Application-level heartbeat:
// Client-side heartbeat (when ping/pong not accessible)
class HeartbeatWebSocket {
constructor(url) {
this.url = url;
this.heartbeatInterval = 25000;
this.heartbeatTimeout = 35000;
this.connect();
}
connect() {
this.ws = new WebSocket(this.url);
this.ws.onopen = () => {
this.startHeartbeat();
};
this.ws.onmessage = (event) => {
const message = JSON.parse(event.data);
if (message.type === 'pong') {
this.resetTimeout();
} else {
// Handle other messages
}
};
this.ws.onclose = () => {
this.stopHeartbeat();
};
}
startHeartbeat() {
this.pingInterval = setInterval(() => {
this.ws.send(JSON.stringify({ type: 'ping' }));
// Set timeout for pong response
this.pongTimeout = setTimeout(() => {
console.log('Pong timeout, closing connection');
this.ws.close();
}, this.heartbeatTimeout - this.heartbeatInterval);
}, this.heartbeatInterval);
}
resetTimeout() {
clearTimeout(this.pongTimeout);
}
stopHeartbeat() {
clearInterval(this.pingInterval);
clearTimeout(this.pongTimeout);
}
}
Scaling WebSocket Servers
Key Concepts
- Sticky Sessions: Route client to same server for connection duration
- Pub/Sub Backend: Redis, RabbitMQ, or Kafka for cross-server messaging
- Horizontal Scaling: Multiple WebSocket server instances behind load balancer
- Connection State: Externalise state to shared storage for failover
- Cluster Awareness: Broadcast messages across all server instances
graph TB
subgraph "Scaled WebSocket Architecture"
LB[Load Balancer<br/>Sticky Sessions]
subgraph "WebSocket Servers"
WS1[WS Server 1]
WS2[WS Server 2]
WS3[WS Server 3]
end
Redis[(Redis<br/>Pub/Sub)]
LB --> WS1
LB --> WS2
LB --> WS3
WS1 <--> Redis
WS2 <--> Redis
WS3 <--> Redis
end
Client1[Client A] --> LB
Client2[Client B] --> LB
Client3[Client C] --> LB
style Redis fill:#ffebee
style LB fill:#e8f5e9
Common Patterns
Redis Pub/Sub for multi-server broadcasting:
const WebSocket = require('ws');
const Redis = require('ioredis');
const wss = new WebSocket.Server({ port: 8080 });
const redisPub = new Redis();
const redisSub = new Redis();
// Store local connections by room
const rooms = new Map();
// Subscribe to Redis channel
redisSub.subscribe('ws-broadcast', (err) => {
if (err) console.error('Redis subscribe error:', err);
});
// Handle Redis messages
redisSub.on('message', (channel, message) => {
const { room, data, excludeServerId } = JSON.parse(message);
// Skip if message from this server
if (excludeServerId === process.env.SERVER_ID) return;
// Broadcast to local clients in room
const clients = rooms.get(room) || new Set();
clients.forEach((ws) => {
if (ws.readyState === WebSocket.OPEN) {
ws.send(data);
}
});
});
// Broadcast to all servers
function broadcastToRoom(room, data) {
redisPub.publish('ws-broadcast', JSON.stringify({
room,
data,
excludeServerId: process.env.SERVER_ID
}));
// Also send to local clients
const clients = rooms.get(room) || new Set();
clients.forEach((ws) => {
if (ws.readyState === WebSocket.OPEN) {
ws.send(data);
}
});
}
wss.on('connection', (ws, req) => {
const room = new URL(req.url, 'ws://localhost').searchParams.get('room');
// Join room
if (!rooms.has(room)) rooms.set(room, new Set());
rooms.get(room).add(ws);
ws.on('message', (message) => {
broadcastToRoom(room, message.toString());
});
ws.on('close', () => {
rooms.get(room)?.delete(ws);
});
});
Load balancer configuration (nginx):
# nginx.conf
upstream websocket {
# Sticky sessions using IP hash
ip_hash;
server ws1.example.com:8080;
server ws2.example.com:8080;
server ws3.example.com:8080;
}
server {
listen 443 ssl;
server_name ws.example.com;
ssl_certificate /etc/ssl/certs/server.crt;
ssl_certificate_key /etc/ssl/private/server.key;
location /ws {
proxy_pass http://websocket;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
# Timeout settings
proxy_read_timeout 86400s;
proxy_send_timeout 86400s;
}
}
Examples
Socket.IO with Redis adapter:
const { Server } = require('socket.io');
const { createAdapter } = require('@socket.io/redis-adapter');
const { createClient } = require('redis');
async function createServer() {
const io = new Server(3000);
const pubClient = createClient({ url: 'redis://localhost:6379' });
const subClient = pubClient.duplicate();
await Promise.all([pubClient.connect(), subClient.connect()]);
// Use Redis adapter for horizontal scaling
io.adapter(createAdapter(pubClient, subClient));
io.on('connection', (socket) => {
socket.join('room1');
socket.on('message', (data) => {
// Broadcasts to all servers via Redis
io.to('room1').emit('message', data);
});
});
return io;
}
Security Considerations
Key Concepts
- WSS (WebSocket Secure): WebSocket over TLS, always use in production
- Origin Validation: Check
Originheader to prevent CSRF attacks - Authentication: Token-based auth via query params or initial message
- Rate Limiting: Protect against DoS by limiting connections and messages
- Input Validation: Sanitise all incoming data to prevent injection attacks
flowchart TD
A[Client Connection] --> B{WSS?}
B -->|No| C[Reject: Use HTTPS]
B -->|Yes| D{Valid Origin?}
D -->|No| E[Reject: Invalid Origin]
D -->|Yes| F{Authenticated?}
F -->|No| G[Reject: Unauthorised]
F -->|Yes| H{Rate Limited?}
H -->|Yes| I[Reject: Too Many Requests]
H -->|No| J[Accept Connection]
style C fill:#ffcdd2
style E fill:#ffcdd2
style G fill:#ffcdd2
style I fill:#ffcdd2
style J fill:#c8e6c9
Common Patterns
Origin validation and authentication:
const WebSocket = require('ws');
const jwt = require('jsonwebtoken');
const wss = new WebSocket.Server({
port: 8080,
verifyClient: ({ origin, req }, callback) => {
// Validate origin
const allowedOrigins = [
'https://example.com',
'https://app.example.com'
];
if (!allowedOrigins.includes(origin)) {
callback(false, 403, 'Forbidden origin');
return;
}
// Validate token from query string
const url = new URL(req.url, 'ws://localhost');
const token = url.searchParams.get('token');
if (!token) {
callback(false, 401, 'No token provided');
return;
}
try {
const decoded = jwt.verify(token, process.env.JWT_SECRET);
req.user = decoded;
callback(true);
} catch (error) {
callback(false, 401, 'Invalid token');
}
}
});
wss.on('connection', (ws, req) => {
console.log(`User ${req.user.id} connected`);
ws.on('message', (message) => {
// User info available from req.user
});
});
Rate limiting:
const WebSocket = require('ws');
class RateLimiter {
constructor(maxRequests, windowMs) {
this.maxRequests = maxRequests;
this.windowMs = windowMs;
this.clients = new Map();
}
isRateLimited(clientId) {
const now = Date.now();
const clientData = this.clients.get(clientId) || { count: 0, resetTime: now + this.windowMs };
if (now > clientData.resetTime) {
clientData.count = 0;
clientData.resetTime = now + this.windowMs;
}
clientData.count++;
this.clients.set(clientId, clientData);
return clientData.count > this.maxRequests;
}
}
const wss = new WebSocket.Server({ port: 8080 });
const connectionLimiter = new RateLimiter(5, 60000); // 5 connections per minute
const messageLimiter = new RateLimiter(100, 60000); // 100 messages per minute
wss.on('connection', (ws, req) => {
const clientIp = req.socket.remoteAddress;
// Check connection rate limit
if (connectionLimiter.isRateLimited(clientIp)) {
ws.close(1008, 'Too many connections');
return;
}
ws.on('message', (message) => {
// Check message rate limit
if (messageLimiter.isRateLimited(clientIp)) {
ws.send(JSON.stringify({ error: 'Rate limited' }));
return;
}
// Process message
});
});
Examples
Secure WebSocket server with Helmet-style headers:
const https = require('https');
const fs = require('fs');
const WebSocket = require('ws');
// Create HTTPS server
const server = https.createServer({
cert: fs.readFileSync('/path/to/cert.pem'),
key: fs.readFileSync('/path/to/key.pem')
});
const wss = new WebSocket.Server({ server });
// Input validation and sanitisation
function validateMessage(data) {
try {
const message = JSON.parse(data);
// Type checking
if (typeof message.type !== 'string') return null;
if (message.type.length > 50) return null;
// Sanitise text content
if (message.text) {
message.text = message.text
.substring(0, 1000) // Max length
.replace(/[<>]/g, ''); // Basic XSS prevention
}
return message;
} catch {
return null;
}
}
wss.on('connection', (ws) => {
ws.on('message', (data) => {
const message = validateMessage(data.toString());
if (!message) {
ws.send(JSON.stringify({ error: 'Invalid message format' }));
return;
}
// Process validated message
});
});
server.listen(443);
Common Libraries and Frameworks
Key Concepts
- ws: Lightweight, fast WebSocket implementation for Node.js
- Socket.IO: Feature-rich library with fallbacks, rooms, and namespaces
- uWebSockets.js: High-performance C++ WebSocket server with Node.js bindings
- sockjs: WebSocket emulation with fallback transports
- Autobahn: WAMP (Web Application Messaging Protocol) implementation
graph LR
subgraph "Library Comparison"
ws[ws<br/>Lightweight<br/>Standards-compliant]
sio[Socket.IO<br/>Feature-rich<br/>Auto-fallback]
uws[uWebSockets.js<br/>High performance<br/>Low memory]
end
style ws fill:#e3f2fd
style sio fill:#fff3e0
style uws fill:#e8f5e9
Common Patterns
ws - Basic server:
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
wss.on('connection', (ws) => {
ws.on('message', (message) => {
// Broadcast to all clients
wss.clients.forEach((client) => {
if (client.readyState === WebSocket.OPEN) {
client.send(message.toString());
}
});
});
});
Socket.IO - Server with rooms and namespaces:
const { Server } = require('socket.io');
const io = new Server(3000, {
cors: {
origin: 'https://example.com',
methods: ['GET', 'POST']
}
});
// Namespace for chat
const chat = io.of('/chat');
chat.on('connection', (socket) => {
// Authentication middleware
const token = socket.handshake.auth.token;
// Join room
socket.on('join', (room) => {
socket.join(room);
socket.to(room).emit('user-joined', socket.id);
});
// Send to specific room
socket.on('message', ({ room, text }) => {
chat.to(room).emit('message', {
from: socket.id,
text,
timestamp: Date.now()
});
});
// Private message
socket.on('private', ({ to, text }) => {
socket.to(to).emit('private', {
from: socket.id,
text
});
});
// Acknowledgement
socket.on('action', (data, callback) => {
// Process action
callback({ status: 'ok' });
});
});
Socket.IO - Client:
import { io } from 'socket.io-client';
const socket = io('https://example.com/chat', {
auth: {
token: 'your-jwt-token'
},
reconnection: true,
reconnectionAttempts: 5,
reconnectionDelay: 1000
});
socket.on('connect', () => {
socket.emit('join', 'general');
});
socket.on('message', (data) => {
console.log(`${data.from}: ${data.text}`);
});
// With acknowledgement
socket.emit('action', { type: 'subscribe' }, (response) => {
console.log('Server acknowledged:', response);
});
// Error handling
socket.on('connect_error', (error) => {
console.error('Connection error:', error);
});
Examples
uWebSockets.js - High-performance server:
const uWS = require('uWebSockets.js');
const app = uWS.App();
app.ws('/*', {
// Configuration
compression: uWS.SHARED_COMPRESSOR,
maxPayloadLength: 16 * 1024 * 1024,
idleTimeout: 120,
// Handlers
open: (ws) => {
ws.subscribe('broadcast');
console.log('Client connected');
},
message: (ws, message, isBinary) => {
const text = Buffer.from(message).toString();
// Broadcast to all subscribers
app.publish('broadcast', text, isBinary);
},
close: (ws, code, message) => {
console.log('Client disconnected');
}
});
app.listen(8080, (token) => {
if (token) {
console.log('Listening on port 8080');
}
});
Python WebSocket server with websockets library:
import asyncio
import websockets
import json
connected_clients = set()
# websockets 14+ (modern asyncio API): handler takes only the connection.
# The request path is on websocket.request.path if you need it.
async def handler(websocket):
# Register client
connected_clients.add(websocket)
try:
async for message in websocket:
data = json.loads(message)
# Broadcast to all clients
broadcast_message = json.dumps({
'type': 'message',
'payload': data
})
# websockets.broadcast() is the idiomatic fan-out helper; it skips
# closing/closed connections for you (replaces the removed .open check).
websockets.broadcast(connected_clients, broadcast_message)
finally:
connected_clients.remove(websocket)
async def main():
async with websockets.serve(handler, 'localhost', 8080):
await asyncio.Future() # Run forever
if __name__ == '__main__':
asyncio.run(main())
Use Cases
Key Concepts
- Real-time Chat: Instant message delivery with presence indicators
- Live Notifications: Push updates without polling
- Collaborative Editing: Synchronise document changes across users
- Live Dashboards: Stream metrics and data visualisations
- Gaming: Low-latency game state synchronisation
- Financial Trading: Real-time price updates and order execution
graph TB
subgraph "WebSocket Use Cases"
Chat[Real-time Chat<br/>Instant messaging<br/>Typing indicators]
Notify[Notifications<br/>Push alerts<br/>Activity feeds]
Collab[Collaboration<br/>Document editing<br/>Whiteboarding]
Data[Live Data<br/>Stock tickers<br/>Sports scores]
end
style Chat fill:#e3f2fd
style Notify fill:#fff3e0
style Collab fill:#e8f5e9
style Data fill:#f3e5f5
Examples
Real-time chat with typing indicators:
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
const rooms = new Map();
wss.on('connection', (ws) => {
let currentRoom = null;
let username = null;
ws.on('message', (data) => {
const message = JSON.parse(data);
switch (message.type) {
case 'join':
username = message.username;
currentRoom = message.room;
if (!rooms.has(currentRoom)) {
rooms.set(currentRoom, new Set());
}
rooms.get(currentRoom).add(ws);
broadcast(currentRoom, {
type: 'user-joined',
username
});
break;
case 'message':
broadcast(currentRoom, {
type: 'message',
username,
text: message.text,
timestamp: Date.now()
});
break;
case 'typing':
broadcast(currentRoom, {
type: 'typing',
username
}, ws); // Exclude sender
break;
case 'stop-typing':
broadcast(currentRoom, {
type: 'stop-typing',
username
}, ws);
break;
}
});
ws.on('close', () => {
if (currentRoom && rooms.has(currentRoom)) {
rooms.get(currentRoom).delete(ws);
broadcast(currentRoom, {
type: 'user-left',
username
});
}
});
});
function broadcast(room, message, exclude = null) {
const clients = rooms.get(room);
if (!clients) return;
const data = JSON.stringify(message);
clients.forEach((client) => {
if (client !== exclude && client.readyState === WebSocket.OPEN) {
client.send(data);
}
});
}
Live notifications system:
// Server
const WebSocket = require('ws');
const Redis = require('ioredis');
const wss = new WebSocket.Server({ port: 8080 });
const redis = new Redis();
const userConnections = new Map();
wss.on('connection', (ws, req) => {
const userId = req.user.id; // From auth middleware
// Store connection
if (!userConnections.has(userId)) {
userConnections.set(userId, new Set());
}
userConnections.get(userId).add(ws);
// Send unread notifications
redis.lrange(`notifications:${userId}`, 0, -1)
.then((notifications) => {
notifications.forEach((n) => ws.send(n));
});
ws.on('close', () => {
userConnections.get(userId)?.delete(ws);
});
});
// Function to send notification (called from your app)
function sendNotification(userId, notification) {
const message = JSON.stringify({
type: 'notification',
...notification,
timestamp: Date.now()
});
// Store in Redis
redis.lpush(`notifications:${userId}`, message);
redis.ltrim(`notifications:${userId}`, 0, 99); // Keep last 100
// Send to connected clients
const connections = userConnections.get(userId);
if (connections) {
connections.forEach((ws) => {
if (ws.readyState === WebSocket.OPEN) {
ws.send(message);
}
});
}
}
Live dashboard with metrics streaming:
// Server streaming system metrics
const WebSocket = require('ws');
const os = require('os');
const wss = new WebSocket.Server({ port: 8080 });
// Collect metrics
function getMetrics() {
const cpus = os.cpus();
const totalMemory = os.totalmem();
const freeMemory = os.freemem();
return {
timestamp: Date.now(),
cpu: {
count: cpus.length,
usage: cpus.map((cpu) => {
const total = Object.values(cpu.times).reduce((a, b) => a + b);
return ((total - cpu.times.idle) / total * 100).toFixed(2);
})
},
memory: {
total: totalMemory,
free: freeMemory,
used: totalMemory - freeMemory,
percentage: ((1 - freeMemory / totalMemory) * 100).toFixed(2)
},
uptime: os.uptime()
};
}
// Broadcast metrics every second
setInterval(() => {
const metrics = JSON.stringify({
type: 'metrics',
data: getMetrics()
});
wss.clients.forEach((client) => {
if (client.readyState === WebSocket.OPEN) {
client.send(metrics);
}
});
}, 1000);
wss.on('connection', (ws) => {
// Send initial metrics immediately
ws.send(JSON.stringify({
type: 'metrics',
data: getMetrics()
}));
});
Quick Reference
| Task | Code/Pattern |
|---|---|
| Create connection | new WebSocket('wss://host/path') |
| Send text | ws.send('message') |
| Send JSON | ws.send(JSON.stringify(obj)) |
| Receive message | ws.onmessage = (e) => console.log(e.data) |
| Close connection | ws.close(1000, 'reason') |
| Check state | ws.readyState === WebSocket.OPEN |
| Server broadcast | wss.clients.forEach(c => c.send(msg)) |
| Ping client (server) | ws.ping() |
| Handle pong | ws.on('pong', heartbeat) |
| Binary mode | ws.binaryType = 'arraybuffer' |
WebSocket Ready States
| State | Value | Description |
|---|---|---|
| CONNECTING | 0 | Connection in progress |
| OPEN | 1 | Connection established |
| CLOSING | 2 | Closing handshake in progress |
| CLOSED | 3 | Connection closed |
Close Codes
| Code | Name | Description |
|---|---|---|
| 1000 | Normal Closure | Clean close |
| 1001 | Going Away | Server/client shutting down |
| 1002 | Protocol Error | Protocol violation |
| 1003 | Unsupported Data | Invalid data type received |
| 1006 | Abnormal Closure | Connection lost unexpectedly |
| 1008 | Policy Violation | Message violates policy |
| 1011 | Internal Error | Server encountered error |
Common Issues and Solutions
| Issue | Cause | Solution |
|---|---|---|
| Connection closes immediately | Missing Upgrade headers | Ensure proxy passes Upgrade and Connection headers |
| 403 Forbidden | Origin validation failing | Add client origin to server's allowed origins list |
| Connection timeout | Proxy/firewall timeout | Implement heartbeat mechanism; increase proxy timeout |
| Messages not received | JSON parse errors | Validate message format; use try-catch around JSON.parse |
| Memory leak | Connections not cleaned up | Remove event listeners and clear intervals on close |
| Cannot scale horizontally | No message broker | Implement Redis Pub/Sub or similar for cross-server communication |
| Client reconnects constantly | Server rejecting connections | Check authentication; verify rate limiting; check server logs |
| Mixed content error | HTTP page using WSS | Use HTTPS for the page or WSS for WebSocket |
| Binary data corrupted | Wrong binary type | Set ws.binaryType = 'arraybuffer' before receiving |
| High latency | Geographic distance | Deploy WebSocket servers in multiple regions; use CDN with WebSocket support |
Debugging Tips
// Enable debug logging for ws library
const WebSocket = require('ws');
process.env.DEBUG = 'ws';
// Log all WebSocket events
const ws = new WebSocket('wss://example.com');
['open', 'close', 'error', 'message', 'ping', 'pong'].forEach((event) => {
ws.on(event, (...args) => {
console.log(`[WS ${event}]`, ...args);
});
});
// Monitor connection state
setInterval(() => {
console.log('ReadyState:', ws.readyState);
console.log('Buffered:', ws.bufferedAmount);
}, 5000);
Performance Optimisation
// Batch messages to reduce overhead
class MessageBatcher {
constructor(ws, flushInterval = 100) {
this.ws = ws;
this.queue = [];
this.interval = setInterval(() => this.flush(), flushInterval);
}
send(message) {
this.queue.push(message);
}
flush() {
if (this.queue.length === 0) return;
this.ws.send(JSON.stringify({
type: 'batch',
messages: this.queue
}));
this.queue = [];
}
destroy() {
clearInterval(this.interval);
this.flush();
}
}