IncomingEvent — caixa de entrada de eventos do accounts

TLDR: o nectar-charges passa a receber webhooks do accounts num endpoint próprio, grava cada evento em incoming_events antes de responder e o processa num job do Solid Queue com até 3 tentativas. É o inbox pattern, a ponta de recebimento do outbox (OutgoingEvent) do accounts. O primeiro evento assinado é o INSTALLMENT_OVERDUE, mas o use case de domínio dele fica para uma spec seguinte.

Contexto

O accounts publica eventos financeiros para os sistemas assinantes via WebhookSubscription (ver accounts/.project/docs/rules/synapse/outgoing_event_delivery.md, R-002). Cada entrega é um POST com Authorization: Bearer <token> e o envelope:

json { "event_id": "<uuid do IncomingEvent de origem no accounts>", "event_type": "INSTALLMENT_OVERDUE", "occurred_at": "2026-09-24T10:05:00+00:00", "data": { "customer_email": "...", "account_status": "...", "purchase": { "...": "..." } } }

O accounts considera a entrega feita com qualquer 2xx e tenta de novo (até 3 vezes) em qualquer outra resposta. Um reenvio pelo admin dele manda o mesmo event_id. Então o nectar-charges precisa:

  • garantir que tudo o que chegou fica gravado, para dar rastreabilidade e permitir reprocesso;
  • responder rápido, sem depender do sucesso do processamento de domínio;
  • deduplicar pelo event_id;
  • ter uma política de tentativas no processamento.

A estrutura espelha o IncomingEvent do synapse no accounts (accounts/.project/docs/rules/synapse/incoming_event_processing.md, R-001), adaptada ao Rails e ao Solid Queue, que já roda neste serviço (config/queue.yml, processo jobs no Procfile).

Objetivos

  • Endpoint POST /api/v1/incoming_events autenticado por token compartilhado com o accounts.
  • Tabela incoming_events com o envelope inteiro, status e contagem de tentativas.
  • Dedup por event_id: reentrega responde 2xx sem gravar nem processar de novo.
  • ProcessIncomingEventJob: um job por evento, 3 tentativas com 5s de intervalo, failed como estado final.
  • Registry event_type → use case pronto para receber o use case do INSTALLMENT_OVERDUE.

Fora de escopo

  • O use case de domínio do INSTALLMENT_OVERDUE: o que criar ou atualizar em Customer, Debit e Installment. O payload atual não fecha com o modelo daqui (só customer_email, sem name e phone; Debit sem chave para purchase.external_id). Vai numa spec própria.
  • Enquanto esse use case não existir, todo INSTALLMENT_OVERDUE recebido termina failed com no use case registered for event: INSTALLMENT_OVERDUE, igual ao R-001. O evento fica gravado e é reprocessável quando o use case entrar.
  • Tela ou ação de reprocesso. O resgate de um failed é pelo console (ver “Reprocesso”).
  • Cadastro da WebhookSubscription no accounts (é registro no admin de lá, sem código).
  • Alerta automático (Sentry, Slack) na falha definitiva.
  • Assinatura HMAC do corpo: o accounts não assina, só manda Bearer.

Mudanças

Todos os caminhos são relativos a modules/backend/.

Migration e modelo — db/migrate/<ts>_create_incoming_events.rb, app/models/incoming_event.rb

Coluna Tipo Regra
id bigint PK, padrão do projeto
event_id string, not null event_id do envelope; índice único — é o dedup
event_type string, not null ex.: INSTALLMENT_OVERDUE
occurred_at datetime, not null occurred_at do envelope
payload jsonb, not null o envelope inteiro, como chegou
status string, not null, default pending enumerize pending · processed · failed
attempts integer, not null, default 0 falhas de processamento já registradas
error_message text motivo da última falha
processed_at datetime preenchido em processed e em failed
created_at / updated_at datetime created_at é o momento do recebimento

Índices: event_id (único) e [status, created_at].

Modelo:

  • MAX_ATTEMPTS = 3.
  • mark_processed! → processed + processed_at.
  • register_failure!(message) → attempts += 1, grava error_message; ao chegar em MAX_ATTEMPTS, vira failed com processed_at.

Autenticação — app/controllers/api/v1/incoming_events_controller.rb

  • Herda de Api::V1::BaseController com skip_before_action :authenticate_request!, porque quem chama é o accounts, não um usuário com JWT.
  • before_action próprio compara o Bearer com ENV.fetch("ACCOUNTS_WEBHOOK_TOKEN") usando ActiveSupport::SecurityUtils.secure_compare. Token ausente ou diferente → 401.
  • O mesmo valor é cadastrado como token na WebhookSubscription do nectar-charges no accounts.

Ingestão — app/use_cases/incoming_events/create.rb

IncomingEvents::Create < Micro::Case, chamado pelo controller:

  1. Valida presença de event_id, event_type, occurred_at e data → senão Failure(:invalid_payload).
  2. Se já existe IncomingEvent com esse event_id → Success(result: { duplicate: true }), sem gravar nem enfileirar.
  3. Numa transação: cria o IncomingEvent e chama ProcessIncomingEventJob.perform_later(id). Como o Solid Queue grava o job no mesmo banco, os dois entram juntos ou nenhum entra: não existe evento pending sem job.
  4. Corrida de duas entregas simultâneas: o RecordNotUnique do índice é tratado como duplicata.

Respostas do controller:

Situação HTTP
Evento novo gravado 202 Accepted
event_id já recebido 200 OK
Token ausente ou errado 401
Envelope sem campo obrigatório 422

O 422 é consumido como falha pelo accounts, que tenta mais duas vezes e desiste — é o comportamento desejado para envelope malformado.

Processamento — app/jobs/process_incoming_event_job.rb, app/use_cases/incoming_events/process.rb

ProcessIncomingEventJob (fila default):

  • self.enqueue_after_transaction_commit = false, explícito, para o enfileiramento ficar dentro da transação da ingestão (passo 3 acima) independentemente do default do Rails.
  • Carrega o evento; se não estiver pending, sai sem fazer nada (job repetido é inofensivo).
  • Chama IncomingEvents::Process dentro de uma transação. Success → mark_processed!.
  • Failure ou exceção → a transação do domínio é revertida, e register_failure! roda fora dela, para a tentativa não se perder junto com o rollback.
  • Se, depois da falha, o evento continua pending, reenfileira o mesmo job com set(wait: 5.seconds). A contagem vive na linha (attempts), não no executions do ActiveJob — é ela que dá a rastreabilidade.

IncomingEvents::Process < Micro::Case:

  • USE_CASES = {} — o registry event_type → classe. Nesta spec fica vazio.
  • event_type sem entrada → Failure(:no_use_case) com mensagem no use case registered for event: <event_type>, tratada como qualquer outra falha.
  • Com entrada → chama o use case passando data do envelope e devolve o resultado dele.

Tabela de decisão (igual ao R-001 do accounts):

Situação attempts antes status depois Reenfileira?
Use case devolve Success qualquer processed não
Falha ou exceção 0 ou 1 pending, attempts + 1 sim, em 5s
Falha ou exceção 2 failed, attempts: 3 não
Sem use case registrado 0 ou 1 pending, attempts + 1 sim, em 5s
Sem use case registrado 2 failed, attempts: 3 não

Reprocesso

Sem interface nesta spec. Pelo console:

ruby event.update!(status: :pending, attempts: 0, error_message: nil, processed_at: nil) ProcessIncomingEventJob.perform_later(event.id)

Rota e configuração

  • config/routes.rb: resources :incoming_events, only: [ :create ] em api/v1.
  • ACCOUNTS_WEBHOOK_TOKEN no .env de exemplo e nos secrets de staging e produção.

Seeds

Nenhum seed: incoming_events é dado operacional, não de cenário.

Como verificar

  • Testes de modelo: mark_processed!, register_failure! abaixo e no limite de tentativas.
  • Testes do IncomingEvents::Create: evento novo cria e enfileira; event_id repetido não cria nem enfileira; envelope incompleto falha.
  • Testes do controller: 202, 200 na duplicata, 401 sem token e com token errado, 422.
  • Testes do job: cada linha da tabela de decisão, incluindo o rollback das escritas do use case com attempts persistido, e o job sem efeito quando o evento não está pending.
  • Ponta a ponta em staging: cadastrar a WebhookSubscription do nectar-charges para INSTALLMENT_OVERDUE no accounts, disparar uma parcela atrasada e ver a linha em incoming_events terminar failed com no use case registered… depois de 3 tentativas, e o OutgoingEvent do accounts como delivered.
  • Antes do deploy, confirmar que as tabelas solid_queue_* existem no banco principal de produção. db/queue_schema.rb existe mas o database.yml não declara banco queue separado, e o NegativationQueueJob já depende disso.

Documentação

  • Criar .project/docs/rules/events/incoming_event_processing.md com a regra de ingestão, dedup e tentativas (o equivalente ao R-001 do accounts).
  • Registrar a nova regra e esta spec no .project/docs/README.md.