Skip to main content

Arbitragem em Tempo Real via WebSocket

O fluxo WebSocket é completamente independente do fluxo REST existente. Ambos rodam em paralelo — o WebSocket adiciona latência ultrabaixa sem afetar a estabilidade do processamento REST de 19 segundos.

🎯 Visão Geral

O sistema possui dois fluxos de integração complementares:

🏗️ Arquitetura dos Componentes

Componentes Principais

WebSocketExchangeHandler

Interface contratual para cada exchange. Implementa Connect(), Name() e HealthCheck(). Cada exchange tem sua própria implementação.

RealtimePriceCache

Cache thread-safe em memória. Preserva volume REST quando o tick WebSocket não o fornece. Fornece GetAll() no mesmo formato do fluxo REST.

RealtimeArbitrageEngine

Motor de detecção. Debounce de 200ms por símbolo evita processamento excessivo. Reutiliza CompareBaseQuotePrices() do fluxo REST sem duplicação de código.

SSEBroadcaster

Fan-out de Server-Sent Events para todos os browsers conectados. Cada cliente recebe um canal dedicado com backpressure automático.

WebSocketManager

Singleton orquestrador. Gerencia reconexão automática com backoff exponencial (1s → 2s → 4s → 8s → 16s máx).

🔄 Fluxo de Dados Detalhado

1. Conexão e Seed de Volume

Ao conectar, o handler WebSocket da Binance:
  1. Chama a REST API pública para buscar snapshot de volume (24h) de todos os pares
  2. Alimenta o RealtimePriceCache com os dados iniciais
  3. Abre conexão WebSocket: wss://stream.binance.com:9443/ws
  4. Subscreve o stream !bookTicker (todos os pares, bid/ask em tempo real)
  5. Inicia timer de 60s para refresh periódico do volume via REST

2. Processamento de Ticks

3. Reconexão Automática

⚡ Ativação

O fluxo WebSocket é controlado por variável de ambiente para garantir zero impacto em produção enquanto não estiver pronto:
Se ENABLED_WEBSOCKET não estiver definido ou for qualquer outro valor, o fluxo WebSocket não é iniciado e nenhum recurso adicional é consumido.

📡 Endpoints SSE

Cinco endpoints HTTP expõem os dados do fluxo WebSocket:

Streams por Caso de Uso

/v1/realtime/prices

Dashboard global — todos os pares de todas as exchanges. Volume alto (~900 KB/s em 6 exchanges). Use com moderação.

/v1/realtime/arbitrage

Oportunidades — somente pares com spread positivo detectado. Payload compacto, ideal para alertas.

/v1/lookers/{id}/stream

Monitoramento pessoal — preços ask/bid apenas das cryptos configuradas no Looker do usuário (máx. 6 pares). Com campo stale para detectar dados desatualizados.

/v1/calculators/stream

PNL da posição — PNL calculado a cada 1s para uma posição spot-futuro aberta. Usa dados de mercado do MongoDB com cache de 5s.

Consumo no Frontend (EventSource)

A API EventSource do browser reconecta automaticamente ao servidor em caso de desconexão, sem necessidade de lógica adicional no frontend.

🗄️ Persistência

Oportunidades detectadas pelo motor WebSocket são salvas na coleção operations_realtime do MongoDB:
A coleção operations_realtime é projetada para ter TTL curto — os dados são efêmeros por natureza, servindo principalmente como auditoria e buffer para clientes que se reconectam.

🔌 Exchanges Suportadas

Fase 1 e 2 — Concluídas ✅

Fase 3 — Planejada 🔄

Fase 4 — REST-only (sem WebSocket público)

🧪 Interface do Handler

Para adicionar suporte a uma nova exchange, implemente a interface WebSocketExchangeHandler:
Registre o novo handler no WebSocketManager:

📊 Monitoramento

Use o endpoint /v1/realtime/status para verificar a saúde de cada handler:
  • Fase 1 + 2 concluídas: 6 exchanges com WebSocket ativo (Binance, OKX, Bybit, KuCoin, Gate.io, HTX)
  • Looker Stream: GET /v1/lookers/{id}/stream — preços filtrados por Looker
  • Calculator Stream: GET /v1/calculators/stream — PNL em tempo real da posição
  • 🔄 Fase 3: Bitget, Kraken, Bitfinex, BingX, MEXC, Crypto.com
  • 🔄 Fase 4: Avaliar alternativas para exchanges REST-only