jvinhit//lab

Search posts

Type to search across journal entries.

navigate open esc close

Node.js Production Engineering 19 — Thiết kế Realtime Systems

Chọn polling, SSE hay WebSocket từ product contract; triển khai delivery, reconnect, authentication, presence, backpressure, scale-out và graceful draining với Node.js.

Một dashboard hiển thị “đơn hàng đã thanh toán”. Người dùng mất mạng ba giây đúng lúc event được phát. Kết nối tự reconnect, giao diện có chấm xanh trở lại — nhưng trạng thái vẫn là “đang xử lý” vì server không lưu event đã bỏ lỡ.

Kết nối lại thành công không đồng nghĩa khôi phục trạng thái thành công. Realtime production là bài toán delivery, ordering, recovery và capacity trước khi là bài toán chọn thư viện WebSocket.

Sau bài này, bạn sẽ có thể:

  • viết realtime contract từ nhu cầu sản phẩm thay vì mặc định WebSocket;
  • chọn polling, SSE, WebSocket hoặc Socket.IO theo hướng dữ liệu và recovery;
  • triển khai SSE/WebSocket có cleanup, size limit, schema validation và backpressure;
  • giải thích ordering và delivery guarantee thực sự của Socket.IO;
  • scale nhiều instance mà hiểu lúc nào cần sticky session và adapter;
  • thiết kế reconnect, presence, token expiry, deploy draining và observability.

1. Bắt đầu bằng realtime contract

Trước công nghệ, ghi rõ sáu câu hỏi:

  1. Hướng dữ liệu: chỉ server → client hay hai chiều?
  2. Freshness: 100 ms, 5 giây hay 1 phút vẫn chấp nhận?
  3. Delivery: có được mất event không? Duplicate có gây hại không?
  4. Ordering: theo một document/user hay toàn hệ thống?
  5. Recovery: reconnect cần replay event hay chỉ refetch snapshot mới nhất?
  6. Scale: bao nhiêu kết nối đồng thời, message/giây và bytes/giây?

Ví dụ:

Stock ticker       freshness < 1s, server→client, bỏ intermediate update được
Order status       không được kẹt trạng thái, reconnect có thể refetch snapshot
Collaborative doc  hai chiều, ordering theo document, cần conflict/recovery model
Presence           chấp nhận eventual, TTL, disconnect không luôn đáng tin

Hai feature đều “realtime” nhưng reliability contract hoàn toàn khác. Đừng bắt chat, price tick và payment status dùng cùng một delivery policy chỉ vì cùng đi qua socket.

2. Chọn transport nhỏ nhất đáp ứng contract

Cơ chếHướngĐiểm mạnhChi phí/giới hạnHợp với
Pollingclient → server định kỳđơn giản, cache/proxy quen thuộcrequest thừa, freshness theo intervaltrạng thái đổi chậm
Long pollingserver giữ request tới khi có dữ liệutương thích HTTP rộngnhiều request lifecyclefallback/legacy
SSEserver → clientHTTP text stream, browser reconnect, event idmột chiều, text, browser API hạn chế headernotification, feed, progress
WebSockethai chiềumột kết nối duplex, overhead message thấptự lo protocol, auth lifecycle, recoverychat, collaboration, game
Socket.IOhai chiều + abstractionroom, ack, reconnect, fallback, adapterskhông phải raw WebSocket protocol, thêm state/overheadứng dụng cần feature framework

Nếu client chỉ cần biết trạng thái mới nhất và độ trễ 10 giây chấp nhận được, polling có conditional request có thể là lựa chọn tốt hơn một connection sống nhiều giờ. Nếu chỉ server push, SSE thường ít moving part hơn WebSocket.

3. SSE đúng: stream, resume và proxy behavior

Native EventSource tự reconnect và gửi Last-Event-ID khi server cung cấp id. Muốn resume, server vẫn phải giữ event log hoặc map cursor tới dữ liệu bền.

import { once } from 'node:events';
import type { Request, Response } from 'express';

app.get('/events', async (req: Request, res: Response) => {
  const actor = await authenticateRequest(req);
  const after = req.get('last-event-id') ?? undefined;
  const abort = new AbortController();

  res.writeHead(200, {
    'Content-Type': 'text/event-stream; charset=utf-8',
    'Cache-Control': 'no-cache, no-transform',
    Connection: 'keep-alive',
    'X-Accel-Buffering': 'no', // Nginx: đừng buffer stream
  });
  res.flushHeaders();

  const heartbeat = setInterval(() => {
    res.write(': heartbeat\n\n'); // comment SSE, giữ proxy connection sống
  }, 15_000);

  // Cleanup theo response/socket đóng, không theo request body đã đọc xong.
  res.once('close', () => {
    clearInterval(heartbeat);
    abort.abort();
  });

  try {
    for await (const event of eventStore.subscribe({
      actorId: actor.id,
      after,
      signal: abort.signal,
    })) {
      const frame = [
        `id: ${event.id}`,
        `event: ${event.type}`,
        `data: ${JSON.stringify(event.data)}`,
        '',
        '',
      ].join('\n');

      if (!res.write(frame)) await once(res, 'drain');
    }
  } catch (error) {
    if (!abort.signal.aborted) logger.error({ error }, 'SSE stream failed');
  }
});

res.write() trả false nghĩa là buffer đã vượt high-water mark; chờ drain là backpressure. Heartbeat trong code minh họa chưa phối hợp drain riêng; ở throughput cao, gom mọi write qua một connection writer để không có hai producer ghi đồng thời.

Native EventSource trong browser không cho đặt arbitrary Authorization header. Thường dùng cookie HttpOnly cùng kiểm tra Origin, hoặc signed URL sống rất ngắn — không đặt access token dài hạn trong query vì URL dễ rơi vào log/history.

4. Raw WebSocket với ws: boundary phải phòng thủ

WebSocket bắt đầu bằng HTTP upgrade rồi trở thành message channel dài hạn. Browser gửi Origin; server phải allowlist để chống cross-site WebSocket hijacking khi auth dựa trên cookie.

import { WebSocket, WebSocketServer, type RawData } from 'ws';
import { z } from 'zod';

const incomingMessage = z.discriminatedUnion('type', [
  z.object({
    type: z.literal('chat.send'),
    roomId: z.string(),
    body: z.string().max(2_000),
  }),
  z.object({
    type: z.literal('typing.set'),
    roomId: z.string(),
    active: z.boolean(),
  }),
]);

const wss = new WebSocketServer({
  noServer: true,
  maxPayload: 64 * 1024,
  perMessageDeflate: false,
});
const actorBySocket = new WeakMap<WebSocket, Actor>();

httpServer.on('upgrade', async (request, socket, head) => {
  try {
    assertAllowedOrigin(request.headers.origin);
    const actor = await authenticateUpgrade(request);

    wss.handleUpgrade(request, socket, head, (ws) => {
      actorBySocket.set(ws, actor);
      wss.emit('connection', ws, request);
    });
  } catch {
    socket.write('HTTP/1.1 401 Unauthorized\r\nConnection: close\r\n\r\n');
    socket.destroy();
  }
});

wss.on('connection', (ws: WebSocket) => {
  const actor = actorBySocket.get(ws);
  if (!actor) return ws.close(1011, 'missing connection context');

  let chain = Promise.resolve();
  let queued = 0;

  async function handleIncoming(raw: RawData): Promise<void> {
    let input: unknown;
    try {
      input = JSON.parse(raw.toString());
    } catch {
      return ws.close(1008, 'invalid JSON');
    }

    const parsed = incomingMessage.safeParse(input);
    if (!parsed.success) return ws.close(1008, 'invalid message');

    if (!(await roomPolicy.canPublish(actor.id, parsed.data.roomId))) {
      return ws.close(1008, 'forbidden');
    }

    await handleMessage(actor, parsed.data);
  }

  ws.on('message', (raw) => {
    if (queued >= 32) return ws.close(1013, 'message backlog exceeded');
    queued += 1;

    // EventEmitter không await async listener. Nối promise để giữ thứ tự
    // theo connection, bắt rejection và giới hạn backlog trong bộ nhớ.
    chain = chain
      .then(() => handleIncoming(raw))
      .catch((error) => {
        logger.error({ error, actorId: actor.id }, 'WebSocket message failed');
        ws.close(1011, 'message processing failed');
      })
      .finally(() => {
        queued -= 1;
      });
  });
});

Đoạn trên rút gọn để tập trung boundary. Bản production còn cần rate limit theo actor/connection, heartbeat, structured error, trace context và shutdown. Không bao giờ tin userId, role hoặc room membership trong payload client.

perMessageDeflate có lợi với payload lớn/lặp, nhưng dùng CPU/memory và có rủi ro amplification; đo trước khi bật. maxPayload mặc định của thư viện có thể lớn hơn business cần rất nhiều, nên đặt rõ.

5. Socket.IO: protocol ứng dụng có sẵn primitive

Socket.IO thêm event, ack, room, reconnect và transport fallback. Nó không phải raw WebSocket server; client WebSocket thuần không nói chuyện trực tiếp với Socket.IO protocol.

interface ClientToServerEvents {
  'chat:send': (
    input: { roomId: string; clientMessageId: string; body: string },
    ack: (
      result: { ok: true; eventId: string } | { ok: false; code: string }
    ) => void
  ) => void;
}

interface ServerToClientEvents {
  'chat:message': (event: ChatEvent) => void;
}

interface SocketData {
  actor: Actor;
}

const io = new Server<
  ClientToServerEvents,
  ServerToClientEvents,
  {},
  SocketData
>(httpServer, {
  cors: { origin: allowedOrigins, credentials: true },
  maxHttpBufferSize: 64 * 1024,
  connectionStateRecovery: {
    maxDisconnectionDuration: 2 * 60 * 1000,
    skipMiddlewares: false,
  },
});

io.use(async (socket, next) => {
  try {
    socket.data.actor = await authenticateHandshake(socket.handshake);
    next();
  } catch {
    next(new Error('unauthorized'));
  }
});

io.on('connection', (socket) => {
  socket.on('chat:send', async (input, ack) => {
    if (!(await roomPolicy.canPublish(socket.data.actor.id, input.roomId))) {
      return ack({ ok: false, code: 'FORBIDDEN' });
    }

    const { event, created } = await chat.appendIdempotently({
      ...input,
      actorId: socket.data.actor.id,
    });

    if (created) io.to(`room:${input.roomId}`).emit('chat:message', event);
    ack({ ok: true, eventId: event.id });
  });
});

clientMessageId cho phép server deduplicate khi client retry vì không nhận ack. Khi store trả event cũ, server không broadcast lại; recipient vẫn deduplicate theo event.id để phòng duplicate từ các đường delivery khác. Ack chỉ chứng minh server callback đã trả, không mặc định chứng minh mọi recipient đã render message.

6. Ordering và delivery: ghi rõ phạm vi

Socket.IO đảm bảo thứ tự event đã đến trên một connection, kể cả lúc transport upgrade. Mặc định message arrival là at-most-once: connection đứt giữa lúc gửi có thể mất event, và server không tự lưu mọi event cho client đang offline.

Để server → client có khả năng replay:

1. Event có id tăng theo stream/document.
2. Server persist event trước khi emit.
3. Client chỉ lưu offset sau khi apply thành công.
4. Reconnect gửi offset cuối.
5. Server replay event sau offset hoặc yêu cầu refetch snapshot.
6. Client deduplicate theo event id.
io.on('connection', async (socket) => {
  const offset = socket.handshake.auth.offset as string | undefined;

  for (const event of await chat.eventsAfter(offset, { limit: 1_000 })) {
    socket.emit('chat:message', event);
  }
});

Connection state recovery của Socket.IO có thể phục hồi room/data và missed packet trong một cửa sổ cấu hình, nhưng recovery không phải lúc nào cũng thành công. Nó hoạt động với in-memory adapter mặc định; Redis Pub/Sub adapter chuẩn không hỗ trợ feature này. Client vẫn cần fallback full resync.

Ordering toàn cục rất đắt và hiếm khi cần. Thường chỉ cần ordering theo roomId/document/aggregate. Khi nhiều producer ghi cùng stream, dùng sequence/version do một owner cấp hoặc conflict-resolution model phù hợp.

7. Scale-out: load balancing và message forwarding là hai việc khác nhau

Khi có nhiều instance:

client A ─▶ instance 1 ─┐
                        ├─ adapter/pub-sub ─▶ broadcast tới room ở mọi instance
client B ─▶ instance 2 ─┘

Adapter chuyển event giữa instance; nó không mặc định lưu lịch sử cho offline recovery.

import { createAdapter } from '@socket.io/redis-adapter';
import { createClient } from 'redis';

const pubClient = createClient({ url: process.env.REDIS_URL });
const subClient = pubClient.duplicate();

await Promise.all([pubClient.connect(), subClient.connect()]);
io.adapter(createAdapter(pubClient, subClient));

Sticky session: khi nào cần?

  • Socket.IO còn dùng HTTP long-polling: các HTTP request của cùng session phải tới instance đã tạo session → cần sticky session hoặc cơ chế đồng bộ session tương đương.
  • Chỉ dùng WebSocket/WebTransport: một connection sống ở một instance → round-robin handshake được, không cần sticky session cho các frame sau.

Tắt polling giảm compatibility ở mạng/proxy hạn chế. Đây là trade-off cần test với người dùng thật, không phải tối ưu mặc định.

Redis Pub/Sub adapter không hỗ trợ Connection State Recovery hay replay khi Redis/service ngắt. Nếu recovery là requirement, dùng Redis Streams adapter tương thích hoặc event store riêng và kiểm ma trận feature/version của adapter.

8. Presence là lease, không phải boolean tuyệt đối

Disconnect event không luôn tới: laptop ngủ, NAT timeout hoặc process chết. Presence production thường là lease có TTL:

online nếu heartbeat/lease được renew trong 30s
away nếu quá 30s
offline nếu quá 2 phút

Một user có nhiều tab/device, nên userId → boolean không đủ. Lưu nhiều connection/session lease và suy ra trạng thái user. Presence là eventual; UI nên diễn đạt “hoạt động gần đây” thay vì hứa chính xác tuyệt đối.

Không ghi một Redis write cho mọi heartbeat của hàng triệu connection mà chưa capacity plan. Có thể aggregate theo process, dùng sorted set/bucket hoặc chỉ bật presence cho room đang quan sát.

9. Backpressure và slow consumer

Nếu server tạo 1.000 message/giây nhưng client đọc 100, buffer tăng cho tới khi memory hoặc latency vỡ. Policy tùy loại dữ liệu:

  • coalesce: chỉ giữ price/progress mới nhất;
  • drop: bỏ typing indicator cũ;
  • batch: gom analytics update trong 100 ms;
  • disconnect: đóng client quá chậm để bảo vệ hệ thống;
  • persist + resume: chat/order event không được mất.

Với ws, theo dõi bufferedAmount; với stream API, tôn trọng write()/drain. Đặt per-connection outbound queue hữu hạn và metric slow-consumer disconnect. “Gửi thành công vào buffer” không đồng nghĩa client đã nhận.

10. Security lifecycle của kết nối dài

Handshake auth chỉ là thời điểm đầu. Access token có thể hết hạn, user bị khóa hoặc permission room thay đổi trong khi socket vẫn sống.

Baseline:

  • TLS (wss://) và allowlist Origin;
  • cookie SameSite/CSRF model hoặc token sống ngắn trong handshake;
  • re-auth/disconnect khi token hết hạn hay session bị revoke;
  • authorization trên mỗi loại event/resource, không chỉ lúc connect;
  • schema, size, rate và concurrency limit;
  • không phản chiếu error nội bộ/stack;
  • audit event nhạy cảm;
  • quota connection theo account/IP có cân nhắc NAT.

CORS config của Socket.IO không thay thế authorization và không bảo vệ raw WebSocket theo cách nhiều người kỳ vọng. Origin validation ở handshake mới xử lý browser cross-site case; non-browser client có thể giả Origin, nên vẫn cần authentication thực.

11. Heartbeat, reconnect và deploy draining

Heartbeat phát hiện half-open connection. Reconnect dùng exponential backoff + jitter để tránh hàng trăm nghìn client cùng kết nối lại sau deploy.

Khi deploy:

  1. instance chuyển ready = false để load balancer không nhận connection mới;
  2. gửi signal “server draining” hoặc đóng với code/reason đã quy ước;
  3. client reconnect có jitter sang instance khỏe;
  4. chờ connection cũ đóng tới deadline rồi terminate;
  5. client resume bằng offset hoặc refetch snapshot.

Nếu đóng toàn fleet cùng lúc, reconnect storm có thể hạ bản release mới. Rolling/canary và connection drain là phần của deployment plan realtime.

12. Hợp đồng với frontend

Team cần thống nhất một event envelope:

type EventEnvelope<T> = {
  id: string;
  type: string;
  version: number;
  occurredAt: string;
  streamId: string;
  sequence: number;
  data: T;
};

Và state machine client:

connecting → live → reconnecting → recovering → live
                                  └─ recovery failed → resyncing → live

UI không nên chỉ có isConnected. recovering/stale giúp tránh hiển thị dữ liệu cũ như dữ liệu live. Mutation gửi qua socket cần clientMessageId, ack timeout và UX cho trạng thái pending/failed.

13. Capacity và observability

Theo dõi:

  • active connection theo instance/region/transport;
  • connect, auth failure, reconnect và recovery success rate;
  • inbound/outbound message + bytes/giây;
  • event processing/ack latency;
  • per-connection buffered bytes, slow-consumer drop/disconnect;
  • room size và adapter publish latency;
  • event replay gap, oldest unrecovered offset;
  • event-loop lag, RSS và file descriptor usage.

Connection count không đủ. 50.000 connection idle có thể rẻ hơn 5.000 connection phát payload lớn 20 lần/giây. Load test phải mô phỏng message distribution, reconnect storm và slow client, không chỉ mở socket.

14. Checklist trước khi ship

  • Product contract ghi freshness, delivery, ordering và recovery.
  • Transport là lựa chọn nhỏ nhất đáp ứng contract.
  • Message có schema/version/size limit; authorization theo resource.
  • Có event id/offset hoặc chiến lược full resync khi reconnect.
  • Duplicate và out-of-order được test.
  • Backpressure policy phân biệt event được drop và event phải persist.
  • Sticky session chỉ bật khi transport/session model cần.
  • Adapter, event store và recovery responsibility được tách rõ.
  • Token expiry/revocation và Origin validation có test.
  • Deploy drain và reconnect storm có runbook/metric.

15. Bài thực hành: order status + support chat

Dùng hai feature để buộc hai contract khác nhau:

  1. Order status bằng SSE: event id, Last-Event-ID, heartbeat, cleanup trên res.once('close', ...), Nginx buffering off.
  2. Support chat bằng Socket.IO: room authorization, clientMessageId, ack, event persistence và offset replay.
  3. Chạy hai instance với Redis adapter; test polling có/không sticky session và WebSocket-only.
  4. Kill connection giữa lúc gửi; refresh tab; tắt Redis; deploy rolling. Ghi expected state/recovery cho từng case.
  5. Tạo slow client và chứng minh outbound queue bị giới hạn.
  6. Dashboard recovery success, oldest gap, buffered bytes và reconnect rate.

Bài đạt khi UI không nói “live” trong lúc đang recover và không mất business state sau một disconnect có chủ đích.

Tài liệu chính thức

Phần tiếp theo

Realtime và microservices tạo nhiều failure path mà log rời rạc không thể giải thích. Phần cuối xây observability từ SLI/SLO: structured logs, metrics có cardinality kiểm soát, distributed traces, sampling và OpenTelemetry Collector — rồi dùng chúng trong một incident workflow hoàn chỉnh.