Appearance
Execução do flow-ai-engine
Este documento descreve o fluxo técnico do flow-ai-engine, desde o bootstrap do processo até o runFlow() executar blocos, publicar mensagens e persistir a sessão no Redis.
Arquivos de referência
services/flow-ai-engine/src/index.tsservices/flow-ai-engine/src/bootstrap.tsservices/flow-ai-engine/src/consumer.tsservices/flow-ai-engine/src/handlers/flow.tsservices/flow-ai-engine/src/flows-cache.tsservices/flow-ai-engine/src/runtime/execute.tsservices/flow-ai-engine/src/runtime/actions.tsservices/flow-ai-engine/src/runtime/input.tsservices/flow-ai-engine/src/runtime/publish.tspackages/flow-ai-runtime/src/conditions.ts(avaliação de condições, compartilhada com oflow-ai-agent)
Bootstrap do engine
O flow-ai-engine não sobe servidor HTTP — é um processo de consumo de Redis Streams.
O index.ts segue o padrão comum a todos os serviços: await initTelemetry("engine") registra a instrumentação (OpenTelemetry) e resolve a config da UI antes de importar dinamicamente ./bootstrap.js, garantindo que ioredis/pg entrem já instrumentados.
O bootstrap.ts então:
- Chama
setupStreams()para garantir streams e consumer groups. - Chama
bootstrapFlowsCache()para buscar flows ativos no Postgres, validar cadaFlowDefinitionpublicada e gravar tudo na chave Redisflows. - Assina
bot:invalidate(pub/sub) para recarregar o cache de flows quando o core publica/republica um flow. - Inicia o subscriber de inatividade (
startInactivitySubscriber()). - Instala os handlers de graceful shutdown (
installShutdownHandlers,drainTimeoutMs=35spara acomodar uma tool em execução) e registraprisma.$disconnectemonPostDrain. - Inicia
consumeFlow()em loop infinito.
Entrada via stream:flow
O engine recebe FlowStreamEvent, uma união discriminada por kind. Cada variante dispara um ponto de entrada diferente do runtime (ver Pontos de entrada do run):
kind | Produtor | Variante |
|---|---|---|
userInput | orchestrator (forward de mensagem do contato, mode=flow) | traz messages: WhatsAppIncomingMessage[] |
humanSessionEnded | core (ao fechar um ticket humano) | retoma o HumanAttendanceBlock; jumpTarget? reposiciona |
agentSessionEnded | agent (ao encerrar a sessão de agente) | retoma o AgentBlock; jumpTarget? reposiciona |
externalJump | core (endpoint POST /sessions/:sessionId/redirect) | traz flowId, blockId, executeOnEntry |
Campos comuns a todas as variantes: sessionId, chatId, to, phoneNumberId, timestamp. A variante userInput:
ts
{
kind: "userInput"
sessionId: string
chatId: string
to: string
phoneNumberId: string
messages: WhatsAppIncomingMessage[]
timestamp: number
}Política de erro do consumer:
| Caso | Ação |
|---|---|
Sem payload | Loga warning e faz XACK. |
| JSON inválido ou shape inválido | Loga warning e faz XACK. |
Erro no handleFlowEvent() | Loga erro e não faz XACK. |
| Processamento concluído | Faz XACK em stream:flow/engine. |
Pontos de entrada do run
handleFlowEvent() carrega a sessão, resolve o flow ativo no cache e despacha por kind:
userInput→runFlow(ctx)humanSessionEnded→runFlowFromHumanSessionEnded(ctx)agentSessionEnded→runFlowFromAgentSessionEnded(ctx)externalJump→runFlowFromExternalJump(ctx)humanSessionEnded/agentSessionEndedcomjumpTargetválido →runFlowFromSessionEndedWithJump(ctx, jump)
Todos terminam chamando walkBlocks(ctx) (Fase 2).
Contexto de execução
runFlow() recebe um ExecutionContext mutável:
ts
type ExecutionContext = {
session: SessionState
event: FlowStreamEvent
activeFlow: CachedFlow
/** Cache completo de flows — necessário para resolver o flow destino em redirectToBot. */
flowsCache: FlowsCachePayload
/** Objeto de input estruturado do batch atual (null se sem mensagens ou humanSessionEnded). */
currentInput: InputObject | null
/** Sinaliza que uma ação de redirect foi executada; o loop principal deve trocar de flow. */
redirectSignal?: { executeOnEntry: boolean; deleted?: boolean }
/** Contador de hops de redirect na execução atual; protege contra loops entre flows. */
redirectHopCount: number
/** Variáveis estáticas do router — disponíveis via {{router.key}}. Nunca gravadas na sessão. */
routerVars: Record<string, string>
/** Variáveis estáticas do flow — disponíveis via {{flow.key}}. Nunca gravadas na sessão. */
flowVars: Record<string, string>
/** TTL em segundos para renovação da flag de debug. Presente apenas quando debug está ativo. */
debugTtlSeconds?: number
}Fase 1: processar input pendente
Roda apenas quando session.currentBlockId aponta para um bloco existente (usuário respondendo a um bloco que pausou aguardando input).
- Carrega o bloco atual.
- Se o bloco é
humanAttendanceouagent(commodecorrespondente), é race condition — retorna sem alterar a sessão. O blocobacktem tratamento próprio (mapeia a seleção da lista para o item de header/histórico/footer). - Extrai
currentInputcomoInputObjectda primeira mensagem do batch. - Grava
currentInput.contentemsession.variables[input.saveInputAs](se configurado). - Executa
block.outputActions. - Executa
globalActions.outputActions(envelope global — saída). - Avalia
outputConditionsem ordem — primeira que casar vence. - Se nenhuma casar, usa
defaultOutput. - Navega para o target escolhido.
Quando session.currentBlockId está vazio, a sessão entra no onboardingBlockId; se aponta para um bloco inexistente, cai no defaultExceptionBlockId.
Global Actions (envelope por bloco)
FlowDefinition aceita um campo opcional globalActions: { inputActions, outputActions }. Essas actions envolvem todos os blocos do flow como um envelope — sem interferir na lógica individual de cada bloco.
Modelo de envelope (ordem de execução)
→ globalActions.inputActions ← abre o envelope (scope: global)
→ block.inputActions (scope: block)
→ publicação de content
[aguarda input do usuário, se não for bypass]
→ block.outputActions (scope: block)
← globalActions.outputActions ← fecha o envelope (scope: global)
→ evaluateOutputConditions → navigate- Aplica-se a
StandardBlockeHumanAttendanceBlock. - As global actions têm acesso ao mesmo
ExecutionContextdos blocos — interpolação idêntica. session.currentBlockesession.currentFlowsão atualizados antes das global inputActions, tornando{{session.currentBlock.name}}disponível nelas.- Campo opcional — flows sem
globalActionscontinuam funcionando normalmente.
Casos de uso típicos
- Registrar eventos de analytics a cada transição (
recordEvent). - Gravar variáveis de auditoria com o nome do bloco/flow atual.
- Chamar um endpoint externo em toda entrada ou saída de bloco.
scope no debugger
Eventos de action no debugger carregam scope: "global" | "block" para distinguir a origem. O debugger agrupa ações globais em itens separados na coluna esquerda ("Ações Globais → Entrada" e "Ações Globais → Saída"), fora do contexto do bloco.
Fase 2: caminhar pelo flow
Entra em blocos e continua enquanto o bloco atual puder avançar sem esperar input.
- Detecta loop por bloco repetido no mesmo run.
- Interrompe se passar de
MAX_BYPASS_HOPS = 30. - Busca o bloco por
session.currentBlockId. - Atualiza
session.currentFlowesession.currentBlockna entrada do bloco. - Se o bloco é
agent, mudasession.mode = "agent", gravasession.agentStatee devolve umAgentStreamEventpara o handler publicar emstream:agent(ver Bloco de agente). OAgentBlocknão tem inputActions. - Executa
globalActions.inputActions(envelope global — entrada). - Executa
block.inputActions. - Se o bloco é
humanAttendance, mudasession.mode = "human", publica o handoff emstream:helpdesk(publishHandoff), estende oexpiryTimeoutpara o tempo restante da janela Meta (mínimo de 1h) e retorna mantendo ocurrentBlockId. - Para bloco standard, publica o
contentdo canal da sessão (resolvido pelo prefixo dosessionId:ig-/wc-/wa) nostream:outgoing-{meta,chat,ig}correspondente (block.content[channel], com fallback paracontent.waquando o canal está vazio). - Se
input.bypass=false, pausa aguardando input. - Se
input.bypass=true, executablock.outputActions, depoisglobalActions.outputActions, escolhe target e continua.
O bloco back (lista interativa de histórico) também é tratado aqui: monta as seções de header/histórico/footer e pausa aguardando seleção.
Actions suportadas
| Action | Efeito |
|---|---|
setContextVariable | Interpola o valor e grava em session.variables. |
setContactVariable | Interpola o valor e grava em contact.contactVariables no Postgres. |
executeScript | Executa código JS em V8 Isolate (isolated-vm). Ver Ação executeScript. |
httpCall | Requisição HTTP com fetch nativo, timeout configurável, headers e variáveis de saída. |
callEndpoint | Chama um endpoint cadastrado na plataforma (wrapper sobre httpCall com resolução de URL). |
executeTool | Executa uma ferramenta publicada (mini-flow) da plataforma; emite stream:tool-executions. |
recordEvent | Registra um evento de analytics vinculado ao contato; emite stream:analytics. |
acceptOptIn / revokeOptIn | Altera o opt-in do contato. |
redirectToBot | Redireciona a sessão para outro flow do mesmo router. Ver Redirecionamentos. |
returnToFlow | Desempilha o último entry de session.redirectHistory e retorna ao flow de origem. |
As primitivas de execução (interpolate, evaluateConditions/pickOutputTarget, executeInSandbox, resolveParamMappingValue, executeTool) vivem em packages/flow-ai-runtime e são compartilhadas entre o engine e o flow-ai-agent.
Actions com I/O externo (executeScript, httpCall, callEndpoint, executeTool) têm um campo onError: continue grava a variável de erro e segue; blocking faz freeze + handoff (ou, se errorHandoff === false, envia a errorMessage e cai no defaultExceptionBlockId).
Todas as actions aceitam um campo opcional conditions: Condition[] — se as condições falharem, a action é pulada (emite action_skip no debugger).
Comparators suportados em conditions
A avaliação de condições vive em packages/flow-ai-runtime/src/conditions.ts (evaluateConditions + pickOutputTarget), compartilhada entre o engine e o flow-ai-agent.
equals,notEquals,contains,notContainsexists,notExists,matchesstartsWith,endsWithisEmpty,isNotEmptygreaterThan,lessThan,greaterOrEqual,lessOrEqual(numéricos)
Saídas do run
| Condição | Resultado |
|---|---|
| Bloco pausou aguardando input | deleted=false; handler regrava sessão no Redis. |
Bloco humanAttendance | deleted=false; session.mode vira human; publishHandoff → stream:helpdesk; handler regrava sessão. |
Bloco agent | deleted=false; session.mode vira agent; handler publica AgentStreamEvent em stream:agent. |
Target endAttendance (reset) | deleted=true; handler remove a sessão do Redis e publica SummaryRequestedEvent em stream:summary. |
Target endAttendance (soft) | deleted=false; publica a endAttendanceMessage (se houver), mantém a sessão sem currentBlockId. |
| Loop / hop limit / bloco inexistente / erro em action, condition ou publish | Freeze + handoff: freezeAndReturn grava session.frozenByError e retorna runtimeError; o handler persiste SessionRuntimeError, marca a sessão frozen (markSessionFrozen), publica a mensagem de erro ao contato, muda session.mode = "human" e faz publishHandoff → stream:helpdesk. |
Erro fora do runFlow() no handler | Consumer não faz ACK; entry de stream:flow fica pendente. |
Congelamento por erro de runtime
Quando o loop principal detecta um problema, a sessão é congelada (não redirecionada). O campo session.frozenByError.phase identifica a origem:
inputActions · publishContent · outputActions · loop · hopLimit (MAX_BYPASS_HOPS = 30) · blockNotFound · phase1
Uma sessão frozen permanece em mode: "human" (não existe mode: "frozen"). A exceção a esse congelamento é uma action blocking com errorHandoff === false: em vez de congelar, o engine envia a errorMessage e redireciona para o defaultExceptionBlockId via redirectSignal, mantendo mode: "flow". (Roteamentos defensivos — input com shape inesperado, target variable, bloco de retomada inválido, jump cross-router — também caem no defaultExceptionBlockId, sem congelar.)
Streams emitidos pelo engine
Além do outbound (stream:outgoing-meta/-chat/-ig), o engine também emite:
stream:agent— ao pausar numAgentBlock(passa o controle aoflow-ai-agent).stream:helpdesk— handoff para atendimento humano (humanAttendanceou freeze por erro).stream:summary— ao deletar a sessão emendAttendancereset (SummaryRequestedEvent).stream:tool-executions— uma linha porexecuteTool(métrica alimentada pelo Studio).stream:analytics— a cadarecordEvent.
Nota:
stream:outgoing(sem sufixo) é legado/morto — os streams de saída vivos sãostream:outgoing-meta,stream:outgoing-chatestream:outgoing-ig.
Bloco de agente (AgentBlock)
Ao entrar num AgentBlock, o engine pausa o flow e delega o atendimento ao flow-ai-agent:
- Muda
session.mode = "agent"e ajusta oexpiryTimeout(usaAgentBlock.inputTimeout). - Inicializa
session.agentState(agente ativo,agentChain,configSnapshot,turnCount). - Devolve um
AgentStreamEventnoRunResult; o handler o publica emstream:agentapós persistir a sessão.
O flow-ai-agent roda o loop de tool-calling do LLM e, ao encerrar, publica um AgentSessionEndedEvent (kind: "agentSessionEnded") de volta em stream:flow — o engine então avalia as outputConditions do AgentBlock e segue o flow.