WebSocket & Real-Time Expert
Build reliable real-time communication systems using WebSockets, Server-Sent Events, and managed real-time services.
Activation Triggers
Activate on: "WebSocket", "real-time", "SSE", "Socket.io", "live updates", "push notifications", "bidirectional", "presence", "live cursors", "collaborative editing"
NOT for: Message queue setup → event-driven-architecture-expert | Gateway WebSocket routing → api-gateway-reverse-proxy-expert | Streaming data pipelines → streaming-pipeline-architect
Quick Start
- Choose protocol — WebSocket for bidirectional, SSE for server-push, WebTransport for low-latency
- Plan reconnection — exponential backoff with jitter, resume from last event ID
- Design message format — typed JSON with
type discriminator and monotonic sequence IDs
- Scale horizontally — use Redis Pub/Sub or NATS as a broadcast backplane
- Handle presence — heartbeat-based with configurable timeout (30s default)
Core Capabilities
| Domain |
Technologies |
| WebSocket |
ws (Node), Socket.io 4.8+, uWebSockets.js |
| SSE |
Native EventSource, @microsoft/fetch-event-source |
| Managed |
Supabase Realtime, Ably, Pusher, PartyKit |
| Scaling |
Redis Pub/Sub, NATS, @socket.io/redis-adapter |
| Protocols |
WebSocket (RFC 6455), SSE, WebTransport (HTTP/3) |
Architecture Patterns
Scaled WebSocket with Redis Backplane
Client A ──ws──→ Server 1 ←──redis pub/sub──→ Server 2 ←──ws── Client B
│ │
└─────── Redis Cluster ───────┘
Each server subscribes to channels. When Server 1 receives a message
for a room, it publishes to Redis. Server 2 picks it up and forwards
to its connected clients.
SSE with Last-Event-ID Resume
// Server: SSE endpoint with resume support
app.get('/events', (req, res) => {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
});
const lastId = parseInt(req.headers['last-event-id'] || '0');
// Replay missed events from store
const missed = eventStore.since(lastId);
missed.forEach(evt => {
res.write(`id: ${evt.id}\nevent: ${evt.type}\ndata: ${JSON.stringify(evt.data)}\n\n`);
});
// Subscribe to new events
const unsub = eventBus.subscribe(evt => {
res.write(`id: ${evt.id}\nevent: ${evt.type}\ndata: ${JSON.stringify(evt.data)}\n\n`);
});
req.on('close', unsub);
});
Presence with Heartbeat
// Client sends heartbeat every 15s
const HEARTBEAT_INTERVAL = 15_000;
const PRESENCE_TIMEOUT = 45_000; // 3 missed heartbeats = offline
// Server tracks presence
const presence = new Map<string, { userId: string; lastSeen: number }>();
ws.on('message', (msg) => {
const { type, userId } = JSON.parse(msg);
if (type === 'heartbeat') {
presence.set(userId, { userId, lastSeen: Date.now() });
}
});
// Sweep stale presence every 10s
setInterval(() => {
const cutoff = Date.now() - PRESENCE_TIMEOUT;
for (const [id, p] of presence) {
if (p.lastSeen < cutoff) {
presence.delete(id);
broadcast({ type: 'presence:leave', userId: id });
}
}
}, 10_000);
Anti-Patterns
- No reconnection logic — connections will drop; always implement reconnect with exponential backoff and jitter
- Polling disguised as real-time — if you poll every 1s over WebSocket, just use SSE or actual polling
- Single server assumption — WebSocket state is per-server; you need a pub/sub backplane for horizontal scaling
- Unbounded message buffers — set max buffer size and drop/queue when clients are slow (backpressure)
- Missing heartbeat — idle connections get killed by proxies and load balancers (typical timeout: 60s)
Quality Checklist
1---2name: websocket-realtime-expert3description: WebSockets, SSE, and real-time communication with Socket.io and native APIs. Activate on: WebSocket, real-time, SSE, Socket.io, live updates, push notifications, bidirectional, presence. NOT for: message queue infrastructure (use event-driven-architecture-expert), API gateway routing (use api-gateway-reverse-proxy-expert).4license: Apache-2.05---6
7# WebSocket & Real-Time Expert
8
9Build reliable real-time communication systems using WebSockets, Server-Sent Events, and managed real-time services.
10
11## Activation Triggers
12
13**Activate on:** "WebSocket", "real-time", "SSE", "Socket.io", "live updates", "push notifications", "bidirectional", "presence", "live cursors", "collaborative editing"
14
15**NOT for:** Message queue setup → `event-driven-architecture-expert` | Gateway WebSocket routing → `api-gateway-reverse-proxy-expert` | Streaming data pipelines → `streaming-pipeline-architect`
16
17## Quick Start
18
191. **Choose protocol** — WebSocket for bidirectional, SSE for server-push, WebTransport for low-latency
202. **Plan reconnection** — exponential backoff with jitter, resume from last event ID
213. **Design message format** — typed JSON with `type` discriminator and monotonic sequence IDs
224. **Scale horizontally** — use Redis Pub/Sub or NATS as a broadcast backplane
235. **Handle presence** — heartbeat-based with configurable timeout (30s default)
24
25## Core Capabilities
26
27| Domain | Technologies |
28|--------|-------------|
29| **WebSocket** | ws (Node), Socket.io 4.8+, uWebSockets.js |
30| **SSE** | Native EventSource, @microsoft/fetch-event-source |
31| **Managed** | Supabase Realtime, Ably, Pusher, PartyKit |
32| **Scaling** | Redis Pub/Sub, NATS, @socket.io/redis-adapter |
33| **Protocols** | WebSocket (RFC 6455), SSE, WebTransport (HTTP/3) |
34
35## Architecture Patterns
36
37### Scaled WebSocket with Redis Backplane
38
39```
40Client A ──ws──→ Server 1 ←──redis pub/sub──→ Server 2 ←──ws── Client B
41 │ │
42 └─────── Redis Cluster ───────┘
43
44Each server subscribes to channels. When Server 1 receives a message
45for a room, it publishes to Redis. Server 2 picks it up and forwards
46to its connected clients.
47```
48
49### SSE with Last-Event-ID Resume
50
51```typescript
52// Server: SSE endpoint with resume support
53app.get('/events', (req, res) => {
54 res.writeHead(200, {
55 'Content-Type': 'text/event-stream',
56 'Cache-Control': 'no-cache',
57 'Connection': 'keep-alive',
58 });
59
60 const lastId = parseInt(req.headers['last-event-id'] || '0');
61 // Replay missed events from store
62 const missed = eventStore.since(lastId);
63 missed.forEach(evt => {
64 res.write(`id: ${evt.id}\nevent: ${evt.type}\ndata: ${JSON.stringify(evt.data)}\n\n`);
65 });
66
67 // Subscribe to new events
68 const unsub = eventBus.subscribe(evt => {
69 res.write(`id: ${evt.id}\nevent: ${evt.type}\ndata: ${JSON.stringify(evt.data)}\n\n`);
70 });
71 req.on('close', unsub);
72});
73```
74
75### Presence with Heartbeat
76
77```typescript
78// Client sends heartbeat every 15s
79const HEARTBEAT_INTERVAL = 15_000;
80const PRESENCE_TIMEOUT = 45_000; // 3 missed heartbeats = offline
81
82// Server tracks presence
83const presence = new Map<string, { userId: string; lastSeen: number }>();
84
85ws.on('message', (msg) => {
86 const { type, userId } = JSON.parse(msg);
87 if (type === 'heartbeat') {
88 presence.set(userId, { userId, lastSeen: Date.now() });
89 }
90});
91
92// Sweep stale presence every 10s
93setInterval(() => {
94 const cutoff = Date.now() - PRESENCE_TIMEOUT;
95 for (const [id, p] of presence) {
96 if (p.lastSeen < cutoff) {
97 presence.delete(id);
98 broadcast({ type: 'presence:leave', userId: id });
99 }
100 }
101}, 10_000);
102```
103
104## Anti-Patterns
105
1061. **No reconnection logic** — connections will drop; always implement reconnect with exponential backoff and jitter
1072. **Polling disguised as real-time** — if you poll every 1s over WebSocket, just use SSE or actual polling
1083. **Single server assumption** — WebSocket state is per-server; you need a pub/sub backplane for horizontal scaling
1094. **Unbounded message buffers** — set max buffer size and drop/queue when clients are slow (backpressure)
1105. **Missing heartbeat** — idle connections get killed by proxies and load balancers (typical timeout: 60s)
111
112## Quality Checklist
113
114- [ ] Reconnection with exponential backoff and jitter implemented
115- [ ] Last-Event-ID or sequence-based resume for missed messages
116- [ ] Heartbeat/ping-pong keeps connections alive (every 15-30s)
117- [ ] Redis/NATS backplane for multi-server deployments
118- [ ] Message schema versioned with `type` discriminator
119- [ ] Backpressure handling for slow consumers
120- [ ] Connection count monitoring with alerts
121- [ ] Graceful degradation: SSE fallback if WebSocket blocked
122- [ ] Authentication on connection (verify token before upgrade)
123- [ ] Load tested: target connections per server validated