Appearance
Fluxo outbound completo
Do engine publicar uma resposta até a mensagem aparecer no WhatsApp do usuário — e como o status de entrega retorna ao sistema.
Visão geral
O caminho outbound é mais complexo que o inbound porque envolve fan-out: o mesmo evento é consumido em paralelo por dois consumers independentes, e o resultado de cada um alimenta streams distintas.
flow-ai-engine
│ publish OutgoingStreamEvent
├─── stream:outgoing-meta ──────────────────────────────────┐
│ │ │
│ group: persist group: send
│ flow-ai-orchestrator flow-ai-meta-api
│ → upsert Message (pending) → chama Meta Cloud API
│ → publica StatusStreamEvent
│ │
│ stream:status
│ │
│ ┌──────────────┼──────────────────────┐
│ group: status group: core-status group: core-campaigns-status
│ orchestrator flow-ai-core flow-ai-core
│ → update Message → Socket.io → desk → update CampaignRecipient
│ (sent/failed/read)
│
├─── stream:outgoing-chat ─── (WebChat: group persist + group webchat no gateway)
│
└─── stream:outgoing-ig ───── (Instagram: só group send no flow-ai-ig-api)Nota: existem três streams outbound vivas —
stream:outgoing-meta,stream:outgoing-chatestream:outgoing-ig. A stream antigastream:outgoing(singular) é legado: a constante ainda existe emflow-ai-redis, mas não tem produtores nem consumer groups ativos.
Nunca há chamada HTTP direta do engine para a Meta. O engine só publica numa stream e esquece — a entrega é de responsabilidade de quem consome.
Fase 1 — engine: publicação do evento outgoing
1.1 Geração do internalId
Antes de publicar qualquer mensagem, o engine gera um internalId para cada OutgoingMessage usando cuid():
typescript
type OutgoingMessage = (WhatsAppOutgoingMessage | InstagramOutgoingMessage) & {
internalId: string // CUID gerado localmente pelo publicador
contextMessageId?: string // wamid da mensagem citada (balão de citação no WhatsApp)
}Esse ID é a chave de correlação que liga os três momentos de vida da mensagem: persistência inicial (pending), confirmação de envio (sent/failed) e status final de entrega (delivered/read). Ele existe porque a mensagem é processada por dois consumers independentes em ordem não determinística — sem o internalId, não haveria como saber qual registro de banco atualizar.
Por que não usar o wamid da Meta como chave? O wamid só existe após o envio bem-sucedido. Se o sistema persistir a mensagem depois de receber o wamid, perde-se a garantia de durabilidade: uma falha entre o envio e a persistência deixaria a mensagem enviada sem registro. A ordem certa é persistir primeiro (com internalId), enviar depois, e correlacionar pelo internalId no status.
1.2 Roteamento por canal
O engine publica em streams diferentes conforme o canal da sessão:
Prefixo do sessionId | Stream destino |
|---|---|
wa-{phoneNumberId}-{from} | stream:outgoing-meta |
wc-{channelId}-{userId} | stream:outgoing-chat |
ig-{igUserId} | stream:outgoing-ig |
Os três produtores da stream outbound (engine, flow-ai-agent e flow-ai-core) escolhem a stream inspecionando o prefixo do sessionId: ig-* → stream:outgoing-ig, wc-* → stream:outgoing-chat, o restante → stream:outgoing-meta. Cada um faz esse startsWith inline no ponto de publicação. (No engine existe também um helper resolveContentChannel(sessionId) que devolve ig/wc/wa, mas ele serve à conversão de mídia — não à escolha da stream.)
Nota: a stream
stream:outgoing-igsó tem o consumer groupsendativo (noflow-ai-ig-api). O grouppersistestá registrado emSTREAM_GROUP_MAP, mas o loop de persist do orchestrator lê apenasoutgoing-metaeoutgoing-chat— ou seja, DMs de Instagram não ganham a linhapendingdo persist consumer (perde-se a garantia de durabilidade "persistir antes de enviar"). A Message ainda acaba no Postgres: oflow-ai-ig-apipublica umStatusStreamEvent(sent/send_failed) após o envio, e o consumer de status do orchestrator fazupsertporinternalId— o ramocreatecria a Message nesse momento (ver Fase 3.1). A exceção são as respostas a comentário (igCommentTarget): oflow-ai-ig-apinão emiteStatusStreamEventpara elas (são ações de moderação one-shot sem Chat), então essas nunca são persistidas.
typescript
type OutgoingStreamEvent = {
sessionId: string
chatId: string
to: string // número/identificador do destinatário
phoneNumberId: string // remetente (conta Meta / igUserId / wc-{channelId})
source: "flow" | "agent" | "dispatch" // engine, agente (humano ou IA) ou campanha
agentUserId?: string // populado quando source="agent" e o remetente é humano
aiAgentId?: string // populado quando source="agent" e o remetente é IA
typing?: boolean // envia typing indicator antes das mensagens
lastIncomingMessageId?: string // wamid da última mensagem do contato (typing)
ticketId?: string // ticket ativo (AI ou humano) para taguear a Message
igCommentTarget?: { commentId: string; mode: "private" | "public" } // resposta a comentário IG
messages: OutgoingMessage[]
timestamp: number
}O campo source distingue mensagens automáticas do engine ("flow"), mensagens de agente ("agent" — humano com agentUserId ou IA com aiAgentId) e disparos de campanha ("dispatch"). Ele é fundamental para o desk saber quais bolhas exibir como "enviado por agente" e para a lógica de status do consumer do core. Note que a stream outbound tem vários produtores: além do engine, o flow-ai-agent e o flow-ai-core também publicam nessas streams.
1.3 Conversão de content para OutgoingMessage
O runtime/publish.ts converte ContentMessage[] (formato interno do flow) para o formato OutgoingMessage. Para mídias (image, video, audio, document) no canal WhatsApp (channel === "wa"), isso inclui uma chamada à Meta Graph API para obter o mediaId a partir da URL — a Meta exige que mídias sejam referenciadas por ID, não por URL direta. Para Instagram e WebChat, o campo mediaId carrega a própria URL pública (esses canais aceitam a URL direta e não passam pela API de mídia da Meta).
Fase 2 — fan-out: persist e send em paralelo
O mesmo OutgoingStreamEvent é lido independentemente por dois consumer groups no mesmo stream. Redis Streams garante que cada group recebe cada entry exatamente uma vez — os dois consumers progridem na stream sem interferência.
2.1 Consumer persist (flow-ai-orchestrator)
O persist consumer grava a mensagem no Postgres com status pending antes que ela seja enviada:
typescript
// sender derivado de event.source + agentUserId — NÃO é event.source direto:
// agentUserId presente → "human"
// source === "agent" → "agent"
// caso contrário → "flow" (inclui source === "dispatch")
const sender = event.agentUserId ? "human" : event.source === "agent" ? "agent" : "flow"
// upsert pelo internalId
await prisma.message.upsert({
where: { id: msg.internalId },
create: {
id: msg.internalId,
chatId: event.chatId,
ticketId: event.ticketId ?? null,
sender,
type: msg.type,
content: extractContent(msg),
metadata: extractMetadata(msg), // botões, seções, mediaId, filename
status: "pending",
sentByUserId: event.agentUserId ?? null,
sentByAiAgentId: event.aiAgentId ?? null,
},
update: {
// idempotente: não sobrescreve status nem wamid se já foram atualizados
sender,
type: msg.type,
content: extractContent(msg),
metadata: extractMetadata(msg),
},
})Por que persistir com status pending antes de enviar? Se o send consumer falhar ou o processo reiniciar após o envio mas antes da persistência, a mensagem teria sido enviada ao usuário sem registro no banco. Persistir primeiro garante que qualquer mensagem que chegou ao usuário tem um registro — mesmo que o status ainda não seja sent. O pending é o estado de durabilidade: "sabemos que essa mensagem existe, mas ainda não confirmamos o envio".
O update parcial do upsert é intencional: se por alguma condição de corrida o consumer de status processar o StatusStreamEvent antes do persist processar o OutgoingStreamEvent (teoricamente impossível mas defensivamente codado), o update não sobrescreve o wamid e o status já gravados.
2.2 Consumer send (flow-ai-meta-api)
O send consumer envia a mensagem à Meta e publica o resultado em stream:status:
Passo 1 — Resolver credenciais:
typescript
const credentials = await getWaCredentials(event.phoneNumberId)
// → { accessToken, phoneNumberId, displayName, wabaId? }E se o cache estiver vazio (cache miss)?
O consumer não faz ACK. A entry fica na pending list do consumer group e será reentregue automaticamente. Isso ocorre na situação em que o flow-ai-core ainda não populou wa:credentials:{phoneNumberId} — por exemplo, após um restart com warm-up incompleto. A entry continua reentregue até o cache ser populado, sem perda de mensagem.
Passo 2 — Enviar para a Meta:
typescript
const client = new WhatsAppClient()
const result = await client.send(credentials, message)Passo 3 — Publicar resultado em stream:status:
typescript
// Sucesso
xadd(STREAMS.STATUS, {
payload: JSON.stringify({
internalId: msg.internalId,
wamid: result.messages[0].id,
status: "sent",
chatId: event.chatId,
source: event.source,
type: msg.type,
content: extractContent(msg),
recipientId: event.to,
timestamp: Date.now(),
})
})
// Falha (4xx/5xx ou rede)
xadd(STREAMS.STATUS, {
payload: JSON.stringify({
internalId: msg.internalId,
status: "send_failed",
error: { code, title, message },
...
})
})Tanto no sucesso quanto na falha, o send consumer sempre faz ACK. A falha de envio não é recuperável com retry imediato (a Meta já recusou ou o erro de rede vai persistir por segundos) — a falha é convertida em evento durável em stream:status e publicada como send_failed. O retry, se necessário, é decisão de negócio (ex: reenvio manual pelo agente).
Fase 3 — stream:status: atualização de banco e desk em tempo real
stream:status recebe dois tipos de evento:
| Origem | Evento | Consumer groups |
|---|---|---|
flow-ai-meta-api (após envio) | sent ou send_failed | status (orchestrator) + core-status (core) + core-campaigns-status (core) |
flow-ai-meta-api (webhook da Meta) | delivered, read, played, failed | status (orchestrator) + core-status (core) + core-campaigns-status (core) |
São três consumer groups em stream:status, cada um lê cada entry independentemente:
status(orchestrator) — persiste/atualiza oMessage.statusno Postgres;core-status(core) — retransmite o status ao desk via Socket.io;core-campaigns-status(core) — atualizaCampaignRecipient.status(sent/delivered/read/failed).
Nota:
stream:statuse o webhook de status são um contrato WhatsApp. Instagram e WebChat também publicam eventos aqui via seus gateways, mas o ciclodelivered/read/playeddocumentado abaixo reflete a semântica da Meta.
3.1 Consumer status (flow-ai-orchestrator) — atualiza o banco
Para sent:
typescript
await prisma.message.upsert({
where: { id: event.internalId },
create: { ..., whatsappMessageId: event.wamid, status: "sent" },
update: { whatsappMessageId: event.wamid, status: "sent" },
})Para delivered / read / played / failed (webhook Meta):
typescript
await prisma.message.update({
where: { whatsappMessageId: event.wamid },
data: { status: mappedStatus },
})Se o wamid for desconhecido (pode ocorrer em cenários de migração ou mensagens antigas), o consumer faz ACK com warning — não há dado suficiente para atualizar.
3.2 Consumer core-status (flow-ai-core) — atualiza o desk via Socket.io
Este consumer é responsável por entregar atualizações de status em tempo real ao desk do agente. A retransmissão do status ao desk só ocorre para mensagens de agente (source: "agent", ou Message.sender === "agent" no caso dos eventos de webhook resolvidos por wamid):
- Busca o ticket aberto (
waiting/assignedcomassignedUserId) associado aochatId. - Emite
ticket:message-statusvia Socket.io para o agente assignado.
O mapeamento de status para o desk (mapStatusForDesk) é:
| Status na stream | Status no desk |
|---|---|
sent ou delivered | "delivered" |
read | "read" |
played | "played" |
send_failed ou failed | "failed" |
Caso especial — janela de 24h expirada: Se o status for send_failed com código 131026 (erro Meta para janela de conversa expirada), o core fecha o ticket automaticamente via handleWindowExpired. Esse ramo é avaliado antes do filtro source === "agent", ou seja, dispara para qualquer source. O handler: fecha o ticket aberto mais recente do chat com closeKind: "meta_window_expired" (sem userId — ação do sistema); publica humanSessionEnded em stream:flow para o engine retomar o flow, se possível; e emite ticket:window_expired via Socket.io ao agente atribuído. O usuário ficou 24h sem interagir, a janela Meta expirou, e mensagens proativas não são mais possíveis sem template.
Ciclo de vida de uma mensagem no banco
engine cria internalId
│
▼
persist consumer → Message { status: "pending" }
│
▼
send consumer → Meta Cloud API
│
├── sucesso → StatusStreamEvent { status: "sent", wamid }
│ │
│ status consumer → Message { status: "sent", whatsappMessageId: wamid }
│
└── falha → StatusStreamEvent { status: "send_failed" }
│
status consumer → Message { status: "failed" }
(depois, via webhook da Meta)
wamid recebido → StatusStreamEvent { status: "delivered" | "read" | "played" }
│
status consumer → Message { status: "delivered" | "read" | "played" }O internalId é a ponte entre o engine e o banco nos estágios pending → sent/failed. O wamid assume o papel de chave nos estágios posteriores (delivered, read) porque esses eventos chegam da Meta com apenas o wamid — sem o internalId interno.
Por que dois consumers no mesmo stream
O fan-out entre persist e send implementa separação de responsabilidades sem coordenação:
- O persist consumer não sabe se o send teve sucesso — ele só garante durabilidade.
- O send consumer não sabe se a persistência ocorreu — ele só garante entrega à Meta.
- Os dois falham e reiniciam independentemente.
Se fosse uma chamada sequencial (persist → send), uma falha no send poderia deixar o consume travado, impedindo novas mensagens de fluir. Com fan-out, cada consumer tem sua own pending list — um travamento em um não afeta o outro.
Arquivos relevantes
| Arquivo | Papel |
|---|---|
services/flow-ai-engine/src/runtime/publish.ts | Converte ContentMessage → OutgoingMessage, gera internalId, publica |
services/flow-ai-engine/src/runtime/publish-helpdesk.ts | Publica handoff para stream:helpdesk |
services/flow-ai-orchestrator/src/consumers/persist.ts | Upsert Message com status pending |
services/flow-ai-orchestrator/src/consumers/status.ts | Atualiza status de Message no Postgres |
services/flow-ai-meta-api/src/consumers/outgoing.ts | Consome stream:outgoing-meta (group send), envia à Meta, publica StatusStreamEvent |
services/flow-ai-ig-api/src/consumers/outgoing.ts | Consome stream:outgoing-ig (group send), envia à IG Graph API (não persiste) |
services/flow-ai-core/src/helpdesk/status-consumer.ts | Consumer core-status: emite status ao desk via Socket.io |
services/flow-ai-core/src/helpdesk/handlers/window-expired.ts | Fecha ticket em send_failed 131026 (janela de 24h) |
packages/flow-ai-types/src/stream-events.ts | OutgoingStreamEvent, StatusStreamEvent, OutgoingMessage |
packages/flow-ai-redis/src/constants.ts | STREAMS.OUTGOING_{META,CHAT,IG}, STREAMS.STATUS, GROUPS |
packages/flow-ai-redis/src/streams.ts | STREAM_GROUP_MAP — mapeamento stream → consumer group |