# NATS JetStream

Monetarie PIX 中 NATS JetStream 的配置和管理，实现持久的 at-least-once 消息投递。

## 前提条件

- NATS Server 2.10+ 并启用 JetStream
- 后端 pods 到 NATS 集群的网络连接
- 了解发布/订阅和流的概念

## 架构

Monetarie PIX 使用 NATS JetStream 实现服务间以及与 Core Banking 的异步事件驱动通信。

```mermaid
graph LR
    SPI[SPI 服务] -->|发布| NATS[NATS JetStream<br/>7 个流]
    DICT[Dict 服务] -->|发布| NATS
    SETTLE[Settlement] -->|发布| NATS
    NATS -->|消费| WORKERS[NATS Workers]
    NATS <-->|事件| CORE[Core Banking]
```

## 流

| 流 | 主题 | 保留 | 最大时间 | 最大字节 | 丢弃策略 |
|----|------|------|----------|----------|----------|
| `MONETARIE_SPI` | `monetarie.spi.>` | 7 天 | 7天 | 已配置 | old |
| `MONETARIE_DICT` | `monetarie.dict.>` | 7 天 | 7天 | 已配置 | old |
| `MONETARIE_SETTLEMENT` | `monetarie.settlement.>` | 7 天 | 7天 | 已配置 | old |
| `MONETARIE_AUDIT` | `monetarie.audit.>` | 90 天 | 90天 | 已配置 | old |
| `MONETARIE_DLQ` | `monetarie.dlq.>` | 90 天 | 90天 | 已配置 | old |
| `MONETARIE_CORE` | `monetarie.core.>` | 7 天 | 7天 | 已配置 | old |

::: warning JETSTREAM + 请求/回复重叠
NATS 请求/回复主题必须在 JetStream 流范围之外。`MONETARIE_SPI` 流捕获 `monetarie.spi.>`，因此任何在 `monetarie.spi.*` 上的请求/回复都会被 JetStream 截获（导致响应损坏）。请使用 `dict.*` 等主题进行请求/回复。
:::

## Workers

| Worker | 流 | 主题 | 批量 | 并发数 |
|--------|-----|------|------|--------|
| 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 文件 | -- | 1 |

### Worker 特性

- **并行处理**：使用可配置 `max_concurrency` 的 `Task.async_stream`
- **自适应背压**：轮询间隔根据批次满载程度调整（最小 50ms，最大 2000ms）
- **死信队列**：失败消息发送到 `monetarie.dlq.*`，包含 3 次重试和指数退避
- **安全 ACK/NAK**：ACK/NAK 失败时使用 Logger.error 以增加可见性
- **投递计数追踪**：安全的 Integer.parse 带回退值，防止因格式错误的头部导致崩溃

## 配置

### 环境变量

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

### NATS 服务器配置

```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 事件

### Core -> PIX（由 CoreEventProcessor 消费）

| 主题 | 描述 |
|------|------|
| `monetarie.core.pix.payment_request` | 创建出站 pacs.008 |
| `monetarie.core.pix.return_request` | 创建 pacs.004 退回 |
| `monetarie.core.pix.key_create` | 通过 DICT 创建 PIX 密钥 |
| `monetarie.core.pix.key_delete` | 通过 DICT 删除 PIX 密钥 |
| `monetarie.core.pix.dict_lookup` | DICT 密钥查询 |
| `monetarie.core.pix.balance_inquiry` | SPI 余额查询 |

### PIX -> Core（由 SPI/DICT workers 发布）

`monetarie.spi.transaction.*` 和 `monetarie.dict.keys.*` 上的事件包含 Core 兼容字段：`type`、`source: "pix"`、`published_at`（ISO 8601）。

### 请求/回复主题（JetStream 之外）

| 主题 | 方向 | 用途 |
|------|------|------|
| `dict.lookup.request` | Core <-> PIX | DICT 密钥查询 |
| `dict.api.request` | Core <-> PIX | 认领/MED/违规 |

## 监控

NATS worker 遥测事件：

| 事件 | 测量值 | 标签 |
|------|--------|------|
| `[: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 |

## 预期结果

配置 NATS 后：

- 7 个 JetStream 流已创建，具有适当的保留策略
- 当 `NATS_ENABLED=true` 时，所有 5+ 个 workers 活跃消费消息
- Core Banking 事件双向流动
- 死信队列捕获失败消息以供调查
- 自适应背压维持最佳吞吐量
