# NATS JetStream

Configuration and management of NATS JetStream for durable, at-least-once message delivery in Monetarie PIX.

## Prerequisites

- NATS Server 2.10+ with JetStream enabled
- Network connectivity from backend pods to NATS cluster
- Understanding of pub/sub and stream concepts

## Architecture

Monetarie PIX uses NATS JetStream for asynchronous event-driven communication between services and with Core Banking.

```mermaid
graph LR
    SPI[SPI Service] -->|publish| NATS[NATS JetStream<br/>7 streams]
    DICT[Dict Service] -->|publish| NATS
    SETTLE[Settlement] -->|publish| NATS
    NATS -->|consume| WORKERS[NATS Workers]
    NATS <-->|events| CORE[Core Banking]
```

## Streams

| Stream | Subjects | Retention | Max Age | Max Bytes | Discard |
|--------|----------|-----------|---------|-----------|---------|
| `MONETARIE_SPI` | `monetarie.spi.>` | 7 days | 7d | Configured | old |
| `MONETARIE_DICT` | `monetarie.dict.>` | 7 days | 7d | Configured | old |
| `MONETARIE_SETTLEMENT` | `monetarie.settlement.>` | 7 days | 7d | Configured | old |
| `MONETARIE_AUDIT` | `monetarie.audit.>` | 90 days | 90d | Configured | old |
| `MONETARIE_DLQ` | `monetarie.dlq.>` | 90 days | 90d | Configured | old |
| `MONETARIE_CORE` | `monetarie.core.>` | 7 days | 7d | Configured | old |

::: warning JETSTREAM + REQUEST/REPLY OVERLAP
NATS request/reply subjects MUST be outside JetStream stream scope. The `MONETARIE_SPI` stream captures `monetarie.spi.>`, so any request/reply on `monetarie.spi.*` gets intercepted by JetStream (corrupt responses). Use subjects like `dict.*` for request/reply.
:::

## Workers

| Worker | Stream | Subjects | Batch | Concurrency |
|--------|--------|----------|-------|-------------|
| InboundProcessor | MONETARIE_SPI | `monetarie.spi.inbound.>` | 100 | 10 |
| OutboundSender | MONETARIE_SPI | `monetarie.spi.outbound.>` | 100 | 10 |
| StatusUpdater | MONETARIE_SPI | `monetarie.spi.status.>` | 100 | 10 |
| ReturnProcessor | MONETARIE_SPI | `monetarie.spi.return.>` | 100 | 5 |
| CoreEventProcessor | MONETARIE_CORE | `monetarie.core.pix.>` | 50 | 5 |
| Scheduler | MONETARIE_SETTLEMENT | `monetarie.settlement.>` | -- | 1 |
| FileImporter | MONETARIE_SETTLEMENT | Settlement files | -- | 1 |

### Worker Features

- **Parallel processing**: `Task.async_stream` with configurable `max_concurrency`
- **Adaptive backpressure**: Poll interval adjusts based on batch fullness (50ms min, 2000ms max)
- **Dead letter queue**: Failed messages sent to `monetarie.dlq.*` with 3 retry attempts and exponential backoff
- **Safe ACK/NAK**: Logger.error on ACK/NAK failures for visibility
- **Delivery count tracking**: Safe Integer.parse with fallback prevents crashes on malformed headers

## Configuration

### Environment Variables

```bash
NATS_HOST=10.10.40.5
NATS_PORT=4222
NATS_USER=optional_user
NATS_PASS=optional_password
NATS_ENABLED=true
```

### NATS Server Configuration

```conf
# nats-server.conf
port: 4222
jetstream {
    store_dir: /data/nats
    max_mem: 4G
    max_file: 100G
}
cluster {
    name: monetarie-pix
    listen: 0.0.0.0:6222
    routes: [
        nats://10.10.40.5:6222
        nats://10.10.40.7:6222
        nats://10.10.40.4:6222
    ]
}
```

## Core Banking Events

### Core -> PIX (consumed by CoreEventProcessor)

| Subject | Description |
|---------|-------------|
| `monetarie.core.pix.payment_request` | Create outbound pacs.008 |
| `monetarie.core.pix.return_request` | Create pacs.004 return |
| `monetarie.core.pix.key_create` | Create PIX key via DICT |
| `monetarie.core.pix.key_delete` | Delete PIX key via DICT |
| `monetarie.core.pix.dict_lookup` | DICT key lookup |
| `monetarie.core.pix.balance_inquiry` | SPI balance query |

### PIX -> Core (published by SPI/DICT workers)

Events on `monetarie.spi.transaction.*` and `monetarie.dict.keys.*` include Core-compatible fields: `type`, `source: "pix"`, `published_at` (ISO 8601).

### Request/Reply Subjects (Outside JetStream)

| Subject | Direction | Purpose |
|---------|-----------|---------|
| `dict.lookup.request` | Core <-> PIX | DICT key lookup |
| `dict.api.request` | Core <-> PIX | Claims/MED/infractions |

## Monitoring

NATS worker telemetry events:

| Event | Measurement | Tags |
|-------|-------------|------|
| `[:pix, :worker, :message_processed]` | duration_ms | worker, subject, status |
| `[:pix, :worker, :batch_completed]` | count, duration_ms | worker |
| `[:pix, :worker, :poll]` | poll_interval_ms | worker, batch_size |

## Expected Outcome

After configuring NATS:

- 7 JetStream streams created with appropriate retention
- All 5+ workers actively consuming messages (when `NATS_ENABLED=true`)
- Core Banking events flowing bidirectionally
- Dead letter queue capturing failed messages for investigation
- Adaptive backpressure maintaining optimal throughput
