Pular para conteúdo

0003 — Arquitetura: integrador + tradutor (banco compartilhado, webhook + sweep)

Status: Aprovado · Responsável: Gustavo Madruga · Atualizado em: 2026-08-11 · Decidido em: 2026-07-15

Contexto

A integração leva preço e estoque dos produtos do X-Adm ao e-commerce do parceiro WebStorm (https://www.webstorm.com.br), via POST /sync/erp. Era preciso decidir por onde os dados trafegam entre o X-Adm e a API do parceiro e quem faz cada parte, com dois requisitos: (a) o X-Adm envia no seu próprio modelo ERP de produto/estoque; (b) precisamos de auditoria do que entrou e do que saiu (com a resposta).

Duas restrições moldam a solução:

  1. Um app dedicado (o tradutor) tem valor além do POST. Sendo um serviço próprio, ele é a casa natural de futuras ferramentas que consultam preço/produto (ex. um bot de WhatsApp que responde consultas). Por isso mantemos o tradutor como app separado, não uma rotina embutida.
  2. Os dois apps compartilham o mesmo banco db_thoms. O tradutor lê a réplica direto no Postgres — não precisa de um canal que transporte o dado até ele.

Decisão

Nuvem de 3 hops — X-Adm → integrador → parceiro — com dois apps sobre o mesmo banco db_thoms:

X-Adm ──POST (modelo ERP)──▶ Integrador (int.thoms.xadm.biz)
   ├─ grava request de entrada
   ├─ consolida a RÉPLICA produto/estoque (db_thoms, com coluna de versão)
   └─ (A) POST "produto X mudou"  ─────────▶ Tradutor (webstorm-ecom.thoms.xadm.biz)
                                              ├─ lê a réplica DIRETO no db_thoms
                                              ├─ mapeia → POST /sync/erp (@Retryable)
                                              └─ grava webstorm_ecom_request (envio + resposta)
   (D) Tradutor também varre a réplica por @Scheduled: versão > última enviada → processa

Blocos e donos

  • Integrador (int.thoms.xadm.biz) — recebe o JSON do X-Adm, grava o log de requests de entrada e consolida a réplica (produto/estoque) no db_thoms. É do projeto integrador (repo xadm/integrador-server, https://docs.xadm.biz/aplicacoes/integrador-server/; app Java/Micronaut multi-tenant, um deploy por cliente), mesmo time; aqui é referenciado.
  • Tradutor (webstorm-ecom.thoms.xadm.biz) — app Micronaut dedicado enxuto: lê a réplica direto no db_thoms, mapeia para o contrato do parceiro, faz o POST /sync/erp e grava webstorm_ecom_request. Micronaut Data/JDBC para as tabelas próprias, Views/JTE (decisão 0025 da casa) para as telas de auditoria/reprocesso. É o mini-projeto deste repo.

Banco compartilhado db_thoms — fronteira de dados

Os dois apps usam o mesmo banco. Para não colidirem:

  • Flyway próprio por app. Cada app tem sua própria history table e seu próprio conjunto de migrations (locations). O tradutor usa uma history dedicada (ex. flyway_schema_history_webstorm_ecom), de modo que migrate de um app não vê nem altera o histórico do outro. Sem isso, um Flyway trata as migrations do outro como "faltando" e quebra.
  • Prefixo webstorm_ecom_*. Todas as tabelas do tradutor são prefixadas (webstorm_ecom_request, etc.) — namespace legível no banco compartilhado, zero colisão de nome.
  • A réplica é do integrador — read-only para o tradutor. O Flyway do tradutor gerencia só as tabelas webstorm_ecom_*. As tabelas da réplica pertencem ao integrador (ele as migra); o tradutor as lê, nunca as migra.
  • A réplica já existe: a tabela estoque. O integrador é um espelho do X-Adm e já mantém estoque, alimentada para a Thoms pelo PUT JSON de saída do X-Adm (igual Sul Plata/OnPetro) — não pelo fluxo de entrada PIED (Maxsul). Identidade do produto na Thoms = cod_prod (chave_est chega null); correlação WebStorm por ean13; colunas nome_prod/venda/venda_pz/ saldo. A coluna pied_codigo_alt (transporte PIED) vem null na Thoms. A integração WebStorm lê essa tabela (read-only) — não cria tabela nem coluna de produto. Tipos reais: § Modelo de dados; modelo de identidade: decisão 0013 do integrador.
  • updated_at na estoque (feito). O sweep (D) varre por coluna de mudança; o integrador adicionou updated_at TIMESTAMPTZ (+ deleted) na estoque, bumpado no apply do X-Adm. O read-model Estoque do tradutor já lê ambas; o sweep usa updated_at > watermark.

Transporte integrador → tradutor: webhook (A) + sweep (D)

Como o dado já está na réplica compartilhada, o integrador só precisa avisar que algo mudou — e mesmo o aviso não precisa ser confiável, porque o banco é a rede de segurança:

  • (A) Webhook — caminho de baixa latência. Ao consolidar, o integrador faz um POST numa URL do tradutor ("produto X mudou", só o id/EAN), protegido por Bearer entre os dois apps e com @Retryable. O tradutor lê a réplica, mapeia, posta e grava. Reusa o padrão de notificação assíncrona (push) que o integrador já opera.
  • (D) Sweep — rede de segurança. Um @Scheduled no tradutor varre estoque por linhas cuja updated_at > última processada (cruzando com webstorm_ecom_request) e processa. Pega tudo que o webhook perdeu (tradutor fora do ar, POST falho, aviso sumido). É o padrão Outbox / Polling Publisher: com produtor e consumidor no mesmo banco, o próprio banco é a fila.

Convergência: (A) e (D) chamam o mesmo código (ler réplica → mapear → POST → gravar). A idempotência por versão no webstorm_ecom_request garante que os dois não duplicam envio.

Conceitos transversais

  • Idempotência por estado/versão da linha: o webstorm_ecom_request registra o que já foi enviado; reprocessar a mesma versão não gera POST duplicado. A API do parceiro é state-based (seta preço/estoque), last-write-wins — só precisamos de at-least-once, que o sweep garante.
  • Ordem = um escritor por vez (advisory lock de cluster). O sweep é single-flight em dois níveis: lock em memória (SingleFlight, coalesce dentro de 1 JVM) e advisory lock do Postgres (SweepClusterLock → pg_try_advisory_lock(hashtext('webstorm_ecom_sweep'))) que serializa o corpo do sweep ENTRE instâncias no db_thoms compartilhado. Envios ao mesmo produto saem seriais, cada um relendo o estado atual da réplica (uma linha por produto), então o parceiro converge para o último valor mesmo com falhas/retries. O payload do parceiro não tem guarda de versão — ele não rejeita escrita velha —, então a ordem depende inteiramente de haver um único escritor por vez; o advisory lock garante isso mesmo com N instâncias idênticas (as demais coalescem ao próximo ciclo/poke). É HA/failover: chave por nome (namespaced, não colide com o integrador no mesmo db_thoms), lock de sessão → se a detentora cai, a sessão morre e o lock solta sozinho, e outra instância assume. Antes o invariante era "rodar em 1 réplica" (só o lock em-memória); o advisory lock removeu essa restrição — escalar horizontal é seguro para o sweep. A guarda de versão por linha continua vindo do pendentes() (updated_at > MAX(enviado com sucesso)), agora livre de corrida entre instâncias.
  • Retry + circuit breaker (parceiro vivo, mas afoga sob carga). Sintoma observado: reprocesso de 1 item = HTTP 200 na hora, mas o sweep em lote estoura ReadTimeout 100% (o parceiro para de responder sob a rajada — provável proteção/rate-limit disparada pelo "forçar reenvio total"). Então: o client re-tenta só transporte (@Retryable, 2 tentativas — retry demais amplifica), com read-timeout de 20s (service-id webstorm-parceiro; não 60s: parceiro que blackholeia só seguraria a conexão mais tempo). Quem realmente protege é o circuit breaker no Sweep: após N lotes seguidos falhos (PARCEIRO_MAX_FALHAS, default 3) ele aborta o ciclo; os pendentes não avançam o watermark e retomam no próximo sweep, dando fôlego ao parceiro em vez de queimar o catálogo inteiro em timeouts. Falha persistente fica gravada no webstorm_ecom_request (auditoria).
  • Lotes: o POST /sync/erp é particionado em lotes configuráveis (PARCEIRO_LOTE, default 10; faixa útil 10–100, limite do endpoint 200), com pacing entre lotes (PARCEIRO_PACING_MS, default 1s). Lote pequeno responde bem antes do read-timeout; o pacing evita o burst martelar o parceiro lote atrás de lote. Calibrar com o parceiro: se o teto for nº de requests (rate-limit) e não tempo por request, subir o lote (menos requests) em vez de descer — tudo via env, sem redeploy.
  • dry_run do parceiro. As respostas 200 vêm com dry_run:true/atualizados:0 (modo homologação da WebStorm) — mesmo o sucesso não altera preço/estoque no e-commerce. A virada para produção é do parceiro (WebStorm), não deste lado — sem ação nossa.
  • Auth: as telas de auditoria ficam atrás do middleware do Coolify. O webhook do integrador é protegido por Bearer app→app (mesma casa).
  • Carga inicial (semear a base do parceiro) segue por CSV/e-mail, independente do fluxo contínuo.

Consequências

  • App dedicado preservado: o tradutor é um serviço próprio (base para o bot de WhatsApp e outras consultas de preço/produto), mas sem as peças pesadas de sync.
  • Zero infra de transporte nova: sem broker, sem serviço de sync. O par banco compartilhado + sweep entrega confiabilidade; o webhook só melhora a latência.
  • Acesso direto ao dado: o tradutor lê a réplica no db_thoms sem intermediário — o mesmo acesso que a futura ferramenta de consulta vai usar.
  • Fronteira de dados explícita: Flyway por app + prefixo webstorm_ecom_* + réplica read-only mantêm dois donos sobre um banco sem pisar um no outro.
  • Fronteira do parceiro inalterada: a API atualiza preço/estoque de variantes existentes (match por EAN, decisão 0002); não cadastra produto novo; hoje responde em dry_run.

Alternativas descartadas

  • PowerSync (ponte de sync no tradutor) — desenho anterior desta decisão. Existe para levar o dado a um consumidor sem acesso ao banco da fonte. Com db_thoms compartilhado, o tradutor lê a réplica direto — a premissa cai. Sobraria SQLite local espelhando dado do mesmo Postgres + bridge Kotlin (PowerSync não tem SDK Java) + serviço de sync: overhead injustificado.
  • Broker de fila (RabbitMQ/Kafka/JMS) — daria durabilidade/retry prontos, mas seria um servidor a mais para fornecer a durabilidade que o Postgres compartilhado já dá. Justifica só se o banco deixar de ser compartilhado ou surgir fan-out para muitos consumidores. É uma troca de mensagem simples (um tipo, um produtor, um consumidor) — não pede broker.
  • LISTEN/NOTIFY do Postgres — nativo e real-time, mas sem persistência (o notify some se ninguém escuta): precisaria do sweep de catch-up de qualquer forma. Complexidade de LISTEN sobre o mesmo D — pior custo/benefício que A+D.
  • Webhook sozinho (sem sweep) — fire-and-forget perde mensagem se o tradutor estiver fora do ar. Por isso o webhook é otimização de latência sobre o sweep, nunca a garantia de entrega.
  • Rotina embutida no integrador (sem app separado) — mais simples para só o POST, mas o tradutor separado é requisito (casa das consultas de preço/produto).
  • Um 2º parceiro de e-commerce não está no escopo — o tradutor implementa o contrato deste parceiro direto; refatora-se se surgir.

Adendo (2026-07) — CSV vira o fluxo de entrada; o tradutor ganha porta de ingestão

O X-Adm decidiu não enviar produto item-a-item via JSON: passa a POSTar o catálogo inteiro em CSV (ProdutosSite.csv, dentro de um .zip). Isto não substitui o transporte integrador→tradutor (webhook + sweep, acima), que segue intacto — muda de onde a réplica é alimentada. Um adapter no tradutor (POST /api/csv/processar) recebe o zip, faz diff contra estado próprio e transcodifica para o PUT/DELETE /api/v1/xadm do integrador (se passando pelo X-Adm). O integrador consolida a réplica e emite EstoqueMudouEvent — daí o fluxo é o mesmo desta decisão (sweep → POST /sync/erp).

Consequências: (a) o tradutor deixa de ser só leitor da réplica — ganha uma porta de ingestão; (b) o adapter é uma ponte temporária (quando o X-Adm postar JSON direto no integrador, deleta-se o adapter e nada mais muda); (c) STATUS=0 do CSV vira soft-delete no integrador → o tradutor envia saldo = 0 (é assim que a ativação/STATUS é tratada — decidido). Detalhe do adapter (diff-store, lotes, at-least-once, paridade de campos): Modelo de dados § Ingestão de produtos via CSV. Linhagem de trabalho: .ia/005-ingestao-csv-catalogo-*.