# NATS JetStream

Configuracao do NATS JetStream 2.10 para mensageria assincrona e processamento de eventos no Monetarie PIX.

## Pre-requisitos

- Cluster NATS JetStream 2.10 com 3 nos (minimo)
- Conectividade de rede entre os pods PIX e o cluster NATS
- Variavel de ambiente `NATS_ENABLED=true` para ativar workers
- Conhecimento basico de NATS JetStream (streams, consumers, subjects)

## Arquitetura

O Monetarie PIX utiliza NATS JetStream para comunicacao assincrona entre servicos internos e com o Core Banking. Os workers processam mensagens em paralelo com backpressure adaptativo.

```mermaid
flowchart TB
    subgraph NATS["Cluster NATS JetStream (3 nos)"]
        S1["MONETARIE_SPI<br/>monetarie.spi.>"]
        S2["MONETARIE_DICT<br/>monetarie.dict.>"]
        S3["MONETARIE_SETTLEMENT<br/>monetarie.settlement.>"]
        S4["MONETARIE_AUDIT<br/>monetarie.audit.>"]
        S5["MONETARIE_DLQ<br/>monetarie.dlq.>"]
        S6["MONETARIE_CORE<br/>monetarie.core.>"]
    end

    subgraph Workers["Workers PIX"]
        IP["InboundProcessor<br/>batch:100 / conc:10"]
        OS["OutboundSender<br/>batch:100 / conc:10"]
        SU["StatusUpdater<br/>batch:100 / conc:10"]
        RP["ReturnProcessor<br/>batch:100 / conc:5"]
        CEP["CoreEventProcessor<br/>batch:50 / conc:5"]
        SCH["Scheduler<br/>settlement cycles"]
        FI["FileImporter<br/>file processing"]
    end

    subgraph External["Externos"]
        CORE["Core Banking"]
        BACEN["BACEN"]
    end

    S1 --> IP & SU
    IP --> S1
    OS --> S1
    RP --> S1
    SU --> S1
    S6 --> CEP
    CEP --> S1 & S2
    SCH --> S3
    FI --> S3
    CORE --> S6
    S1 --> CORE
    S5 -.-> IP & OS & SU & RP & CEP

    style S5 fill:#f66,stroke:#333
    style NATS fill:#e8f5e9,stroke:#333
```

## Cluster NATS

### Topologia de Producao

O cluster NATS roda em VMs dedicadas (nao no GKE):

| No | Hostname | IP | Porta |
|----|----------|-----|-------|
| 1 | nats-dev-1 | 10.10.40.5 | 4222 |
| 2 | nats-dev-2 | 10.10.40.7 | 4222 |
| 3 | nats-dev-3 | 10.10.40.4 | 4222 |

### Configuracao do Servidor NATS

```hcl
# nats-server.conf
server_name: nats-dev-1
listen: 0.0.0.0:4222

jetstream {
  store_dir: /data/jetstream
  max_mem: 4G
  max_file: 50G
}

cluster {
  name: monetarie-nats
  listen: 0.0.0.0:6222
  routes: [
    nats-route://10.10.40.5:6222
    nats-route://10.10.40.7:6222
    nats-route://10.10.40.4:6222
  ]
}

# Limites
max_payload: 8MB
max_pending: 64MB
```

## Streams

O Monetarie PIX define 7 streams JetStream, cada um capturando um conjunto de subjects:

| Stream | Subjects | Retencao | max_bytes | max_msgs | Discard | Proposito |
|--------|----------|----------|-----------|----------|---------|-----------|
| `MONETARIE_SPI` | `monetarie.spi.>` | 7 dias | 10GB | 10.000.000 | old | Eventos de transacoes SPI |
| `MONETARIE_DICT` | `monetarie.dict.>` | 7 dias | 5GB | 5.000.000 | old | Eventos de chaves PIX |
| `MONETARIE_SETTLEMENT` | `monetarie.settlement.>` | 7 dias | 5GB | 5.000.000 | old | Eventos de liquidacao/reconciliacao |
| `MONETARIE_AUDIT` | `monetarie.audit.>` | 90 dias | 50GB | 50.000.000 | old | Trilha de auditoria |
| `MONETARIE_DLQ` | `monetarie.dlq.>` | 90 dias | 10GB | 10.000.000 | old | Dead Letter Queue |
| `MONETARIE_CORE` | `monetarie.core.>` | 7 dias | 5GB | 5.000.000 | old | Eventos Core Banking <-> PIX |

### Criacao de Streams

Os streams sao criados automaticamente pelos workers na inicializacao via `$JS.API.STREAM.CREATE.*`. Exemplo da estrutura enviada:

```json
{
  "name": "MONETARIE_SPI",
  "subjects": ["monetarie.spi.>"],
  "retention": "limits",
  "max_consumers": -1,
  "max_msgs": 10000000,
  "max_bytes": 10737418240,
  "max_age": 604800000000000,
  "storage": "file",
  "discard": "old",
  "num_replicas": 3
}
```

### Subjects por Stream

#### MONETARIE_SPI

| Subject | Direcao | Descricao |
|---------|---------|-----------|
| `monetarie.spi.transaction.created` | PIX -> Core | Transacao criada |
| `monetarie.spi.transaction.accepted` | PIX -> Core | Transacao aceita (ACSP) |
| `monetarie.spi.transaction.settled` | PIX -> Core | Transacao liquidada (STLD) |
| `monetarie.spi.transaction.rejected` | PIX -> Core | Transacao rejeitada (RJCT) |
| `monetarie.spi.transaction.returned` | PIX -> Core | Devolucao processada (RTRN) |
| `monetarie.spi.transaction.cancelled` | PIX -> Core | Transacao cancelada (CANC) |
| `monetarie.spi.inbound.*` | BACEN -> PIX | Mensagens recebidas do SPI |
| `monetarie.spi.outbound.*` | PIX -> BACEN | Mensagens enviadas ao SPI |

#### MONETARIE_DICT

| Subject | Direcao | Descricao |
|---------|---------|-----------|
| `monetarie.dict.keys.created` | PIX -> Core | Chave PIX criada |
| `monetarie.dict.keys.deleted` | PIX -> Core | Chave PIX removida |
| `monetarie.dict.keys.updated` | PIX -> Core | Chave PIX atualizada |
| `monetarie.dict.claims.*` | PIX -> Core | Eventos de portabilidade |
| `monetarie.dict.infractions.*` | PIX -> Core | Eventos de infracoes MED 2.0 |

#### MONETARIE_CORE

| Subject | Direcao | Descricao |
|---------|---------|-----------|
| `monetarie.core.pix.payment_request` | Core -> PIX | Solicitacao de pagamento PIX |
| `monetarie.core.pix.return_request` | Core -> PIX | Solicitacao de devolucao |
| `monetarie.core.pix.key_create` | Core -> PIX | Criar chave PIX |
| `monetarie.core.pix.key_delete` | Core -> PIX | Remover chave PIX |
| `monetarie.core.pix.dict_lookup` | Core -> PIX | Consulta DICT |
| `monetarie.core.pix.balance_inquiry` | Core -> PIX | Consulta de saldo |

::: danger JETSTREAM E REQUEST/REPLY
**Subjects de request/reply DEVEM estar fora do escopo de streams JetStream.** Se um subject `monetarie.spi.*` for usado para request/reply, o JetStream intercepta a resposta e corrompe o fluxo. Use subjects como `dict.lookup.request` e `dict.api.request` para request/reply.
:::

## Consumers

Cada worker possui seu proprio consumer duravel:

| Worker | Consumer | Stream | Filter | AckWait | MaxDeliver |
|--------|----------|--------|--------|---------|------------|
| InboundProcessor | `spi-inbound` | MONETARIE_SPI | `monetarie.spi.inbound.>` | 30s | 5 |
| OutboundSender | `spi-outbound` | MONETARIE_SPI | `monetarie.spi.outbound.>` | 30s | 5 |
| StatusUpdater | `spi-status` | MONETARIE_SPI | `monetarie.spi.transaction.>` | 30s | 5 |
| ReturnProcessor | `spi-return` | MONETARIE_SPI | `monetarie.spi.return.>` | 30s | 5 |
| CoreEventProcessor | `core-events` | MONETARIE_CORE | `monetarie.core.pix.>` | 30s | 5 |
| Scheduler | `settlement-sched` | MONETARIE_SETTLEMENT | `monetarie.settlement.>` | 60s | 3 |
| FileImporter | `settlement-files` | MONETARIE_SETTLEMENT | `monetarie.settlement.files.>` | 60s | 3 |

## Arquitetura de Workers

### BaseWorker

Todos os workers herdam de `Shared.Workers.BaseWorker`, que fornece:

```mermaid
flowchart TB
    BW["BaseWorker"] --> INIT["init/1<br/>Verifica NATS_ENABLED"]
    INIT -->|Habilitado| CONN["Conectar NATS<br/>Criar Stream/Consumer"]
    INIT -->|Desabilitado| STOP["Parar worker"]
    CONN --> POLL["Poll Loop<br/>(poll_interval adaptativo)"]
    POLL --> FETCH["Fetch batch<br/>(batch_size mensagens)"]
    FETCH -->|Mensagens| PROC["Processar<br/>(max_concurrency paralelo)"]
    FETCH -->|Vazio| BACK["Backpressure<br/>Aumentar intervalo"]
    PROC --> ACK["ACK / NAK"]
    ACK -->|Sucesso| POLL
    ACK -->|Falha max_retries| DLQ["Publicar no DLQ<br/>(3 retries com backoff)"]
    DLQ --> POLL
    BACK --> POLL
```

### Parametros dos Workers

| Worker | batch_size | poll_interval | max_concurrency | max_retries |
|--------|-----------|---------------|-----------------|-------------|
| InboundProcessor | 100 | 200ms | 10 | 5 |
| OutboundSender | 100 | 200ms | 10 | 5 |
| StatusUpdater | 100 | 200ms | 10 | 5 |
| ReturnProcessor | 100 | 200ms | 5 | 5 |
| CoreEventProcessor | 50 | 500ms | 5 | 5 |

### Backpressure Adaptativo

O BaseWorker ajusta automaticamente o `poll_interval` baseado na carga:

| Condicao | Intervalo | Descricao |
|----------|-----------|-----------|
| Batch cheio (= batch_size) | 50ms (minimo) | Maxima velocidade de consumo |
| Batch parcial (> 0, < batch_size) | Configurado (200-500ms) | Velocidade normal |
| Batch vazio (0 mensagens) | 2.000ms (maximo) | Reduz polling desnecessario |

## DLQ (Dead Letter Queue)

Mensagens que falham apos `max_retries` sao enviadas para o DLQ com metadados:

```json
{
  "original_subject": "monetarie.spi.inbound.pacs008",
  "original_body": "{...payload original...}",
  "error": "timeout connecting to database",
  "retry_count": 5,
  "worker": "InboundProcessor",
  "failed_at": "2026-02-13T14:30:00Z",
  "trace_id": "abc123def456"
}
```

O subject DLQ segue o padrao: `monetarie.dlq.{subject_original}`.

A publicacao no DLQ usa 3 tentativas com backoff exponencial (1s, 2s, 4s) para garantir que nenhuma mensagem seja perdida silenciosamente.

### Entrega Agendada (JetStream)

O `NatsPublisher` do DICT utiliza o header `Nats-Msg-Deliver-After` do JetStream para agendamento duravel de mensagens, substituindo o volatil `Process.send_after`:

```elixir
# Entrega agendada em 30 segundos (duravel, sobrevive a crashes)
headers = [{"Nats-Msg-Deliver-After", "30s"}]
Gnat.pub(conn, subject, payload, headers: headers)
```

## Monitoramento NATS

### QueueMonitor

O `SettlementService.Monitoring.QueueMonitor` verifica o consumer lag a cada 15 segundos:

```mermaid
flowchart LR
    QM["QueueMonitor<br/>(15s)"] --> CS["Consumer Stats<br/>via $JS.API"]
    CS --> LAG["Calcular Lag<br/>(pending msgs)"]
    LAG --> WS["Broadcast via<br/>WebSocket"]
    WS --> UI["Dashboard<br/>queues:depth"]
```

### Metricas Prometheus

O endpoint `/metrics` expoe metricas NATS:

| Metrica | Tipo | Descricao |
|---------|------|-----------|
| `pix_worker_message_processed_total` | Counter | Mensagens processadas por worker |
| `pix_worker_batch_completed_total` | Counter | Batches completados por worker |
| `pix_worker_poll_duration_seconds` | Histogram | Duracao do poll por worker |
| `pix_nats_consumer_lag` | Gauge | Mensagens pendentes por consumer |

## Troubleshooting

### Verificar Streams

```bash
# Via nats CLI
nats stream ls
nats stream info MONETARIE_SPI
nats consumer ls MONETARIE_SPI

# Verificar consumer lag
nats consumer info MONETARIE_SPI spi-inbound
```

### Problemas Comuns

| Problema | Causa | Solucao |
|----------|-------|---------|
| Workers nao iniciam | `NATS_ENABLED=false` | Definir `NATS_ENABLED=true` |
| CoreEventProcessor timeout | Stream MONETARIE_CORE nao existe | Aguardar Core team criar o stream |
| Mensagens no DLQ | Falha persistente no processamento | Verificar logs do worker, corrigir causa raiz |
| Consumer lag crescente | Workers nao acompanham a carga | Aumentar `max_concurrency` ou adicionar pods |
| Request/reply corrompido | Subject dentro do escopo JetStream | Mover subject para fora do escopo do stream |

## Resultado Esperado

Apos a configuracao do NATS JetStream:

- Os 7 streams estao criados com retencao e limites configurados
- Todos os workers iniciam sem erros quando `NATS_ENABLED=true`
- O consumer lag e monitorado via WebSocket no dashboard
- Mensagens com falha sao redirecionadas para o DLQ com metadados completos
- O backpressure adaptativo reduz polling quando nao ha mensagens
- Os logs mostram `"[Shared.Application] NATS enabled, starting NATS supervisor"`
