Todos os artigos
FT-RETRY-WORKER· 1.0.0Integrações

Pipeline de mensagens — fluxos e cenários

Doc canônico do pipeline de mensagens com 7 cenários (sucesso, retry-success, retry-DLQ, retry-manual, 4xx não-retentável, cert-expirando, cancelar). Estados, classifier de erros, onde ver tudo.

Atualizado em 8/6/2026
Documento canônico de como uma mensagem caminha pelo Lefia, do recebimento até estado final. Cobre os 7 cenários possíveis e mostra os componentes envolvidos em cada um.

Visão geral do pipeline



[parceiro/SAP] → POST /api/edi-hub/receive-message ou /api/webhook/[id]

[insert messages com pipeline_status='received']

[processMessage síncrono em src/lib/pipeline/process-message.ts]
├─ duplicate check (por message_type)
├─ map (se há mapping configurado)
└─ send (se há connection destino)
├─ sucesso → pipeline_status='sent' / dest_status='200'
└─ falha → pipeline_status='error' / dest_status_text=erro

[tryEnqueueRetry em src/lib/queue/retry-orchestrator.ts]
├─ retry_auto_enabled=false → fica em 'error' (usuário)
├─ erro não retentável → fica em 'error'
├─ retry_count<max → 'retry_pending' + job em job_queue
└─ esgotou → 'dead_letter' (se DLQ on) ou 'error'


O worker /api/cron/process-queue puxa jobs prontos da job_queue (RPC dequeue_jobs com FOR UPDATE SKIP LOCKED) e re-executa processMessage pra cada um. Em nova falha, chama tryEnqueueRetry de novo.

Estados de pipeline_status



StatusSignificadoPróximo passo
receivedMensagem entrou no banco; processMessage ainda não rodouPipeline executa
mappedField mapping aplicado, aguardando enviosendToDestination
sentEntregue ao destino com sucessoFinal
errorFalhou ou esgotou retries (sem DLQ)Humano clica reprocessar
retry_pendingFalhou mas retry automático agendadoWorker pega no run_after
dead_letterEsgotou retries com DLQ habilitadoAlerta + revisão manual em /admin/queue
cancelledHumano cancelou retry em vôoFinal
received_observedModo observability — recebeu mas não enviouFinal (modo monitor)


7 cenários de execução



Cenário A — Sucesso direto



Mensagem entra, mapeia, envia, destino retorna 200.

  • pipeline_status: received → mapped → sent
  • dest_status: "200"
  • Webhook event message.sent é disparado pra integrações de saída
  • Logs em message_logs step=sent status=success


Cenário B — Erro 5xx + retry automático ON + sucesso na 2ª tentativa



Integração com retry_auto_enabled=true, destino retorna 503.

  • 1ª tentativa: pipeline_status: received → error (dest='503')
  • tryEnqueueRetry classifica http_5xx (retryable), calcula nextRetryAt (ex: +60s exponential), atualiza:
- pipeline_status: 'retry_pending' - retry_count: 1 - next_retry_at: now+60s - last_retry_error_class: 'http_5xx'
  • Job inserido em job_queue com run_after=now+60s
  • Worker executa, processMessage funciona desta vez → pipeline_status: 'sent'


Cenário C — Erro 5xx + retry automático ON + esgota tentativas → DLQ



Mesmo cenário B, mas todas as N tentativas falham.

  • Tentativas 1..retry_max: retry_pending com backoff (30s, 60s, 120s, ...)
  • Após retry_max: tryEnqueueRetry retorna exhausted_dlq
  • Se retry_dead_letter=true: pipeline_status='dead_letter'
  • Senão: fica em 'error'
  • Cron /api/cron/dlq-alerts agrupa por integração + bucket horário, dispara Slack e email pra RESEND_NOTIFY_EMAIL
  • Aparece em /admin/queue aba "Dead-letter" com botões "Reprocessar" e "Cancelar"


Cenário D — Retry automático OFF (default) → fica parado pro usuário



Integração com retry_auto_enabled=false (o padrão de produto).

  • Falha → pipeline_status: 'error' (sem retry, sem job na fila)
  • Mensagem fica visível no monitor
  • Operador clica "Reprocessar" → /api/edi-hub/reprocess-message reseta status e re-executa pipeline
  • Se reprocessar manual também falha, mesmo fluxo: ou volta pro usuário (auto OFF) ou enfileira retry (auto ON)


Cenário E — Erro 4xx (payload inválido) → não retenta nunca



Destino retorna 400/403/404/422 — problema do payload, não transiente.

  • tryEnqueueRetry classifica como http_4xx ou syntax
  • retryable=false — independente de retry_auto_enabled, não enfileira
  • Atualiza apenas last_retry_error_class
  • Mensagem fica em error pro usuário corrigir o mapping/payload
  • 401 é exceção: classificado como http_401 retryable (token pode ter renovado)
  • 429 também: classificado como http_429, parseia Retry-After do response


Cenário F — Certificado digital expirando (alerta proativo)



Não é fluxo de mensagem, mas afeta integrações que usam mTLS (SEFAZ, etc).

  • Cron /api/cron/cert-expiry-alerts roda 07:00 UTC
  • Varre connections.certificates.client.not_after
  • Buckets: 30, 14, 7, 1 e 0 dias
  • 1 alerta por bucket via Slack (config SLACK_ALERTS_WEBHOOK_URL)
  • UI: integration-card mostra badge "Cert vence em Xd" amarelo (≤30d) ou vermelho (≤7d)
  • Renovação automática: trigger connections_cert_expiry_alerts_reset apaga rows antigas quando not_after muda — buckets resetam sozinhos


Cenário G — Operador cancela retry em vôo



Mensagem em retry_pending (ex: retry_count=2 of 5), usuário decide parar.

  • Em /admin/queue aba "Aguardando" ou "Dead-letter", botão X (cancelar)
  • POST /api/admin/messages/[id]/cancel-retry:
- Apaga jobs pendentes de job_queue (started_at IS NULL)
- pipeline_status: 'cancelled' - next_retry_at: null - Audit event message.retry_cancelled registrado
  • Worker eventual nunca executa o job (porque já foi removido)


Cenário H — Modo observability (read-only)



Integração com mode='observability' — default pra novos tenants.

  • Mensagem chega via webhook/API → processMessage roda normal
  • Mapping aplica (se configurado) → pipeline_status: 'mapped'
  • Antes de chamar sendToDestination, helper skipSend intercepta:
- Atualiza pipeline_status: 'received_observed' - dest_status: 'skipped_observability' - dest_status_text: 'Modo observability — envio ao destino bloqueado' - Log step=observability registrado
- Não chama sendData — destino do cliente fica intocado
  • tryEnqueueRetry no-op (não há erro pra retentar; mensagem foi observada com sucesso)
  • Card da integração mostra badge READ-ONLY azul persistente


Use pra: rollout enterprise progressivo. Cliente roda Lefia em paralelo com iPaaS atual sem risco de duplicar dado. Quando confiar, alterna pra orchestration no edit da integração — UI esconde retry policy enquanto observability (irrelevante).

Cenário I — Reentrega exata bloqueada (idempotency)



Parceiro reenvia exatamente o mesmo payload (timeout, falha de rede do lado dele).

Camadas de proteção (do mais cedo pro mais tarde):

  1. Header HTTP Idempotency-Key — se cliente passou (Stripe-style)
- Sanitizado em sanitizeIdempotencyKey (RFC: ASCII printable, ≤255 chars)
- UNIQUE INDEX (tenant_id, integration_id, idempotency_key) - Bate → retorna mensagem original sem reinserir, response com deduplicated: true, deduplicated_by: 'idempotency_key'
  1. Fingerprint computada — sempre, mesmo sem header
- SHA-256 de (message_type_id + duplicate_key_fields normalizados + hash payload canônico) - Reaproveita os duplicate_key_fields da UI "Tratamento de Duplicatas"
- UNIQUE INDEX (tenant_id, integration_id, fingerprint) - Bate → retorna original com deduplicated_by: 'fingerprint'
  1. Race-condition fallback — se 2 requests concorrentes passaram pelo pré-check
- INSERT cai em 23505 (unique violation) — handler busca o winner e devolve com deduplicated_by: 'race_winner'

Distinção importante: idempotency NÃO substitui a UI "Tratamento de Duplicatas" (duplicate_handling por message_type). Eles cobrem cenários diferentes:
  • Idempotency: mesmo payload byte-a-byte (reentrega exata)
  • duplicate_handling: mesmas chaves de negócio mas conteúdo distinto (ex: 2 PO-12345 com itens diferentes — operador decide accept/reject/overwrite)


Idempotency atua antes do INSERT; duplicate_handling atua dentro do processMessage depois de inserir.

Onde ver o que



O queOnde
Mensagens em retry, jobs pendentes, DLQ/admin/queue (super_admin)
Toggle de retry automático por integraçãoEdit da integração → seção "Política de Retry"
Janela de horário (afeta retry)Edit da integração → seção "Agendamento"
Histórico de logs por mensagem/edi-hub/integracao/[id]
Erros do Lefia em si (não de integração)/admin/errors + Sentry
Alertas Slackcanal configurado em SLACK_ALERTS_WEBHOOK_URL
Email de DLQRESEND_NOTIFY_EMAIL


Plataforma SEFAZ — caso especial



Em "Plataforma" do form de criar/editar integração, escolher SEFAZ automaticamente:
  • Troca o protocolo pra SOAP (webservice nativo)
  • Mostra hint "SEFAZ exige mTLS — anexe o certificado digital A1 (.pfx/.p12)"
  • Campos clientCertPfx (file upload) e clientCertPassword ficam visíveis
  • Server-side: parseia PKCS#12, valida senha, encripta blob+senha (AES-256-GCM), grava em connections.certificates


Erros mapeados pelo classifier



src/lib/integrations/retry-policy.tsclassifyError():

SinalClasseRetentável?
HTTP 5xxhttp_5xxsim
HTTP 429 + Retry-Afterhttp_429sim (respeita header)
HTTP 401http_401sim (uma vez, token pode renovar)
HTTP 422syntaxNÃO (payload/schema inválido)
HTTP 4xx outroshttp_4xxNÃO
ETIMEDOUT, msg "timeout"timeoutsim
ECONNRESET, ECONNREFUSED, ENOTFOUNDconnectionsim
msg com "schema/parse/invalid xml"syntaxNÃO
outrosunknownsim (default seguro)


Privacidade no retry



  • Cert digital: blob e senha sempre encriptados em repouso (src/lib/crypto.ts AES-256-GCM)
  • Payload da mensagem: redaction de PII aplicada antes de qualquer chamada à IA do Lefia (src/lib/ai/redaction.ts) — vide help article pii-redaction-llm
  • Logs de retry: apenas classe do erro + dest_status_text (texto do destino), sem credenciais


Limitações conhecidas (deste commit)



  • Schedule do worker /api/cron/process-queue é diário em Vercel Hobby. Pra retries com delay < 24h o operador precisa cron externo (cron-job.org bate com Bearer CRON_SECRET) ou subir pra Vercel Pro (cron a cada 1min).
  • Idempotency em mensagens de integração não é universal — backlog #60.
  • A3 (token físico) não suportado — backlog (precisa Lefia Agent on-prem).
  • Validação CRL/OCSP/cadeia não é feita no upload do cert (delegada à contraparte).
  • Logging não-estruturado (Pino) — backlog #59.

Histórico de versões

  • 1.0.0

    Retry worker com fila durável em Postgres. Default retry MANUAL; integração liga "retry automático" explicitamente. Classifier 5xx/timeout/conn retentáveis vs 4xx/syntax não. Backoff fixed/linear/exp com cap. Janela horária respeitada. DLQ com alerta Slack+email agregado por integração+hour_bucket. UI /admin/queue. Cancel-retry com audit. SEFAZ→SOAP default + hint mTLS no form. 7 cenários documentados em help.

    5/6/2026 · new