# Eventos NATS

Catalogo completo de subjects NATS JetStream da plataforma FluxiQ PIX, incluindo schemas de payload, configuracao de consumers e tratamento de DLQ.

## Pre-requisitos

- Entendimento basico de NATS JetStream (streams, consumers, subjects)
- Acesso ao cluster NATS (nats CLI)
- Conhecimento do padrao de routing por campo `"event"` nos payloads

## Visao Geral das Streams

```mermaid
graph LR
    subgraph Core Banking
        CORE_PUB[Publisher]
    end

    subgraph FluxiQ PIX
        subgraph Streams JetStream
            S_SPI[MONETARIE_SPI<br/>monetarie.spi.>]
            S_DICT[MONETARIE_DICT<br/>monetarie.dict.>]
            S_SETTLE[MONETARIE_SETTLEMENT<br/>monetarie.settlement.>]
            S_AUDIT[MONETARIE_AUDIT<br/>monetarie.audit.>]
            S_DLQ[MONETARIE_DLQ<br/>monetarie.dlq.>]
            S_CORE[MONETARIE_CORE<br/>monetarie.core.>]
        end

        subgraph Workers
            W_IN[InboundProcessor]
            W_OUT[OutboundSender]
            W_STATUS[StatusUpdater]
            W_RETURN[ReturnProcessor]
            W_CORE[CoreEventProcessor]
            W_SCHED[Scheduler]
            W_FILE[FileImporter]
        end
    end

    CORE_PUB -->|monetarie.core.pix.*| S_CORE
    S_CORE --> W_CORE
    S_SPI --> W_IN
    S_SPI --> W_OUT
    S_SPI --> W_STATUS
    S_SPI --> W_RETURN
    S_SETTLE --> W_SCHED
    S_SETTLE --> W_FILE
```

## Configuracao das Streams

| Stream | Subjects | Retencao | Max Bytes | Max Msgs | Storage | Replicas |
|--------|----------|----------|-----------|----------|---------|----------|
| MONETARIE_SPI | `monetarie.spi.>` | 7 dias | 10 GB | 10M | File | 1 |
| MONETARIE_DICT | `monetarie.dict.>` | 7 dias | 5 GB | 5M | File | 1 |
| MONETARIE_SETTLEMENT | `monetarie.settlement.>` | 7 dias | 5 GB | 5M | File | 1 |
| MONETARIE_AUDIT | `monetarie.audit.>` | 90 dias | 20 GB | 50M | File | 1 |
| MONETARIE_DLQ | `monetarie.dlq.>` | 90 dias | 2 GB | 1M | File | 1 |
| MONETARIE_CORE | `monetarie.core.>` | 7 dias | 5 GB | 5M | File | 1 |

Todas as streams usam `discard: :old` — quando atingem o limite, mensagens mais antigas sao descartadas.

## Subjects: Core -> PIX

Publicados pelo Core Banking, consumidos pelo `CoreEventProcessor`:

### monetarie.core.pix.payment_request

Solicitacao de pagamento PIX enviada pelo Core.

```json
{
  "event": "payment_request",
  "source": "core",
  "published_at": "2026-02-13T14:30:00.000Z",
  "trace_id": "a1b2c3d4e5f6a7b8",
  "data": {
    "amount": 15000,
    "currency": "BRL",
    "debtor": {
      "cpf_cnpj": "12345678901",
      "name": "Joao da Silva",
      "ispb": "12345678",
      "branch": "0001",
      "account": "12345-6"
    },
    "creditor": {
      "key": "98765432100",
      "key_type": "CPF"
    },
    "remittance_info": "Pagamento NF 12345",
    "idempotency_key": "pay-uuid-v4"
  }
}
```

**Handler:** `CoreEventProcessor.handle_payment_request/1` — Cria transacao outbound (pacs.008) no banco de dados e publica evento SPI.

### monetarie.core.pix.return_request

Solicitacao de devolucao PIX.

```json
{
  "event": "return_request",
  "source": "core",
  "published_at": "2026-02-13T14:35:00.000Z",
  "trace_id": "b2c3d4e5f6a7b8c9",
  "data": {
    "original_end_to_end_id": "E1234567820260213143000001",
    "amount": 15000,
    "reason": "MD06",
    "description": "Devolucao solicitada pelo pagador"
  }
}
```

**Handler:** `CoreEventProcessor.handle_return_request/1` — Cria pacs.004 de devolucao.

### monetarie.core.pix.key_create

Solicitacao de registro de chave PIX.

```json
{
  "event": "key_create",
  "source": "core",
  "published_at": "2026-02-13T15:00:00.000Z",
  "data": {
    "key_type": "CPF",
    "key_value": "12345678901",
    "owner": {
      "cpf_cnpj": "12345678901",
      "name": "Joao da Silva",
      "ispb": "12345678",
      "branch": "0001",
      "account": "12345-6",
      "account_type": "CACC"
    }
  }
}
```

**Handler:** `CoreEventProcessor.handle_key_create/1` — Proxy para DICT POST /api/v2/entries.

### monetarie.core.pix.key_delete

Solicitacao de remocao de chave PIX.

```json
{
  "event": "key_delete",
  "source": "core",
  "published_at": "2026-02-13T15:10:00.000Z",
  "data": {
    "key_type": "CPF",
    "key_value": "12345678901",
    "reason": "USER_REQUESTED"
  }
}
```

**Handler:** `CoreEventProcessor.handle_key_delete/1` — Proxy para DICT DELETE /api/v2/entries/:key.

### monetarie.core.pix.dict_lookup

Consulta de chave no DICT.

```json
{
  "event": "dict_lookup",
  "source": "core",
  "published_at": "2026-02-13T14:29:55.000Z",
  "data": {
    "key_type": "CPF",
    "key_value": "12345678901"
  }
}
```

**Handler:** `CoreEventProcessor.handle_dict_lookup/1` — Consulta DICT API e retorna resultado via NATS reply.

::: warning Subjects de request/reply
Os subjects `dict.lookup.request` e `dict.api.request` sao usados para request/reply NATS (NAO JetStream). Eles estao FORA do escopo dos streams JetStream. Se fossem capturados por um stream, o JetStream interceptaria as respostas.
:::

### monetarie.core.pix.balance_inquiry

Consulta de saldo SPI.

```json
{
  "event": "balance_inquiry",
  "source": "core",
  "published_at": "2026-02-13T14:29:50.000Z",
  "data": {
    "ispb": "12345678",
    "account_type": "PI"
  }
}
```

**Handler:** `CoreEventProcessor.handle_balance_inquiry/1` — Consulta saldo SPI e retorna via NATS reply.

## Subjects: PIX -> Core

Publicados pelo FluxiQ PIX, consumidos pelo Core Banking (PixHandler):

### monetarie.spi.transaction.*

Eventos de transacoes SPI:

| Subject | Evento | Quando |
|---------|--------|--------|
| `monetarie.spi.transaction.created` | Transacao criada | pacs.008 recebido ou criado |
| `monetarie.spi.transaction.updated` | Status atualizado | Qualquer mudanca de status |
| `monetarie.spi.transaction.completed` | Liquidacao concluida | Status = STLD |
| `monetarie.spi.transaction.rejected` | Transacao rejeitada | Status = RJCT |
| `monetarie.spi.transaction.returned` | Devolucao processada | pacs.004 confirmado |

Payload padrao:

```json
{
  "type": "transaction.completed",
  "source": "pix",
  "published_at": "2026-02-13T14:30:01.200Z",
  "trace_id": "a1b2c3d4e5f6a7b8",
  "data": {
    "id": "550e8400-e29b-41d4-a716-446655440000",
    "end_to_end_id": "E1234567820260213143000001",
    "status": "STLD",
    "status_id": 7,
    "amount": 15000,
    "message_type": "pacs.008",
    "direction": "OUTBOUND",
    "debtor_ispb": "12345678",
    "creditor_ispb": "53822116",
    "operation_time": "2026-02-13T14:30:00.000Z",
    "settlement_time": "2026-02-13T14:30:01.200Z"
  }
}
```

### monetarie.dict.keys.*

Eventos de chaves DICT:

| Subject | Evento | Quando |
|---------|--------|--------|
| `monetarie.dict.keys.created` | Chave registrada | Nova chave no DICT |
| `monetarie.dict.keys.deleted` | Chave removida | Chave removida do DICT |
| `monetarie.dict.keys.updated` | Chave atualizada | Dados da chave alterados |

Payload:

```json
{
  "type": "key.created",
  "source": "pix",
  "published_at": "2026-02-13T15:00:00.000Z",
  "trace_id": "c3d4e5f6a7b8c9d0",
  "data": {
    "key_type": "CPF",
    "key_value": "12345678901",
    "owner_name": "Joao da Silva",
    "owner_ispb": "12345678",
    "account_branch": "0001",
    "account_number": "12345-6"
  }
}
```

## Subjects Internos

### monetarie.settlement.*

| Subject | Descricao |
|---------|-----------|
| `monetarie.settlement.netting.completed` | Ciclo de netting concluido |
| `monetarie.settlement.reconciliation.started` | Reconciliacao iniciada |
| `monetarie.settlement.reconciliation.completed` | Reconciliacao concluida |
| `monetarie.settlement.fee.calculated` | Tarifas calculadas |

### monetarie.audit.*

| Subject | Descricao |
|---------|-----------|
| `monetarie.audit.xml.archived` | XML arquivado (SHA-256 hash) |
| `monetarie.audit.api_validation.completed` | Validacao XSD concluida |
| `monetarie.audit.user.login` | Login de usuario |
| `monetarie.audit.user.logout` | Logout de usuario |

### monetarie.dlq.*

Dead Letter Queue — mensagens que falharam apos max_retries (5 tentativas):

```json
{
  "original_subject": "monetarie.spi.transaction.created",
  "original_body": "{\"type\":\"transaction.created\",...}",
  "error": "Postgrex.Error: connection refused",
  "worker": "SpiService.Workers.InboundProcessor",
  "failed_at": "2026-02-13T14:31:00.000Z"
}
```

O subject da DLQ e `monetarie.dlq.{original_subject}`, permitindo filtragem por tipo de mensagem.

## Routing Key

O campo `"event"` no payload e a chave de roteamento usada pelo `CoreEventProcessor`:

```elixir
# CoreEventProcessor.process_message/2
def process_message(%{"event" => event} = message, _metadata) do
  case event do
    "payment_request" -> handle_payment_request(message)
    "return_request"  -> handle_return_request(message)
    "key_create"      -> handle_key_create(message)
    "key_delete"      -> handle_key_delete(message)
    "dict_lookup"     -> handle_dict_lookup(message)
    "balance_inquiry" -> handle_balance_inquiry(message)
    unknown ->
      Logger.warning("Unknown event type: #{unknown}")
      :ok  # ACK silencioso — evita redelivery infinita
  end
end
```

::: danger Campo "event" obrigatorio
Se o campo `"event"` estiver ausente ou com valor nao reconhecido, a mensagem sera ACKed silenciosamente sem processamento. Isso e por design para evitar redelivery infinita, mas pode causar perda silenciosa de dados se o publisher enviar eventos com typos.
:::

## Configuracao de Consumers

Cada worker cria um pull consumer com configuracao especifica:

| Worker | Stream | Consumer Name | Filter Subject | Batch | Concurrency |
|--------|--------|---------------|---------------|-------|-------------|
| InboundProcessor | MONETARIE_SPI | inbound-processor | monetarie.spi.transaction.created | 100 | 10 |
| OutboundSender | MONETARIE_SPI | outbound-sender | monetarie.spi.transaction.outbound | 100 | 10 |
| StatusUpdater | MONETARIE_SPI | status-updater | monetarie.spi.transaction.updated | 100 | 10 |
| ReturnProcessor | MONETARIE_SPI | return-processor | monetarie.spi.transaction.return | 100 | 5 |
| CoreEventProcessor | MONETARIE_CORE | core-event-processor | monetarie.core.pix.> | 50 | 5 |
| Scheduler | MONETARIE_SETTLEMENT | settlement-scheduler | monetarie.settlement.schedule.> | 10 | 1 |
| FileImporter | MONETARIE_SETTLEMENT | file-importer | monetarie.settlement.file.> | 10 | 1 |

Configuracao de cada consumer:

```json
{
  "durable_name": "inbound-processor",
  "filter_subject": "monetarie.spi.transaction.created",
  "ack_policy": "explicit",
  "max_deliver": 5,
  "ack_wait": 60000000000,
  "max_ack_pending": 100
}
```

| Parametro | Valor | Descricao |
|-----------|-------|-----------|
| `ack_policy` | explicit | ACK manual obrigatorio |
| `max_deliver` | 5 | Maximo de redeliveries antes de descarte |
| `ack_wait` | 60s | Tempo para ACK antes de redelivery |
| `max_ack_pending` | 100 | Maximo de mensagens sem ACK simultaneas |

## Tratamento de DLQ

Mensagens que falham apos 5 tentativas sao enviadas para a DLQ com metadata de retry:

```mermaid
graph LR
    A[Mensagem] --> B{Processamento}
    B -->|Sucesso| C[ACK]
    B -->|Falha| D{Tentativa < 5?}
    D -->|Sim| E[NAK → Redelivery]
    E --> B
    D -->|Nao| F[ACK + DLQ]
    F --> G[monetarie.dlq.*]

    style C fill:#27ae60,color:#fff
    style F fill:#e74c3c,color:#fff
    style G fill:#e74c3c,color:#fff
```

O DLQ publish tem 3 retries com backoff exponencial (100ms, 200ms, 300ms) para evitar perda silenciosa de dados.

### Monitorar DLQ

```bash
# Verificar mensagens na DLQ
nats consumer info MONETARIE_DLQ dlq-monitor \
  --server=nats://10.10.40.5:4222

# Consumir mensagens da DLQ para reprocessamento manual
nats consumer next MONETARIE_DLQ dlq-monitor --count=10 \
  --server=nats://10.10.40.5:4222
```

### Replay de Mensagens da DLQ

```bash
# Ler mensagem da DLQ
nats consumer next MONETARIE_DLQ dlq-monitor --count=1 \
  --server=nats://10.10.40.5:4222

# Republicar mensagem original no subject correto
# Extrair original_subject e original_body do payload DLQ
nats pub monetarie.spi.transaction.created '{"type":"transaction.created",...}' \
  --server=nats://10.10.40.5:4222
```

## Request/Reply (Fora do JetStream)

Alguns subjects usam NATS request/reply puro (NAO JetStream):

| Subject | Direcao | Descricao |
|---------|---------|-----------|
| `dict.lookup.request` | Core -> PIX | Consulta de chave DICT |
| `dict.api.request` | Core -> PIX | Claims/MED/infracoes |

Estes subjects estao propositalmente fora do escopo dos streams JetStream. Se fossem capturados pelo stream `MONETARIE_SPI` (que escuta `monetarie.spi.>`), o JetStream interceptaria as respostas de request/reply, causando timeouts e respostas corrompidas.

```bash
# Exemplo de request/reply para lookup
nats request dict.lookup.request '{"key_type":"CPF","key_value":"12345678901"}' \
  --server=nats://10.10.40.5:4222 --timeout=5s

# Saida esperada:
# {
#   "status": "ok",
#   "data": {
#     "key_type": "CPF",
#     "key_value": "12345678901",
#     "owner_name": "Joao da Silva",
#     "owner_ispb": "12345678"
#   }
# }
```

## Propagacao de Trace ID

O trace ID (`x-trace-id`) e propagado em todas as mensagens NATS:

1. **Publicacao**: `Shared.Nats.JetStream.publish/3` extrai `trace_id` do `Logger.metadata` e injeta como header NATS
2. **Consumo**: `BaseWorker.set_trace_id/2` extrai o header `x-trace-id` da mensagem e define `Logger.metadata(trace_id: ...)`
3. **Fallback**: Se o header nao existe, tenta extrair do campo `"trace_id"` do payload JSON

## Resultado Esperado

Ao utilizar este catalogo de eventos NATS, voce sera capaz de:

- Publicar e consumir eventos nos 6 streams JetStream corretamente
- Implementar handlers para os 6 tipos de eventos Core -> PIX
- Consumir eventos PIX -> Core com os campos `type`, `source` e `published_at`
- Configurar consumers pull com parametros adequados por worker
- Monitorar e reprocessar mensagens da DLQ
- Entender a separacao entre subjects JetStream e request/reply
- Propagar trace IDs para rastreamento end-to-end
