Skip to content

Redis consumers — conexões e consumer groups

Este documento descreve como os consumers da aplicação se conectam ao Redis, como os consumer groups são criados e qual política de ACK/retry cada serviço segue.

Regra de conexão

O pacote flow-ai-redis expõe dois caminhos de conexão:

UsoAPIMotivo
Comandos não bloqueantesredis singletonUsado para XADD, XACK, GET, SET, DEL e comandos curtos.
Consumers com XREADGROUP BLOCKcreateBlockingClient()Cria redis.duplicate() para o loop bloqueante de cada consumer.

Cada consumer bloqueante precisa de uma conexão dedicada. O comando XREADGROUP ... BLOCK ... segura a conexão TCP enquanto espera mensagens; se vários consumers compartilhassem o singleton, comandos como XADD, XACK e SET ficariam atrasados atrás do bloqueio.

createBlockingClient() (packages/flow-ai-redis/src/client.ts) registra cada conexão dedicada num Set interno (blockingClients), para que o graceful shutdown as feche no momento certo (ver seção Graceful shutdown).

Mapa de streams e consumer groups

As streams e os grupos ficam centralizados em packages/flow-ai-redis/src/constants.ts (STREAMS e GROUPS). A infraestrutura (todas as combinações stream → grupo) é registrada no STREAM_GROUP_MAP de packages/flow-ai-redis/src/streams.ts e criada por setupStreams() no boot de cada serviço.

Existem 18 streams declaradas em STREAMS; uma delas (stream:outgoing) é legado sem grupos ativos (ver nota abaixo). O STREAM_GROUP_MAP tem 23 registros stream → grupo. Note que stream:status tem três grupos e stream:outgoing-meta, stream:outgoing-ig, stream:outgoing-chat e stream:wa-flow-logs têm dois cada.

StreamConsumer groupServiçoEntry pointResponsabilidade
stream:incomingorchestratorflow-ai-orchestratorconsumeIncoming()Persistir incoming, manter SessionState e rotear para stream:flow (mode flow), stream:helpdesk (mode human) ou stream:agent (mode agent).
stream:flowengineflow-ai-engineconsumeFlow()Executar o flow e publicar respostas no stream:outgoing-* escolhido pelo prefixo do sessionId.
stream:outgoing-metapersistflow-ai-orchestratorconsumePersist()Persistir a mensagem de saída como pending no Postgres.
stream:outgoing-metasendflow-ai-meta-apiconsumeOutgoing()Enviar mensagens para a Meta Cloud API e publicar status.
stream:outgoing-chatpersistflow-ai-orchestratorconsumePersist()Persistir mensagem de saída de WebChat (mesmo loop do meta).
stream:outgoing-chatwebchatflow-ai-webchat-gatewaystartOutgoingConsumer()Entregar mensagem à sala wc-{channelId}-{userId} via Socket.io.
stream:outgoing-igpersistflow-ai-orchestratorconsumePersist()Grupo registrado mas não consumido — ver nota de drift abaixo.
stream:outgoing-igsendflow-ai-ig-apiconsumeOutgoingIg()Enviar mensagens para a IG Graph API e publicar status.
stream:helpdeskcoreflow-ai-corestartHelpdeskConsumer()Criar/atribuir tickets, tratar handoff e emitir eventos via Socket.io.
stream:statusstatusflow-ai-orchestratorconsumeStatus()Persistir/atualizar o status das mensagens no Postgres.
stream:statuscore-statusflow-ai-corestartHelpdeskStatusConsumer()Retransmitir atualizações de entrega ao desk em tempo real (não persiste).
stream:statuscore-campaigns-statusflow-ai-corestartCampaignStatusConsumer()Atualizar CampaignRecipient.status (sent/delivered/read/failed).
stream:wa-flow-logscore-wa-flow-logsflow-ai-corestartWhatsAppFlowLogConsumer()Persistir em whatsapp_flow_webhook_logs.
stream:wa-flow-logscore-wa-flow-deferredflow-ai-corestartWhatsAppFlowDeferredScheduler()Agendar jobs deferidos (endpointConfig.deferredJobs).
stream:analyticscore-analyticsflow-ai-corestartAnalyticsConsumer()Persistir em event_occurrences.
stream:templatescore-templatesflow-ai-corestartTemplateStatusConsumer()Atualizar status/quality de templates e emitir via Socket.io.
stream:tool-executionscore-tool-executionsflow-ai-corestartToolExecutionsConsumer()Persistir em tool_execution_log.
stream:agentagentflow-ai-agentstartConsumer()Executar o loop LLM de tool-calling (consumer serial único).
stream:llm-usagecore-llm-usageflow-ai-corestartLlmUsageConsumer()Persistir uso de tokens em llm_usage.
stream:agent-turnscore-agent-turnsflow-ai-corestartAgentTurnsConsumer()Persistir turns de agente em agent_turns (AI Debugger).
stream:summarycore-summaryflow-ai-corestartSummaryConsumer()Gerar resumos via LLM e persistir em Chat/Ticket/Contact.
stream:ig-commentsig-commentsflow-ai-ig-apiconsumeIgComments()Casar keyword com gatilhos e disparar DM/resposta pública.
stream:ig-comment-activationscore-ig-comment-activationsflow-ai-corestartIgCommentActivationsConsumer()Persistir em instagram_comment_activations.

Nota (legado): a constante STREAMS.OUTGOING (stream:outgoing) ainda existe em constants.ts, mas não aparece no STREAM_GROUP_MAP — não tem nenhum consumer group nem produtor ativo. As três streams de saída vivas são stream:outgoing-meta, stream:outgoing-chat e stream:outgoing-ig, uma por canal. O engine escolhe qual usar via resolveContentChannel(sessionId) (services/flow-ai-engine/src/runtime/publish.ts): ig-* → ig, wc-* → chat, caso contrário meta.

Nota (drift): o grupo persist está registrado para stream:outgoing-ig (streams.ts:152), mas o consumePersist() só lê [OUTGOING_META, OUTGOING_CHAT] (services/flow-ai-orchestrator/src/consumers/persist.ts:186). Ou seja, o orchestrator não persiste as mensagens de saída do Instagram — elas são enviadas pela IG Graph API, mas não gravadas em Message no Postgres.

Padrão operacional de um consumer

A maioria dos loops de consumo usa o helper runStreamLoop() de packages/flow-ai-redis/src/shutdown.ts, que centraliza o ciclo de vida da conexão bloqueante e a integração com o graceful shutdown:

  1. Cria o blockingClient com createBlockingClient() (auto-registrado para teardown).
  2. Registra o loop no drain via registerLoop().
  3. Enquanto !isShuttingDown(), lê entries com xreadgroup(group, consumer, streams, count, blockMs, blockingClient) (padrão count=10, blockMs=2000).
  4. Para cada batch, chama onBatch(results) — o consumer parseia fields.payload, valida o shape esperado, executa o handler de negócio e faz xack() conforme a política.
  5. Erro de leitura no Redis → loga e retenta em 1s (sem ACK).
  6. Erro dentro do onBatch → loga e segue em 1s (não derruba o loop; entries não-ackadas caem no PEL e são recuperadas depois).
  7. No shutdown, o BLOCK corrente retorna em até blockMs, a flag é vista, o loop sai limpo e chama loop.done().

Cada mensagem é processada dentro de runWithMessageContext(stream, group, id, fields, fn) (flow-ai-telemetry), que reconstrói o contexto de trace propagado no XADD.

Exceção: o consumeOutgoing() do flow-ai-meta-api (services/flow-ai-meta-api/src/consumers/outgoing.ts) não usa runStreamLoop(). Como ele precisa atender simultaneamente o Redis principal e um Redis dev (roteamento do número de teste TEST_PHONE_NUMBER_ID), ele dirige registerLoop() / isShuttingDown() diretamente, através de uma abstração OutgoingTransport. O transporte dev (createDevTransport) chama blockingClient.xreadgroup("GROUP", ...) cru do ioredis, porque os helpers de flow-ai-redis operam sobre o singleton e não têm overload de xadd/xack para um cliente alternativo.

Política de ACK e pendência

CasoAçãoRacional
Entry sem payloadXACKPayload não recuperável.
JSON inválidoXACKPayload não recuperável.
Shape inválidoXACKPayload não corresponde ao contrato da stream.
Erro de leitura no RedisSem ACK, sleep 1sFalha fora de uma entry específica.
Erro no handler de negócioSem ACKPermite reentrega via pending list (PEL).
Cache miss de credenciais no meta-api/ig-apiSem ACKEntry deve ser reprocessada quando o flow-ai-core popular wa:credentials:{phoneNumberId} / ig:credentials:{igUserId}.
Falha de envio para a Meta/IGPublica send_failed e XACKA falha foi convertida em evento durável em stream:status.

Fan-out em stream:outgoing-meta

stream:outgoing-meta tem dois consumer groups independentes:

  • persist, consumido pelo flow-ai-orchestrator (consumePersist()), grava a mensagem no Postgres como pending.
  • send, consumido pelo flow-ai-meta-api (consumeOutgoing()), envia a mensagem para a Meta.

Esses grupos não competem entre si. A mesma entry é entregue uma vez para cada grupo, permitindo durabilidade e envio em paralelo. O mesmo padrão vale para stream:outgoing-chat (persist + webchat) e stream:outgoing-ig (persist — dormente — + send).

Fluxo de credenciais WhatsApp

O consumer flow-ai-meta-api/src/consumers/outgoing.ts não acessa o token no Postgres. Ele lê o cache Redis wa:credentials:{phoneNumberId}, preenchido pelo flow-ai-core.

  1. flow-ai-core faz warm-up no boot com makeWhatsAppNumberService().warmupCache().
  2. CRUD de WhatsApp Numbers cifra o token no Postgres e grava a versão decifrada no Redis.
  3. flow-ai-meta-api consome stream:outgoing-meta (group send).
  4. Antes de enviar, busca wa:credentials:{phoneNumberId} via getWaCredentials().
  5. Se houver miss, não faz ACK; a entry fica pendente até o cache ser populado.

Graceful shutdown

O shutdown determinístico vive em packages/flow-ai-redis/src/shutdown.ts. Cada serviço chama installShutdownHandlers({ service, log, drainTimeoutMs }) uma vez no bootstrap, depois de registrar os cleanups e iniciar os loops. O handler de SIGINT/SIGTERM orquestra a ordem de teardown:

  1. Levanta a flag isShuttingDown() — os loops param de pegar batch novo.
  2. Pré-drain (onPreDrain): para de aceitar trabalho novo — Fastify app.close(), Socket.io io.close(). Não afeta o stream em voo.
  3. Espera os loops de consumo saírem (cada um termina o batch atual + xack e encerra), limitado por drainTimeoutMs (padrão 15s).
  4. Pós-drain (onPostDrain): libera recursos usados no processamento, ex. prisma.$disconnect().
  5. Fecha as conexões bloqueantes (quitBlockingClients()) — loops já saíram, sem reads em voo.
  6. Fecha a conexão de pub/sub compartilhada (closeSubscriberConnection()).
  7. Fecha o Redis singleton compartilhado por último (quitRedis()) — xack/xadd usam ele.

Um watchdog força process.exit(1) se o drain estourar drainTimeoutMs + 5s. O kill_timeout do pm2 precisa ser maior que isso, senão o pm2 manda SIGKILL antes do drain terminar.

Helpers relacionados:

  • registerLoop() — um loop de formato próprio (ex. o consumeOutgoing() do meta-api) registra-se e recebe done() para chamar ao sair. Loops via runStreamLoop() já fazem isso automaticamente.
  • runPollingLoop({ intervalMs, tick }) — para workers de polling (não-stream), ex. startIdleTicketAbandonWorker. Usa interruptibleDelay(), que acorda imediatamente no shutdown para não segurar o drain com um intervalo longo.

Arquivos de referência

  • packages/flow-ai-redis/src/constants.ts (STREAMS, GROUPS)
  • packages/flow-ai-redis/src/streams.ts (STREAM_GROUP_MAP, setupStreams, xreadgroup, xack, xadd)
  • packages/flow-ai-redis/src/client.ts (createBlockingClient, quitBlockingClients, quitRedis)
  • packages/flow-ai-redis/src/shutdown.ts (installShutdownHandlers, runStreamLoop, runPollingLoop, registerLoop, onPreDrain, onPostDrain)
  • services/flow-ai-orchestrator/src/consumer.ts (consumeIncoming)
  • services/flow-ai-orchestrator/src/consumers/persist.ts (consumePersist)
  • services/flow-ai-orchestrator/src/consumers/status.ts (consumeStatus)
  • services/flow-ai-engine/src/consumer.ts (consumeFlow)
  • services/flow-ai-meta-api/src/consumers/outgoing.ts (consumeOutgoing)
  • services/flow-ai-ig-api/src/consumers/outgoing.ts (consumeOutgoingIg)
  • services/flow-ai-agent/src/consumer.ts (startConsumer)
  • services/flow-ai-webchat-gateway/src/consumer.ts (startOutgoingConsumer)
  • services/flow-ai-core/src/helpdesk/consumer.ts (startHelpdeskConsumer)
  • services/flow-ai-core/src/consumers/*.consumer.ts (demais consumers do core)

flow-ai — plataforma proprietária de atendimento via WhatsApp