Appearance
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:
| Uso | API | Motivo |
|---|---|---|
| Comandos não bloqueantes | redis singleton | Usado para XADD, XACK, GET, SET, DEL e comandos curtos. |
Consumers com XREADGROUP BLOCK | createBlockingClient() | 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.
| Stream | Consumer group | Serviço | Entry point | Responsabilidade |
|---|---|---|---|---|
stream:incoming | orchestrator | flow-ai-orchestrator | consumeIncoming() | Persistir incoming, manter SessionState e rotear para stream:flow (mode flow), stream:helpdesk (mode human) ou stream:agent (mode agent). |
stream:flow | engine | flow-ai-engine | consumeFlow() | Executar o flow e publicar respostas no stream:outgoing-* escolhido pelo prefixo do sessionId. |
stream:outgoing-meta | persist | flow-ai-orchestrator | consumePersist() | Persistir a mensagem de saída como pending no Postgres. |
stream:outgoing-meta | send | flow-ai-meta-api | consumeOutgoing() | Enviar mensagens para a Meta Cloud API e publicar status. |
stream:outgoing-chat | persist | flow-ai-orchestrator | consumePersist() | Persistir mensagem de saída de WebChat (mesmo loop do meta). |
stream:outgoing-chat | webchat | flow-ai-webchat-gateway | startOutgoingConsumer() | Entregar mensagem à sala wc-{channelId}-{userId} via Socket.io. |
stream:outgoing-ig | persist | flow-ai-orchestrator | consumePersist() | Grupo registrado mas não consumido — ver nota de drift abaixo. |
stream:outgoing-ig | send | flow-ai-ig-api | consumeOutgoingIg() | Enviar mensagens para a IG Graph API e publicar status. |
stream:helpdesk | core | flow-ai-core | startHelpdeskConsumer() | Criar/atribuir tickets, tratar handoff e emitir eventos via Socket.io. |
stream:status | status | flow-ai-orchestrator | consumeStatus() | Persistir/atualizar o status das mensagens no Postgres. |
stream:status | core-status | flow-ai-core | startHelpdeskStatusConsumer() | Retransmitir atualizações de entrega ao desk em tempo real (não persiste). |
stream:status | core-campaigns-status | flow-ai-core | startCampaignStatusConsumer() | Atualizar CampaignRecipient.status (sent/delivered/read/failed). |
stream:wa-flow-logs | core-wa-flow-logs | flow-ai-core | startWhatsAppFlowLogConsumer() | Persistir em whatsapp_flow_webhook_logs. |
stream:wa-flow-logs | core-wa-flow-deferred | flow-ai-core | startWhatsAppFlowDeferredScheduler() | Agendar jobs deferidos (endpointConfig.deferredJobs). |
stream:analytics | core-analytics | flow-ai-core | startAnalyticsConsumer() | Persistir em event_occurrences. |
stream:templates | core-templates | flow-ai-core | startTemplateStatusConsumer() | Atualizar status/quality de templates e emitir via Socket.io. |
stream:tool-executions | core-tool-executions | flow-ai-core | startToolExecutionsConsumer() | Persistir em tool_execution_log. |
stream:agent | agent | flow-ai-agent | startConsumer() | Executar o loop LLM de tool-calling (consumer serial único). |
stream:llm-usage | core-llm-usage | flow-ai-core | startLlmUsageConsumer() | Persistir uso de tokens em llm_usage. |
stream:agent-turns | core-agent-turns | flow-ai-core | startAgentTurnsConsumer() | Persistir turns de agente em agent_turns (AI Debugger). |
stream:summary | core-summary | flow-ai-core | startSummaryConsumer() | Gerar resumos via LLM e persistir em Chat/Ticket/Contact. |
stream:ig-comments | ig-comments | flow-ai-ig-api | consumeIgComments() | Casar keyword com gatilhos e disparar DM/resposta pública. |
stream:ig-comment-activations | core-ig-comment-activations | flow-ai-core | startIgCommentActivationsConsumer() | Persistir em instagram_comment_activations. |
Nota (legado): a constante
STREAMS.OUTGOING(stream:outgoing) ainda existe emconstants.ts, mas não aparece noSTREAM_GROUP_MAP— não tem nenhum consumer group nem produtor ativo. As três streams de saída vivas sãostream:outgoing-meta,stream:outgoing-chatestream:outgoing-ig, uma por canal. O engine escolhe qual usar viaresolveContentChannel(sessionId)(services/flow-ai-engine/src/runtime/publish.ts):ig-*→ ig,wc-*→ chat, caso contrário meta.
Nota (drift): o grupo
persistestá registrado parastream:outgoing-ig(streams.ts:152), mas oconsumePersist()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 emMessageno 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:
- Cria o
blockingClientcomcreateBlockingClient()(auto-registrado para teardown). - Registra o loop no drain via
registerLoop(). - Enquanto
!isShuttingDown(), lê entries comxreadgroup(group, consumer, streams, count, blockMs, blockingClient)(padrãocount=10,blockMs=2000). - Para cada batch, chama
onBatch(results)— o consumer parseiafields.payload, valida o shape esperado, executa o handler de negócio e fazxack()conforme a política. - Erro de leitura no Redis → loga e retenta em 1s (sem ACK).
- 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). - No shutdown, o
BLOCKcorrente retorna em atéblockMs, a flag é vista, o loop sai limpo e chamaloop.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()doflow-ai-meta-api(services/flow-ai-meta-api/src/consumers/outgoing.ts) não usarunStreamLoop(). Como ele precisa atender simultaneamente o Redis principal e um Redis dev (roteamento do número de testeTEST_PHONE_NUMBER_ID), ele dirigeregisterLoop()/isShuttingDown()diretamente, através de uma abstraçãoOutgoingTransport. O transporte dev (createDevTransport) chamablockingClient.xreadgroup("GROUP", ...)cru do ioredis, porque os helpers deflow-ai-redisoperam sobre o singleton e não têm overload dexadd/xackpara um cliente alternativo.
Política de ACK e pendência
| Caso | Ação | Racional |
|---|---|---|
Entry sem payload | XACK | Payload não recuperável. |
| JSON inválido | XACK | Payload não recuperável. |
| Shape inválido | XACK | Payload não corresponde ao contrato da stream. |
| Erro de leitura no Redis | Sem ACK, sleep 1s | Falha fora de uma entry específica. |
| Erro no handler de negócio | Sem ACK | Permite reentrega via pending list (PEL). |
Cache miss de credenciais no meta-api/ig-api | Sem ACK | Entry deve ser reprocessada quando o flow-ai-core popular wa:credentials:{phoneNumberId} / ig:credentials:{igUserId}. |
| Falha de envio para a Meta/IG | Publica send_failed e XACK | A 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 peloflow-ai-orchestrator(consumePersist()), grava a mensagem no Postgres comopending.send, consumido peloflow-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.
flow-ai-corefaz warm-up no boot commakeWhatsAppNumberService().warmupCache().- CRUD de WhatsApp Numbers cifra o token no Postgres e grava a versão decifrada no Redis.
flow-ai-meta-apiconsomestream:outgoing-meta(groupsend).- Antes de enviar, busca
wa:credentials:{phoneNumberId}viagetWaCredentials(). - 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:
- Levanta a flag
isShuttingDown()— os loops param de pegar batch novo. - Pré-drain (
onPreDrain): para de aceitar trabalho novo — Fastifyapp.close(), Socket.ioio.close(). Não afeta o stream em voo. - Espera os loops de consumo saírem (cada um termina o batch atual +
xacke encerra), limitado pordrainTimeoutMs(padrão 15s). - Pós-drain (
onPostDrain): libera recursos usados no processamento, ex.prisma.$disconnect(). - Fecha as conexões bloqueantes (
quitBlockingClients()) — loops já saíram, sem reads em voo. - Fecha a conexão de pub/sub compartilhada (
closeSubscriberConnection()). - Fecha o Redis singleton compartilhado por último (
quitRedis()) —xack/xaddusam 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. oconsumeOutgoing()do meta-api) registra-se e recebedone()para chamar ao sair. Loops viarunStreamLoop()já fazem isso automaticamente.runPollingLoop({ intervalMs, tick })— para workers de polling (não-stream), ex.startIdleTicketAbandonWorker. UsainterruptibleDelay(), 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)