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.
O worker
Mensagem entra, mapeia, envia, destino retorna 200.
Integração com
Mesmo cenário B, mas todas as N tentativas falham.
Integração com
Destino retorna 400/403/404/422 — problema do payload, não transiente.
Não é fluxo de mensagem, mas afeta integrações que usam mTLS (SEFAZ, etc).
Mensagem em
-
Integração com
- Não chama
Use pra: rollout enterprise progressivo. Cliente roda Lefia em paralelo com iPaaS atual sem risco de duplicar dado. Quando confiar, alterna pra
Parceiro reenvia exatamente o mesmo payload (timeout, falha de rede do lado dele).
Camadas de proteção (do mais cedo pro mais tarde):
- UNIQUE INDEX
- UNIQUE INDEX
Distinção importante: idempotency NÃO substitui a UI "Tratamento de Duplicatas" (
Idempotency atua antes do INSERT; duplicate_handling atua dentro do
Em "Plataforma" do form de criar/editar integração, escolher SEFAZ automaticamente:
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
| Status | Significado | Próximo passo |
|---|---|---|
received | Mensagem entrou no banco; processMessage ainda não rodou | Pipeline executa |
mapped | Field mapping aplicado, aguardando envio | sendToDestination |
sent | Entregue ao destino com sucesso | Final |
error | Falhou ou esgotou retries (sem DLQ) | Humano clica reprocessar |
retry_pending | Falhou mas retry automático agendado | Worker pega no run_after |
dead_letter | Esgotou retries com DLQ habilitado | Alerta + revisão manual em /admin/queue |
cancelled | Humano cancelou retry em vôo | Final |
received_observed | Modo observability — recebeu mas não enviou | Final (modo monitor) |
7 cenários de execução
Cenário A — Sucesso direto
Mensagem entra, mapeia, envia, destino retorna 200.
pipeline_status: received → mapped → sentdest_status: "200"- Webhook event
message.senté disparado pra integrações de saída - Logs em
message_logsstep=sentstatus=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') tryEnqueueRetryclassificahttp_5xx(retryable), calculanextRetryAt(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_queuecomrun_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_pendingcom backoff (30s, 60s, 120s, ...) - Após retry_max:
tryEnqueueRetryretornaexhausted_dlq - Se
retry_dead_letter=true:pipeline_status='dead_letter' - Senão: fica em
'error' - Cron
/api/cron/dlq-alertsagrupa por integração + bucket horário, dispara Slack e email praRESEND_NOTIFY_EMAIL - Aparece em
/admin/queueaba "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-messagereseta 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.
tryEnqueueRetryclassifica comohttp_4xxousyntaxretryable=false— independente de retry_auto_enabled, não enfileira- Atualiza apenas
last_retry_error_class - Mensagem fica em
errorpro usuário corrigir o mapping/payload - 401 é exceção: classificado como
http_401retryable (token pode ter renovado) - 429 também: classificado como
http_429, parseiaRetry-Afterdo 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-alertsroda 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_resetapaga rows antigas quandonot_aftermuda — 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/queueaba "Aguardando" ou "Dead-letter", botão X (cancelar) - POST
/api/admin/messages/[id]/cancel-retry:
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 →
processMessageroda normal - Mapping aplica (se configurado) →
pipeline_status: 'mapped' - Antes de chamar
sendToDestination, helperskipSendintercepta:
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
tryEnqueueRetryno-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):
- Header HTTP
Idempotency-Key— se cliente passou (Stripe-style)
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'
- Fingerprint computada — sempre, mesmo sem header
(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'
- Race-condition fallback — se 2 requests concorrentes passaram pelo pré-check
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 que | Onde |
|---|---|
| Mensagens em retry, jobs pendentes, DLQ | /admin/queue (super_admin) |
| Toggle de retry automático por integração | Edit 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 Slack | canal configurado em SLACK_ALERTS_WEBHOOK_URL |
| Email de DLQ | RESEND_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) eclientCertPasswordficam 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.ts → classifyError():| Sinal | Classe | Retentável? |
|---|---|---|
| HTTP 5xx | http_5xx | sim |
| HTTP 429 + Retry-After | http_429 | sim (respeita header) |
| HTTP 401 | http_401 | sim (uma vez, token pode renovar) |
| HTTP 422 | syntax | NÃO (payload/schema inválido) |
| HTTP 4xx outros | http_4xx | NÃO |
| ETIMEDOUT, msg "timeout" | timeout | sim |
| ECONNRESET, ECONNREFUSED, ENOTFOUND | connection | sim |
| msg com "schema/parse/invalid xml" | syntax | NÃO |
| outros | unknown | sim (default seguro) |
Privacidade no retry
- Cert digital: blob e senha sempre encriptados em repouso (
src/lib/crypto.tsAES-256-GCM) - Payload da mensagem: redaction de PII aplicada antes de qualquer chamada à IA do Lefia (
src/lib/ai/redaction.ts) — vide help articlepii-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 comBearer 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