Explorer
Node.js

Chapter 19: WebSockets & Real-Time Communication: Socket.IO, Namespaces, Redis Adapter, and SSE

Chapter 19: WebSockets & Real-Time Communication: Socket.IO, Namespaces, Redis Adapter, and SSE

The HTTP protocol is fundamentally unidirectional and client-driven: the client sends a request, the server computes a response, and the connection terminates.

For traditional CRUD applications, this request-response cycle is ideal. However, modern applications demand real-time, instantaneous interactivity:

  • Collaborative code editors where multiple learners type in the same buffer simultaneously.
  • Live test case execution streams and judge worker status updates.
  • Real-time chat, notification badges, and instant messaging.
  • Live competitive programming leaderboards.

When servers need to push data to clients spontaneously, traditional HTTP falls short.

In this chapter, you will master the architecture of real-time communication in Node.js:

  1. The evolution from Short Polling and Long Polling to Server-Sent Events (SSE) and WebSockets.
  2. The low-level WebSocket Protocol Handshake (RFC 6455) and HTTP Upgrade mechanism.
  3. Building real-time systems with Socket.IO: Namespaces, Rooms, and Acknowledgments.
  4. Authenticating Real-Time Sockets with JWTs during the handshake.
  5. Horizontal Scaling with the Redis Adapter: Distributing real-time events across multi-container clusters.
  6. Server-Sent Events (SSE): When unidirectional streaming over standard HTTP is superior to WebSockets.

1. The Evolution of Real-Time Web

Real-World Analogy: How People Communicate

  • Short Polling is like calling your friend on the phone every 5 seconds asking: "Are you ready yet? No. Are you ready yet? No. Are you ready yet? Yes!" It consumes immense energy and generates hundreds of wasted phone dials.
  • Server-Sent Events (SSE) is like listening to an FM Radio broadcast. The studio DJ speaks into the microphone and music streams continuously into your car speakers. You cannot talk back to the DJ through your car radio, but for live traffic updates, stock tickers, and AI token streaming, it is simple, lightweight, and crystal clear.
  • WebSockets is like a live two-way phone call. Once connected, both sides hold open a continuous, dedicated audio channel and can speak back and forth at the exact same millisecond with zero setup lag.
TEXT
1. Short Polling:
   Client  ──► GET /status (Every 2s) ──► Server (No new data)
   Client  ──► GET /status (Every 2s) ──► Server (No new data)
   Client  ──► GET /status (Every 2s) ──► Server (200 OK: Data!)
   ⚠️ Massive overhead: Hundreds of redundant HTTP headers per minute.

2. Long Polling:
   Client  ──► GET /status ────────────► Server holds socket open...
   Client  ◄── 200 OK (After 15s) ◄───── Server emits event when ready
   Client  ──► Immediately re-opens ───► Server holds socket open...
   ⚠️ Reconnection latency and connection churn.

3. Server-Sent Events (SSE):
   Client  ──► GET /events (Accept: text/event-stream)
   Server  ──► Persistent HTTP/2 Stream (Server pushes events at will)
   ✅ Native browser auto-reconnect, lightweight, text-based.
   ❌ Unidirectional only (Server-to-Client).

4. WebSockets:
   Client  ──► GET /socket (Upgrade: websocket)
   Server  ◄── 101 Switching Protocols
   ══════════════════════════════════════════════════════════════════════════
   Bi-directional, Full-Duplex TCP Socket (Raw frames: 2-10 bytes overhead)
   ══════════════════════════════════════════════════════════════════════════
   ✅ True bi-directional, lowest latency, binary and text support.

Technology Comparison Matrix

Feature Short Polling Long Polling Server-Sent Events (SSE) WebSockets
Direction Client → Server Client → Server Server → Client (Unidirectional) Full-Duplex (Bi-directional)
Protocol HTTP HTTP HTTP / HTTP/2 WebSocket (ws:// / wss://)
Connection Disconnects each time Reconnects after message Persistent single HTTP connection Persistent single TCP socket
Data Types Text / JSON Text / JSON UTF-8 Text (Event Streams) Binary & UTF-8 Text
Header Overhead ~800 bytes per poll ~800 bytes per cycle Single initial HTTP header ~2 to 10 bytes per frame
Best Used For Low-frequency checks Legacy fallback Stock tickers, AI token streaming, notifications Live chat, gaming, collaborative editing

2. The WebSocket Protocol Handshake (RFC 6455)

A WebSocket connection does not start as a raw TCP socket—it begins as an HTTP request called a WebSocket Handshake. This allows WebSocket traffic to traverse standard web ports (80 for HTTP, 443 for HTTPS) and pass through corporate firewalls and reverse proxies.

The Handshake Sequence

TEXT
Client (Browser)                                    Node.js Server
       │                                                   │
       │  1. HTTP GET /socket                              │
       │     Upgrade: websocket                            │
       │     Connection: Upgrade                           │
       │     Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==   │
       │     Sec-WebSocket-Version: 13                     │
       ───────────────────────────────────────────────────►│
       │                                                   │ (Validates key & auth)
       │  2. HTTP/1.1 101 Switching Protocols              │
       │     Upgrade: websocket                            │
       │     Connection: Upgrade                           │
       │     Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
       │◄──────────────────────────────────────────────────│
       │                                                   │
       ═════════════════════════════════════════════════════
         PERSISTENT FULL-DUPLEX TCP FRAMES (Bidirectional)
       ═════════════════════════════════════════════════════

Handshake Header Breakdown:

  1. Upgrade: websocket: Tells the server to switch protocol from HTTP to WebSocket.
  2. Connection: Upgrade: Signals that the underlying TCP socket should remain open rather than closing.
  3. Sec-WebSocket-Key: A random 16-byte base64 string generated by the client.
  4. Sec-WebSocket-Accept: The server proves it understands RFC 6455 by concatenating the key with a standard GUID (258EAFA5-E914-47DA-95CA-C5AB0DC85B11), taking the SHA-1 hash, and returning it base64-encoded.
  5. HTTP 101 Switching Protocols: The server confirms the protocol change. From this moment onward, HTTP is discarded and raw binary/text frames stream over the TCP socket.

3. Real-Time Architecture with Socket.IO

While the native ws library provides a bare-bones implementation of RFC 6455, production applications almost universally adopt Socket.IO.

Why Socket.IO?

  1. Fallback Transports: If a corporate proxy blocks WebSocket traffic, Socket.IO automatically falls back to HTTP Long-Polling without breaking the application.
  2. Heartbeats & Auto-Reconnection: Automatically detects broken connections via ping/pong packets and reconnects with exponential backoff.
  3. Connection State Recovery: Can temporarily store packets during brief disconnects and replay them upon reconnection.
  4. Logical Routing: Provides built-in Namespaces (multiplexed paths) and Rooms (channel grouping).
BASH
pnpm add socket.io @socket.io/redis-adapter ioredis

4. Secure Authentication & Connection Middleware

A critical security vulnerability in real-time applications is allowing unauthenticated connections to establish sockets.

Never allow a client to connect and then send an authenticate event over the socket. If an attacker opens 10,000 unauthenticated sockets, your server's memory and file descriptors are exhausted before authentication logic ever runs.

Always authenticate during the initial HTTP handshake:

TS
// src/common/websocket/socket-auth.middleware.ts
import { Socket } from 'socket.io';
import jwt from 'jsonwebtoken';
import { config } from '../config/env';

export interface AuthenticatedSocket extends Socket {
  user?: {
    id: string;
    email: string;
    role: string;
  };
}

/**
 * Socket.IO middleware that verifies JWT tokens before accepting connections.
 */
export const socketAuthMiddleware = (socket: Socket, next: (err?: Error) => void) => {
  // 1. Extract token from handshake auth object or query params
  const token =
    socket.handshake.auth?.token ||
    socket.handshake.headers.authorization?.replace(/^Bearer\s+/i, '');

  if (!token) {
    return next(new Error('Authentication failed: Missing token'));
  }

  try {
    // 2. Cryptographically verify token
    const decoded = jwt.verify(token, config.JWT_ACCESS_SECRET) as any;

    // 3. Attach authenticated identity to socket instance
    (socket as AuthenticatedSocket).user = {
      id: decoded.sub,
      email: decoded.email,
      role: decoded.role,
    };

    next(); // Connection allowed!
  } catch (error) {
    next(new Error('Authentication failed: Invalid or expired token'));
  }
};

5. Logical Channels: Namespaces and Rooms

Socket.IO provides two essential architectural concepts for routing messages efficiently:

TEXT
                      Socket.IO Server
                             │
            ┌────────────────┴────────────────┐
            ▼                                 ▼
   Namespace: /problems              Namespace: /chat
            │                                 │
     ┌──────┴──────┐                   ┌──────┴──────┐
     ▼             ▼                   ▼             ▼
Room: "p_101"  Room: "p_102"     Room: "lobby"  Room: "help"
(Only users    (Only users       (Global chat)  (Support)
 on Two Sum)    on 3Sum)

1. Namespaces

Namespaces partition an application's socket logic over a single shared TCP connection. Instead of opening multiple WebSockets for chat, notifications, and code execution, they multiplex over paths:

  • io.of('/problems')
  • io.of('/notifications')

2. Rooms

Rooms are arbitrary, ephemeral channels that sockets can join and leave. They exist entirely in memory (and Redis) on the server:

  • When a user navigates to problem "two-sum", the client emits join_problem.
  • The server calls socket.join('problem:two-sum').
  • When a submission runs, the server broadcasts exclusively to that room: io.to('problem:two-sum').emit('submission_event', result).

3. Acknowledgments (Request-Response over WebSockets)

Unlike raw WebSocket messages which are strictly fire-and-forget, Socket.IO provides Acknowledgments, allowing you to implement two-way request-response semantics over a persistent socket connection:

  • The sender transmits an event and attaches a callback function (or awaits a Promise).
  • The receiver processes the event and executes the callback, returning data directly to the sender.
  • In Socket.IO v4.4+, you can use Promise-based timeouts via emitWithAck:
TS
// Client-side request with a 5000ms timeout guarantee:
try {
  const ack = await socket.timeout(5000).emitWithAck('run_test_case', {
    problemId: 'p_101',
    solutionCode: userCode,
  });
  console.log('Server confirmed test case queued:', ack.jobId);
} catch (err) {
  // Throws if the server does not acknowledge within 5000ms!
  console.error('Server timed out or unreachable');
}

Production Real-Time Manager Implementation

TS
// src/common/websocket/socket.server.ts
import { Server as HttpServer } from 'node:http';
import { Server as SocketIOServer } from 'socket.io';
import { socketAuthMiddleware, AuthenticatedSocket } from './socket-auth.middleware';
import { logger } from '../logger/logger';

export class RealtimeGateway {
  private static io: SocketIOServer;

  public static initialize(httpServer: HttpServer): SocketIOServer {
    this.io = new SocketIOServer(httpServer, {
      cors: {
        origin: ['http://localhost:3000', 'http://localhost:3002'],
        credentials: true,
      },
      pingTimeout: 20000,
      pingInterval: 25000,
    });

    // 🔒 Enforce authentication on all incoming connections
    this.io.use(socketAuthMiddleware);

    this.registerEventHandlers();
    return this.io;
  }

  private static registerEventHandlers(): void {
    this.io.on('connection', (rawSocket) => {
      const socket = rawSocket as AuthenticatedSocket;
      const user = socket.user!;

      logger.info({ socketId: socket.id, userId: user.id }, 'User connected to WebSocket');

      // Join a personal room for private notifications
      socket.join(`user:${user.id}`);

      // Handle joining collaborative problem workspace
      socket.on('join_problem', (problemSlug: string) => {
        const roomName = `problem:${problemSlug}`;
        socket.join(roomName);
        logger.info({ userId: user.id, roomName }, 'User joined problem room');

        // Notify others in the room
        socket.to(roomName).emit('user_joined', {
          userId: user.id,
          email: user.email,
        });
      });

      // Handle leaving problem workspace
      socket.on('leave_problem', (problemSlug: string) => {
        const roomName = `problem:${problemSlug}`;
        socket.leave(roomName);
        socket.to(roomName).emit('user_left', { userId: user.id });
      });

      // Handle test case execution with Acknowledgments
      socket.on('run_test_case', async (data: { problemSlug: string }, callback) => {
        logger.info({ userId: user.id, problem: data.problemSlug }, 'Test execution requested');
        // Acknowledge receipt back to the sender
        if (typeof callback === 'function') {
          callback({
            status: 'queued',
            jobId: `job_${Date.now()}`,
            timestamp: new Date().toISOString(),
          });
        }
      });

      // Handle disconnection
      socket.on('disconnect', (reason) => {
        logger.info({ socketId: socket.id, userId: user.id, reason }, 'User disconnected');
      });
    });
  }

  /**
   * Public helper to broadcast an event to a specific problem room from controllers or services.
   */
  public static broadcastToProblem(problemSlug: string, event: string, payload: any): void {
    if (this.io) {
      this.io.to(`problem:${problemSlug}`).emit(event, payload);
    }
  }

  /**
   * Public helper to dispatch a private notification to a specific user.
   */
  public static notifyUser(userId: string, event: string, payload: any): void {
    if (this.io) {
      this.io.to(`user:${userId}`).emit(event, payload);
    }
  }
}

6. Horizontal Scaling: The Redis Adapter

By default, Socket.IO stores room memberships and socket instances in local process RAM.

This creates a critical problem in a multi-container cluster:

TEXT
❌ The Multi-Node Problem without Redis Adapter:

  Client A (User 1)               Client B (User 2)
         │                               │
         ▼                               ▼
  ┌──────────────┐                ┌──────────────┐
  │ Node.js      │                │ Node.js      │
  │ Instance 1   │                │ Instance 2   │
  └──────────────┘                └──────────────┘
         │                               ▲
         │ (Emits to "room_123")         │ (NEVER RECEIVES IT!)
         ▼                               │
  Only User 1 gets it! ──────────────────┘

If Client A connects to Node Server 1, and Client B connects to Node Server 2, Server 1 cannot see Server 2's sockets in memory. A broadcast to "room_123" from Server 1 will never reach Client B!

The Solution: @socket.io/redis-adapter

The Redis Adapter connects every Node.js instance to a shared Redis Pub/Sub channel. When any node broadcasts to a room, it publishes the event to Redis, which forwards it to all other Node instances in the cluster:

TEXT
✅ Multi-Node Scaling with Redis Pub/Sub:

  Client A (User 1)                               Client B (User 2)
         │                                               │
         ▼                                               ▼
  ┌──────────────┐                                ┌──────────────┐
  │ Node Server 1│                                │ Node Server 2│
  └──────┬───────┘                                └──────▲───────┘
         │                                               │
         │ 1. Publish to Redis                           │ 3. Forward to Client B
         ▼                                               │
  ┌──────────────────────────────────────────────────────────────┐
  │                   Redis Pub/Sub Channel                      │
  │                    (Cluster Backplane)                       │
  └──────────────────────────────────────────────────────────────┘

Configuring the Redis Adapter

TS
import { createAdapter } from '@socket.io/redis-adapter';
import Redis from 'ioredis';

const pubClient = new Redis(process.env.REDIS_URL || 'redis://localhost:6379');
const subClient = pubClient.duplicate();

// Attach Redis Adapter to Socket.IO Server:
io.adapter(createAdapter(pubClient, subClient));

[!IMPORTANT] Sticky Sessions in Load Balancers: When using Socket.IO with HTTP long-polling fallback, your load balancer (AWS ALB, Nginx, Traefik) must have sticky sessions enabled (cookie-based affinity).

If the client sends Handshake Step 1 (HTTP polling) to Server 1, and Handshake Step 2 (HTTP polling upgrade) hits Server 2, Server 2 will throw a Session ID unknown error. Sticky sessions ensure that a given client's handshake requests always hit the same server until upgraded to WebSockets.


7. Server-Sent Events (SSE): The Lightweight Alternative

WebSockets are powerful, but they add operational complexity: custom protocols, stateful TCP sockets, sticky load balancer rules, and special proxy timeouts.

When you only need server-to-client streaming (such as AI chat token streaming, financial ticker prices, or test runner progress logs), Server-Sent Events (SSE) is simpler, faster, and more robust.

Why SSE is Great:

  1. Standard HTTP: Works over regular HTTP/1.1 or HTTP/2 without protocol upgrades.
  2. Built-in Browser Reconnection: The native browser EventSource API automatically reconnects if the connection drops.
  3. HTTP/2 Multiplexing: Multiple SSE streams multiplex over a single TCP connection without port limits.

Implementing an SSE Endpoint in Express

TS
// src/features/submission/submission-stream.controller.ts
import { Request, Response } from 'express';

export function streamSubmissionProgress(req: Request, res: Response) {
  const { submissionId } = req.params;

  // 1. Set required SSE headers
  res.setHeader('Content-Type', 'text/event-stream');
  res.setHeader('Cache-Control', 'no-cache, no-transform');
  res.setHeader('Connection', 'keep-alive');
  res.setHeader('X-Accel-Buffering', 'no'); // Disable Nginx proxy buffering!

  res.flushHeaders(); // Send headers immediately to establish stream

  // 2. Helper to write SSE formatted data
  const sendEvent = (event: string, data: any) => {
    res.write(`event: ${event}\n`);
    res.write(`data: ${JSON.stringify(data)}\n\n`); // Must end with double newline!
  };

  // 3. Emit initial event
  sendEvent('status', { phase: 'queued', progress: 0 });

  // 4. Listen to progress events from event emitter or Redis
  const interval = setInterval(() => {
    sendEvent('progress', {
      phase: 'compiling',
      timestamp: Date.now(),
    });
  }, 1000);

  // 5. 🔒 CRITICAL: Clean up listeners on client disconnect
  req.on('close', () => {
    clearInterval(interval);
    res.end();
  });
}

8. Production Real-Time Checklist & Summary

CODE
┌────────────────────────────────────────────────────────────────────────────┐
│                  PRODUCTION REAL-TIME ARCHITECTURE CHECKLIST               │
├────────────────────────────────────────────────────────────────────────────┤
│ [ ] Handshake Authentication: JWT verification enforced before socket      │
│     upgrade; unauthenticated sockets rejected at the HTTP boundary.        │
│                                                                            │
│ [ ] Redis Adapter Configured: Multi-node scaling active via Redis Pub/Sub; │
│     events broadcast reliably across all container replicas.               │
│                                                                            │
│ [ ] Sticky Sessions on Reverse Proxy: Load balancer configured with        │
│     cookie-based affinity to support Socket.IO polling handshakes.         │
│                                                                            │
│ [ ] Connection Cleanup: Event listeners (`req.on('close')`, `disconnect`)  │
│     cleaned up promptly to prevent memory leaks and dangling sockets.      │
│                                                                            │
│ [ ] SSE for Unidirectional Data: Server-Sent Events used for AI streaming  │
│     and status updates instead of maintaining heavy bi-directional sockets.│
│                                                                            │
│ [ ] Ping/Pong Timeouts: Heartbeat timers configured (`pingTimeout: 20000`) │
│     to detect half-open dead TCP sockets from dropped mobile networks.     │
└────────────────────────────────────────────────────────────────────────────┘

In the next chapter, we will master Chapter 20: Scaling Node.js: The Cluster Module, Worker Threads, and Multi-Core Architecture, building multi-process architectures, managing cluster lifecycles, and offloading heavy compute tasks to worker threads.

Finished this lesson?

Mark this chapter complete to update your learning streak and unlock the next lesson.