Skip to content

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.ts
  • services/flow-ai-engine/src/bootstrap.ts
  • services/flow-ai-engine/src/consumer.ts
  • services/flow-ai-engine/src/handlers/flow.ts
  • services/flow-ai-engine/src/flows-cache.ts
  • services/flow-ai-engine/src/runtime/execute.ts
  • services/flow-ai-engine/src/runtime/actions.ts
  • services/flow-ai-engine/src/runtime/input.ts
  • services/flow-ai-engine/src/runtime/publish.ts
  • packages/flow-ai-runtime/src/conditions.ts (avaliação de condições, compartilhada com o flow-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:

  1. Chama setupStreams() para garantir streams e consumer groups.
  2. Chama bootstrapFlowsCache() para buscar flows ativos no Postgres, validar cada FlowDefinition publicada e gravar tudo na chave Redis flows.
  3. Assina bot:invalidate (pub/sub) para recarregar o cache de flows quando o core publica/republica um flow.
  4. Inicia o subscriber de inatividade (startInactivitySubscriber()).
  5. Instala os handlers de graceful shutdown (installShutdownHandlers, drainTimeoutMs=35s para acomodar uma tool em execução) e registra prisma.$disconnect em onPostDrain.
  6. 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):

kindProdutorVariante
userInputorchestrator (forward de mensagem do contato, mode=flow)traz messages: WhatsAppIncomingMessage[]
humanSessionEndedcore (ao fechar um ticket humano)retoma o HumanAttendanceBlock; jumpTarget? reposiciona
agentSessionEndedagent (ao encerrar a sessão de agente)retoma o AgentBlock; jumpTarget? reposiciona
externalJumpcore (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:

CasoAção
Sem payloadLoga warning e faz XACK.
JSON inválido ou shape inválidoLoga warning e faz XACK.
Erro no handleFlowEvent()Loga erro e não faz XACK.
Processamento concluídoFaz 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:

  • userInputrunFlow(ctx)
  • humanSessionEndedrunFlowFromHumanSessionEnded(ctx)
  • agentSessionEndedrunFlowFromAgentSessionEnded(ctx)
  • externalJumprunFlowFromExternalJump(ctx)
  • humanSessionEnded/agentSessionEnded com jumpTarget válidorunFlowFromSessionEndedWithJump(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).

  1. Carrega o bloco atual.
  2. Se o bloco é humanAttendance ou agent (com mode correspondente), é race condition — retorna sem alterar a sessão. O bloco back tem tratamento próprio (mapeia a seleção da lista para o item de header/histórico/footer).
  3. Extrai currentInput como InputObject da primeira mensagem do batch.
  4. Grava currentInput.content em session.variables[input.saveInputAs] (se configurado).
  5. Executa block.outputActions.
  6. Executa globalActions.outputActions (envelope global — saída).
  7. Avalia outputConditions em ordem — primeira que casar vence.
  8. Se nenhuma casar, usa defaultOutput.
  9. 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 StandardBlock e HumanAttendanceBlock.
  • As global actions têm acesso ao mesmo ExecutionContext dos blocos — interpolação idêntica.
  • session.currentBlock e session.currentFlow são atualizados antes das global inputActions, tornando {{session.currentBlock.name}} disponível nelas.
  • Campo opcional — flows sem globalActions continuam 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.

  1. Detecta loop por bloco repetido no mesmo run.
  2. Interrompe se passar de MAX_BYPASS_HOPS = 30.
  3. Busca o bloco por session.currentBlockId.
  4. Atualiza session.currentFlow e session.currentBlock na entrada do bloco.
  5. Se o bloco é agent, muda session.mode = "agent", grava session.agentState e devolve um AgentStreamEvent para o handler publicar em stream:agent (ver Bloco de agente). O AgentBlock não tem inputActions.
  6. Executa globalActions.inputActions (envelope global — entrada).
  7. Executa block.inputActions.
  8. Se o bloco é humanAttendance, muda session.mode = "human", publica o handoff em stream:helpdesk (publishHandoff), estende o expiryTimeout para o tempo restante da janela Meta (mínimo de 1h) e retorna mantendo o currentBlockId.
  9. Para bloco standard, publica o content do canal da sessão (resolvido pelo prefixo do sessionId: ig-/wc-/wa) no stream:outgoing-{meta,chat,ig} correspondente (block.content[channel], com fallback para content.wa quando o canal está vazio).
  10. Se input.bypass=false, pausa aguardando input.
  11. Se input.bypass=true, executa block.outputActions, depois globalActions.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

ActionEfeito
setContextVariableInterpola o valor e grava em session.variables.
setContactVariableInterpola o valor e grava em contact.contactVariables no Postgres.
executeScriptExecuta código JS em V8 Isolate (isolated-vm). Ver Ação executeScript.
httpCallRequisição HTTP com fetch nativo, timeout configurável, headers e variáveis de saída.
callEndpointChama um endpoint cadastrado na plataforma (wrapper sobre httpCall com resolução de URL).
executeToolExecuta uma ferramenta publicada (mini-flow) da plataforma; emite stream:tool-executions.
recordEventRegistra um evento de analytics vinculado ao contato; emite stream:analytics.
acceptOptIn / revokeOptInAltera o opt-in do contato.
redirectToBotRedireciona a sessão para outro flow do mesmo router. Ver Redirecionamentos.
returnToFlowDesempilha 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, notContains
  • exists, notExists, matches
  • startsWith, endsWith
  • isEmpty, isNotEmpty
  • greaterThan, lessThan, greaterOrEqual, lessOrEqual (numéricos)

Saídas do run

CondiçãoResultado
Bloco pausou aguardando inputdeleted=false; handler regrava sessão no Redis.
Bloco humanAttendancedeleted=false; session.mode vira human; publishHandoffstream:helpdesk; handler regrava sessão.
Bloco agentdeleted=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 publishFreeze + 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 publishHandoffstream:helpdesk.
Erro fora do runFlow() no handlerConsumer 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 num AgentBlock (passa o controle ao flow-ai-agent).
  • stream:helpdesk — handoff para atendimento humano (humanAttendance ou freeze por erro).
  • stream:summary — ao deletar a sessão em endAttendance reset (SummaryRequestedEvent).
  • stream:tool-executions — uma linha por executeTool (métrica alimentada pelo Studio).
  • stream:analytics — a cada recordEvent.

Nota: stream:outgoing (sem sufixo) é legado/morto — os streams de saída vivos são stream:outgoing-meta, stream:outgoing-chat e stream:outgoing-ig.

Bloco de agente (AgentBlock)

Ao entrar num AgentBlock, o engine pausa o flow e delega o atendimento ao flow-ai-agent:

  1. Muda session.mode = "agent" e ajusta o expiryTimeout (usa AgentBlock.inputTimeout).
  2. Inicializa session.agentState (agente ativo, agentChain, configSnapshot, turnCount).
  3. Devolve um AgentStreamEvent no RunResult; o handler o publica em stream: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.

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