Comunicação em Tempo Real com WebSockets no Node.js
Protocolo WebSocket explicado, servidor ws integrado ao Express, autenticação de conexões via JWT, rooms por ID de usuário, heartbeat/ping-pong e escalonamento horizontal com Redis Pub/Sub.
HTTP é um protocolo request-response: o cliente pergunta, o servidor responde, a conexão se encerra. Para notificações em tempo real — mensagens de chat, status de pedido atualizado, lances em um leilão — essa arquitetura não funciona. O servidor precisa empurrar dados para o cliente sem que ele pergunte.
WebSocket é um protocolo diferente do HTTP: após um handshake inicial (feito sobre HTTP), ele estabelece uma conexão TCP persistente, bidirecional e de baixa latência. O cliente e o servidor podem se enviar dados a qualquer momento sem overhead de headers HTTP em cada mensagem. Neste artigo vamos do protocolo até uma implementação production-ready com autenticação JWT, rooms por usuário e escalonamento horizontal.
HTTP Polling vs SSE vs WebSocket: Quando Usar Cada Um
- Long Polling — o cliente faz uma requisição HTTP que fica em aberto até o servidor ter algo para responder. Simples, mas custoso: reconecta após cada mensagem. Use apenas quando WebSocket não é viável.
- Server-Sent Events (SSE) — conexão HTTP unidirecional: servidor → cliente. Ideal para feeds de notificação, logs em tempo real, progress bars. Mais simples que WebSocket para casos unidirecionais.
- WebSocket — bidirecional e de baixa latência. Use para chat, jogos online, editores colaborativos, dashboards com input do usuário, leilões, trading.
Configuração do Servidor WebSocket com ws
A integração do ws com o servidor HTTP do Express é fundamental: a mesma porta (ex: 3333) atende tanto requisições HTTP normais quanto conexões WebSocket. O browser faz o upgrade da conexão HTTP para WebSocket automaticamente via o header Upgrade: websocket. Autenticar a conexão no momento do handshake é crítico: sem isso, qualquer pessoa pode se conectar ao WebSocket e receber mensagens privadas de outros usuários.
npm install ws jsonwebtoken
npm install -D @types/wsimport { WebSocketServer, WebSocket, type RawData } from 'ws';
import http from 'http';
import { verify } from 'jsonwebtoken';
// Extende WebSocket para adicionar propriedades de negócio
interface AuthenticatedWebSocket extends WebSocket {
userId: string;
isAlive: boolean; // Para heartbeat
}
// Map<userId, Set<conexões>> — um usuário pode ter múltiplas abas abertas
type ClientRegistry = Map<string, Set<AuthenticatedWebSocket>>;
export class AppWebSocketServer {
private wss: WebSocketServer;
private clients: ClientRegistry = new Map();
private heartbeatInterval: ReturnType<typeof setInterval> | null = null;
constructor(httpServer: http.Server) {
this.wss = new WebSocketServer({
server: httpServer,
// Valida o upgrade HTTP antes de aceitar a conexão WebSocket
verifyClient: ({ req }, callback) => {
const token = new URL(
req.url ?? '/',
`http://${req.headers.host}`
).searchParams.get('token');
if (!token) {
callback(false, 401, 'Token não fornecido');
return;
}
try {
const payload = verify(token, process.env.JWT_SECRET as string) as { sub: string };
// Armazena no objeto da requisição para usar em 'connection'
(req as unknown as Record<string, unknown>).userId = payload.sub;
callback(true);
} catch {
callback(false, 401, 'Token inválido ou expirado');
}
},
});
this.wss.on('connection', this.handleConnection.bind(this));
this.startHeartbeat();
}
private handleConnection(
ws: AuthenticatedWebSocket,
req: http.IncomingMessage
): void {
const userId = (req as unknown as Record<string, string>).userId;
ws.userId = userId;
ws.isAlive = true;
// Registra a conexão no registry
if (!this.clients.has(userId)) {
this.clients.set(userId, new Set());
}
this.clients.get(userId)!.add(ws);
console.log(`[WS] Usuário ${userId} conectado. Total: ${this.wss.clients.size}`);
// Responde ao pong do heartbeat
ws.on('pong', () => {
ws.isAlive = true;
});
// Processa mensagens recebidas do cliente
ws.on('message', (data: RawData) => {
this.handleMessage(ws, data);
});
// Remove do registry ao desconectar
ws.on('close', () => {
const userSockets = this.clients.get(userId);
if (userSockets) {
userSockets.delete(ws);
if (userSockets.size === 0) {
this.clients.delete(userId);
}
}
console.log(`[WS] Usuário ${userId} desconectado.`);
});
// Envia confirmação de conexão
this.sendToSocket(ws, { type: 'CONNECTED', userId });
}
private handleMessage(ws: AuthenticatedWebSocket, raw: RawData): void {
try {
const message = JSON.parse(raw.toString()) as { type: string; payload?: unknown };
switch (message.type) {
case 'PING':
this.sendToSocket(ws, { type: 'PONG' });
break;
// Adicione mais handlers conforme necessário
default:
console.warn(`[WS] Tipo de mensagem desconhecido: ${message.type}`);
}
} catch {
ws.close(1008, 'Mensagem inválida');
}
}
// Envia para um socket específico
sendToSocket(ws: WebSocket, data: unknown): void {
if (ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify(data));
}
}
// Envia para todas as conexões de um usuário específico (todas as abas)
sendToUser(userId: string, data: unknown): void {
const userSockets = this.clients.get(userId);
if (!userSockets) return;
const payload = JSON.stringify(data);
for (const ws of userSockets) {
if (ws.readyState === WebSocket.OPEN) {
ws.send(payload);
}
}
}
// Broadcast para todos os usuários conectados
broadcast(data: unknown, excludeUserId?: string): void {
const payload = JSON.stringify(data);
for (const [userId, sockets] of this.clients.entries()) {
if (userId === excludeUserId) continue;
for (const ws of sockets) {
if (ws.readyState === WebSocket.OPEN) {
ws.send(payload);
}
}
}
}
// Heartbeat: detecta e desconecta clientes mortos
private startHeartbeat(): void {
this.heartbeatInterval = setInterval(() => {
this.wss.clients.forEach((ws) => {
const socket = ws as AuthenticatedWebSocket;
if (!socket.isAlive) {
// Não respondeu ao ping anterior — conexão morta
socket.terminate();
return;
}
socket.isAlive = false;
socket.ping();
});
}, 30_000); // Verifica a cada 30 segundos
this.wss.on('close', () => {
if (this.heartbeatInterval) clearInterval(this.heartbeatInterval);
});
}
}O heartbeat (ping/pong) é essencial em produção. Firewalls, load balancers e proxies reversos fecham conexões TCP ociosas após 60-90 segundos. Sem heartbeat, o servidor pensa que o cliente ainda está conectado, mas as mensagens nunca chegam. O mecanismo de ping/pong do protocolo WebSocket nativo é mais eficiente que enviar mensagens JSON de keep-alive.
Integrando com a Aplicação Express
import express from 'express';
import http from 'http';
import { AppWebSocketServer } from './infra/websocket/WebSocketServer';
const app = express();
const httpServer = http.createServer(app);
// Singleton: instância única do WebSocket Server compartilhada
export const wsServer = new AppWebSocketServer(httpServer);
// Middlewares e rotas do Express...
app.use(express.json());
app.use('/api', apiRoutes);
httpServer.listen(3333, () => {
console.log('HTTP + WebSocket server rodando na porta 3333');
});Disparando Eventos de Negócio em Tempo Real
A integração com os Use Cases é simples: após a operação de banco, emite o evento WebSocket para os usuários relevantes:
import { wsServer } from '../../../server';
export class CreateOrderUseCase {
async execute(dto: CreateOrderDTO): Promise<Order> {
const order = await this.ordersRepo.create(dto);
// Notifica o cliente que fez o pedido em todas as suas abas
wsServer.sendToUser(dto.customerId, {
type: 'ORDER_CREATED',
data: {
orderId: order.id,
status: order.status,
total: order.total,
},
});
// Notifica todos os operadores do dashboard
wsServer.broadcast(
{ type: 'NEW_ORDER_RECEIVED', data: { orderId: order.id } },
dto.customerId // Exclui o próprio cliente
);
return order;
}
}Cliente React: Hook useWebSocket
import { useEffect, useRef, useCallback } from 'react';
type MessageHandler = (event: { type: string; data?: unknown }) => void;
export function useWebSocket(accessToken: string | null, onMessage: MessageHandler) {
const wsRef = useRef<WebSocket | null>(null);
const reconnectTimeoutRef = useRef<ReturnType<typeof setTimeout> | null>(null);
const reconnectAttempts = useRef(0);
const connect = useCallback(() => {
if (!accessToken) return;
// Token enviado como query param para passar pelo verifyClient
const url = `${process.env.NEXT_PUBLIC_WS_URL}?token=${accessToken}`;
const ws = new WebSocket(url);
wsRef.current = ws;
ws.onopen = () => {
console.log('[WS] Conectado');
reconnectAttempts.current = 0;
};
ws.onmessage = (event) => {
try {
const parsed = JSON.parse(event.data as string);
onMessage(parsed);
} catch {
console.error('[WS] Mensagem inválida:', event.data);
}
};
ws.onclose = (event) => {
if (event.wasClean) return; // Fechamento intencional
// Reconexão com backoff exponencial
const delay = Math.min(1000 * 2 ** reconnectAttempts.current, 30_000);
reconnectAttempts.current++;
console.log(`[WS] Reconectando em ${delay}ms...`);
reconnectTimeoutRef.current = setTimeout(connect, delay);
};
ws.onerror = (err) => {
console.error('[WS] Erro:', err);
};
}, [accessToken, onMessage]);
useEffect(() => {
connect();
return () => {
// Cleanup: fecha a conexão e cancela a reconexão pendente
if (reconnectTimeoutRef.current) {
clearTimeout(reconnectTimeoutRef.current);
}
wsRef.current?.close(1000, 'Componente desmontado');
};
}, [connect]);
}Escalonamento Horizontal com Redis Pub/Sub
O maior problema de WebSockets com múltiplas instâncias: um cliente conectado na Instância A não pode receber mensagens de eventos que aconteceram na Instância B. A solução é um message broker centralizado — o Redis Pub/Sub:
import Redis from 'ioredis';
import type { AppWebSocketServer } from './WebSocketServer';
// Dois clientes Redis separados: pub e sub não podem compartilhar conexão
const publisher = new Redis(process.env.REDIS_URL ?? 'redis://localhost:6379');
const subscriber = new Redis(process.env.REDIS_URL ?? 'redis://localhost:6379');
const WS_CHANNEL = 'ws:events';
interface WSEvent {
target: 'user' | 'broadcast';
userId?: string;
excludeUserId?: string;
data: unknown;
}
export function initRedisPubSub(wsServer: AppWebSocketServer): void {
// Escuta eventos publicados por QUALQUER instância
subscriber.subscribe(WS_CHANNEL, (err) => {
if (err) console.error('[Redis PubSub] Erro ao subscrever:', err);
});
subscriber.on('message', (_channel: string, message: string) => {
try {
const event = JSON.parse(message) as WSEvent;
if (event.target === 'user' && event.userId) {
// Entrega para o usuário conectado nesta instância (se houver)
wsServer.sendToUser(event.userId, event.data);
} else if (event.target === 'broadcast') {
wsServer.broadcast(event.data, event.excludeUserId);
}
} catch {
console.error('[Redis PubSub] Mensagem inválida:', message);
}
});
}
// Função exportada para os Use Cases publicarem eventos
export async function publishWSEvent(event: WSEvent): Promise<void> {
await publisher.publish(WS_CHANNEL, JSON.stringify(event));
}
// Nos Use Cases, substitua wsServer.sendToUser() por:
// await publishWSEvent({ target: 'user', userId: dto.customerId, data: { type: 'ORDER_CREATED', ... } });
// Agora qualquer instância entrega a mensagem para o usuário certo.Conclusão
WebSockets com a biblioteca ws cobrem a maioria dos casos de uso de tempo real sem o overhead do Socket.IO. O padrão crítico é: autenticar no handshake (não confiar em mensagens pós-conexão), usar heartbeat para detectar conexões mortas, organizar clientes em um registry por userId para entrega direcionada, e usar Redis Pub/Sub desde o início se a aplicação precisar escalar horizontalmente.