Поделиться
Поделиться

Real-time функции — уведомления, чаты, обновление лент, совместное редактирование — стали нормой для современных веб-приложений. Но за кажущейся простотой «просто откройте WebSocket» скрываются нетривиальные вопросы масштабирования, надёжности и безопасности. Разберём всё по порядку.

WebSocket vs SSE vs Long-Polling

Прежде чем открывать WebSocket, стоит убедиться, что это правильный инструмент.

| Параметр | Long-Polling | SSE | WebSocket | |---|---|---|---| | Направление | Только сервер → клиент | Только сервер → клиент | Двустороннее | | Накладные расходы | Высокие (HTTP запрос на каждое сообщение) | Низкие | Минимальные | | Автоматический reconnect | Реализуется вручную | Встроен в EventSource | Реализуется вручную | | Поддержка браузеров | Универсальная | IE не поддерживается | IE10+ | | Прокси/балансировщики | Прозрачно | Нужен X-Accel-Buffering: no | Нужен upgrade | | Подходит для | Редкие обновления, legacy среды | Уведомления, потоки событий | Чат, игры, совместная работа |

Когда не нужен WebSocket

Если клиент только получает данные — новости, уведомления, статус заказа — SSE проще и не требует специального сервера:

// Клиент: SSE
const es = new EventSource('/api/notifications');

es.addEventListener('order_update', (event) => {
  const data = JSON.parse(event.data);
  updateOrderStatus(data);
});

es.addEventListener('error', () => {
  // EventSource автоматически переподключится через 3 секунды
});

// Сервер (Node.js):
app.get('/api/notifications', (req, res) => {
  res.writeHead(200, {
    'Content-Type': 'text/event-stream',
    'Cache-Control': 'no-cache',
    'X-Accel-Buffering': 'no', // важно для nginx
    'Connection': 'keep-alive',
  });

  const sendEvent = (type, data) => {
    res.write(`event: ${type}\ndata: ${JSON.stringify(data)}\n\n`);
  };

  // Подписка на события
  const unsubscribe = eventBus.on('notification', (event) => {
    if (event.userId === req.user.id) {
      sendEvent(event.type, event.data);
    }
  });

  req.on('close', () => {
    unsubscribe();
  });
});

Socket.io: комнаты и пространства имён

Socket.io — самая популярная библиотека для WebSocket в Node.js экосистеме. Две ключевые концепции для организации соединений:

Namespaces (пространства имён)

Разделяют соединения по функциональным областям. Это отдельные логические каналы поверх одного TCP-соединения:

const io = require('socket.io')(server);

// Namespace для чата
const chat = io.of('/chat');
chat.on('connection', (socket) => {
  console.log('Chat user connected');
});

// Namespace для игры
const game = io.of('/game');
game.on('connection', (socket) => {
  console.log('Game user connected');
});
// Клиент подключается к конкретному namespace
const chatSocket = io('/chat');
const gameSocket = io('/game');

Rooms (комнаты)

Rooms — динамические группы внутри namespace. Идеальны для чатов, игровых лобби, многопользовательских документов:

io.on('connection', (socket) => {
  // Присоединение к комнате
  socket.on('join_room', async (roomId) => {
    // Авторизация
    const canJoin = await checkRoomAccess(socket.userId, roomId);
    if (!canJoin) return socket.emit('error', 'access_denied');

    socket.join(roomId);
    // Уведомление других участников
    socket.to(roomId).emit('user_joined', { userId: socket.userId });
    // Подтверждение инициатору
    socket.emit('room_joined', { roomId });
  });

  // Сообщение в комнату
  socket.on('message', (data) => {
    const { roomId, text } = data;
    // Проверяем, что пользователь в комнате
    if (!socket.rooms.has(roomId)) return;
    
    const msg = { id: uuid(), userId: socket.userId, text, ts: Date.now() };
    // Всем в комнате, кроме отправителя
    socket.to(roomId).emit('message', msg);
    // Самому отправителю (с подтверждением)
    socket.emit('message_sent', msg);
  });

  // Выход из комнаты
  socket.on('leave_room', (roomId) => {
    socket.leave(roomId);
    socket.to(roomId).emit('user_left', { userId: socket.userId });
  });

  // Автоматический выход при дисконнекте
  socket.on('disconnect', () => {
    socket.rooms.forEach((room) => {
      socket.to(room).emit('user_left', { userId: socket.userId });
    });
  });
});

Горизонтальное масштабирование через Redis Adapter

Одна нода Node.js может обслуживать ~10,000–50,000 одновременных WebSocket соединений. При большей нагрузке нужно несколько нод. Проблема: если пользователи A и B подключены к разным нодам, как они получают сообщения из одной комнаты?

Решение — Redis Pub/Sub adapter:

const { createServer } = require('http');
const { Server } = require('socket.io');
const { createAdapter } = require('@socket.io/redis-adapter');
const { createClient } = require('redis');

const httpServer = createServer(app);
const io = new Server(httpServer, {
  cors: { origin: process.env.FRONTEND_URL, credentials: true }
});

// Два Redis клиента: pub и sub
const pubClient = createClient({ url: process.env.REDIS_URL });
const subClient = pubClient.duplicate();

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

// Подключаем Redis adapter
io.adapter(createAdapter(pubClient, subClient));

httpServer.listen(3000);

Теперь io.to(roomId).emit(...) работает корректно независимо от того, к какой ноде подключён получатель — Redis Pub/Sub маршрутизирует сообщение.

[Клиент A] ──→ [Нода 1]
                  │ pub: room:xyz → message
               [Redis]
                  │ sub: room:xyz → message
[Клиент B] ──→ [Нода 2] ──→ emit to Клиент B

Sticky sessions

При использовании нескольких нод за балансировщиком необходимы sticky sessions (IP или cookie hash), чтобы handshake и все запросы одного клиента приходили на одну ноду:

# Nginx sticky session по IP
upstream socketio_nodes {
    ip_hash;
    server node1:3000;
    server node2:3000;
    server node3:3000;
}

server {
    location /socket.io/ {
        proxy_pass http://socketio_nodes;
        proxy_http_version 1.1;
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection "upgrade";
        proxy_set_header X-Real-IP $remote_addr;
    }
}

Reconnect с exponential backoff

Нестабильные соединения — нормальная ситуация на мобильных сетях. Клиент должен переподключаться с exponential backoff, чтобы не устроить DDoS серверу при его перезапуске.

class ReliableSocket {
  constructor(url, options = {}) {
    this.url = url;
    this.maxRetries = options.maxRetries || 10;
    this.baseDelay = options.baseDelay || 1000;
    this.maxDelay = options.maxDelay || 30000;
    this.retryCount = 0;
    this.socket = null;
    this.listeners = new Map();
  }

  connect() {
    this.socket = io(this.url, {
      reconnection: false, // управляем переподключением сами
      timeout: 10000,
    });

    this.socket.on('connect', () => {
      this.retryCount = 0;
      console.log('Connected');
      this.emit('connect');
    });

    this.socket.on('disconnect', (reason) => {
      console.log(`Disconnected: ${reason}`);
      this.emit('disconnect', reason);
      if (reason !== 'io client disconnect') {
        this.scheduleReconnect();
      }
    });

    this.socket.on('connect_error', () => {
      this.scheduleReconnect();
    });
  }

  scheduleReconnect() {
    if (this.retryCount >= this.maxRetries) {
      this.emit('give_up');
      return;
    }
    // Exponential backoff с jitter
    const delay = Math.min(
      this.baseDelay * (2 ** this.retryCount) + Math.random() * 1000,
      this.maxDelay
    );
    this.retryCount++;
    console.log(`Reconnect attempt ${this.retryCount} in ${Math.round(delay)}ms`);
    setTimeout(() => this.connect(), delay);
  }
}

Jitter (случайная добавка) критически важен: без него все клиенты после рестарта сервера будут переподключаться одновременно и снова перегрузят его.

Heartbeat: детектирование зависших соединений

TCP-соединения могут «зависнуть» — технически открыты, но данные не передаются. Heartbeat (ping/pong) позволяет детектировать и закрывать такие соединения.

// Socket.io настраивает ping/pong автоматически
const io = new Server(httpServer, {
  pingInterval: 25000, // сервер пингует каждые 25 сек
  pingTimeout: 20000,  // ожидание pong — 20 сек, потом дисконнект
});

Для кастомного heartbeat (например, с полезной нагрузкой):

// Сервер отправляет heartbeat с серверным временем
setInterval(() => {
  io.emit('heartbeat', { serverTime: Date.now() });
}, 30000);

// Клиент отвечает и может синхронизировать время
socket.on('heartbeat', (data) => {
  const latency = Date.now() - data.serverTime;
  socket.emit('heartbeat_ack', { latency });
});

Аутентификация WebSocket соединений

WebSocket не поддерживает кастомные заголовки при установке соединения (только Cookie и Authorization в некоторых реализациях). Распространённые подходы:

Подход 1: Token в query string (не рекомендуется для продакшна)

// Клиент
const socket = io({ auth: { token: localStorage.getItem('jwt') } });

// Сервер — middleware
io.use((socket, next) => {
  const token = socket.handshake.auth.token;
  try {
    const payload = jwt.verify(token, process.env.JWT_SECRET);
    socket.userId = payload.sub;
    next();
  } catch (err) {
    next(new Error('unauthorized'));
  }
});

Подход 2: Cookie (предпочтительно для браузеров)

// HTTP-only cookie устанавливается при логине
res.cookie('session', sessionToken, {
  httpOnly: true,
  secure: true,
  sameSite: 'strict',
  maxAge: 86400000
});

// Сервер читает cookie из handshake
io.use(async (socket, next) => {
  const cookies = cookie.parse(socket.handshake.headers.cookie || '');
  const sessionToken = cookies.session;
  if (!sessionToken) return next(new Error('unauthorized'));

  const session = await sessionStore.get(sessionToken);
  if (!session) return next(new Error('unauthorized'));

  socket.userId = session.userId;
  next();
});

Паттерны обработки сообщений

Acknowledgements (подтверждения)

// Клиент отправляет сообщение и ждёт подтверждения
socket.emit('send_message', messageData, (ack) => {
  if (ack.success) {
    markMessageDelivered(messageData.id);
  } else {
    showError(ack.error);
  }
});

// Сервер подтверждает
socket.on('send_message', async (data, callback) => {
  try {
    const msg = await saveMessage(data);
    callback({ success: true, messageId: msg.id });
    socket.to(data.roomId).emit('new_message', msg);
  } catch (err) {
    callback({ success: false, error: err.message });
  }
});

Rate limiting

const rateLimit = new Map();

io.use((socket, next) => {
  const userId = socket.userId;
  const now = Date.now();
  const userLimit = rateLimit.get(userId) || { count: 0, resetAt: now + 1000 };

  if (now > userLimit.resetAt) {
    userLimit.count = 0;
    userLimit.resetAt = now + 1000;
  }

  userLimit.count++;
  rateLimit.set(userId, userLimit);

  if (userLimit.count > 30) { // 30 сообщений в секунду
    return next(new Error('rate_limit_exceeded'));
  }
  next();
});

Итог

WebSocket и real-time архитектура — это не только открытие соединения. Надёжная система требует: правильного выбора между WS/SSE/polling, грамотной организации через комнаты и namespace, Redis adapter для горизонтального масштабирования, exponential backoff с jitter при переподключении и аутентификации на уровне handshake. Для геолокационных приложений с real-time синхронизацией также смотрите статью Location-based приложения и геоигры.