Metadata-Version: 2.4
Name: pulse-mq
Version: 6.0.0
Summary: 高性能 Client/Server 消息中间件 — ZeroMQ ROUTER/DEALER + bcrypt 认证 + 实时监控
Project-URL: Homepage, https://github.com/haifeng/pulse-mq
Project-URL: Repository, https://github.com/haifeng/pulse-mq
Project-URL: Issues, https://github.com/haifeng/pulse-mq/issues
Author: haifeng
License-Expression: MIT
License-File: LICENSE
Keywords: finance,market-data,mq,pubsub,realtime,zmq
Classifier: Development Status :: 3 - Alpha
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Classifier: Topic :: System :: Networking
Classifier: Typing :: Typed
Requires-Python: >=3.13
Requires-Dist: bcrypt>=4.0
Requires-Dist: loguru>=0.7
Requires-Dist: lz4>=4.0
Requires-Dist: msgspec>=0.18
Requires-Dist: pandas>=2.0
Requires-Dist: pyarrow>=14.0
Requires-Dist: python-snappy>=0.7
Requires-Dist: pyzmq>=26.0
Requires-Dist: zstandard>=0.22
Description-Content-Type: text/markdown

# PulseMQ

面向金融行情的高性能消息中间件，基于 ZeroMQ ROUTER/DEALER 架构。

- **Client/Server 架构** — 服务端（ROUTER）集中路由 + 控制面；客户端（DEALER）发布/订阅
- **PLAIN + bcrypt 认证** — ZAP 认证链：bcrypt 哈希凭据存储，admin token 保护监控接口
- **类型保真** — DataFrame / dict / str / bytes 端到端保真（`_restore_type` 自动还原）
- **完整监控** — Web UI（ECharts）+ REST API + SSE 实时推送，延迟分位、事件流、在线客户端
- **运行期重连** — 断线自动指数退避重连（1s → 2s → 4s ... → 30s 封顶）

---

## 安装

> Python >= 3.13

```bash
pip install pulse-mq
```

依赖项：ZeroMQ、msgspec、python-snappy、lz4、zstandard、pyarrow、pandas、bcrypt、loguru。

```python
import pulsemq            # Python 模块名（无连字符）
from pulsemq import Server
```

> PyPI 分发名是 `pulse-mq`（`pip install` 用），import 名是 `pulsemq`（因 Python 标识符不允许连字符）。与 `python-dateutil` → `import dateutil` 模式一致。

---

## 快速开始

### 启动服务端

```python
from pulsemq import Server

srv = Server(
    data_endpoint="tcp://0.0.0.0:5555",
    control_endpoint="tcp://0.0.0.0:5556",
    admin_endpoint="0.0.0.0:9090",
    credentials={"user1": "pass1"},  # 或 credentials_file
)
await srv.start()
```

首次启动若不存在凭据文件，自动生成默认 `admin` 用户，密码输出到 stderr。
可通过 CLI 管理用户：

```bash
pulsemq users add user1 pass1 --role publisher
pulsemq users list
```

### 服务端内置定时推送

无需外部生产者客户端，直接在 Server 上注册定时回调：

```python
from pulsemq import Server

srv = Server(data_endpoint=..., control_endpoint=..., credentials={"u": "p"})

@srv.producer("market.tick", interval=2.0, serializer="msgpack")
async def gen_tick():
    return {"symbol": "AAPL", "price": 180.5, "volume": 1000}

@srv.producer("market.quote", interval=0.5, serializer="pyarrow", compression="lz4")
async def gen_quote():
    import pandas as pd
    return pd.DataFrame({"price": [10, 20], "vol": [100, 200]})

@srv.burst_producer("bench", serializer="msgpack")
async def bench():
    if not has_more():
        return None
    return {"seq": next_seq()}

await srv.start()   # 自动开始调度
```

| 方法 | 参数 | 说明 |
|------|------|------|
| `srv.producer(topic, interval, serializer, compression)` | `interval` 秒 | 固定间隔定时推送 |
| `srv.burst_producer(topic, serializer, compression)` | 无间隔 | 连续推送，返回 `None` 停止 |

回调返回值支持 DataFrame / dict / str / bytes，自动编码为协议帧并路由到所有匹配的消费者。

### 生产者

```python
import asyncio
from pulsemq.client import ProducerClient

async def main():
    prod = ProducerClient(
        "tcp://127.0.0.1:5555", "tcp://127.0.0.1:5556",
        username="user1", password="pass1",
    )
    await prod.start()
    await prod.publish("market.stock.AAPL", {"price": 180.5, "volume": 1000})
    await prod.stop()

asyncio.run(main())
```

### 消费者

```python
from pulsemq.client import ConsumerClient

async def main():
    cons = ConsumerClient(
        "tcp://127.0.0.1:5555", "tcp://127.0.0.1:5556",
        username="user2", password="pass2",
    )
    await cons.start()

    def on_msg(msg):
        print(msg.topic, msg.payload, msg.timestamp_ns)

    await cons.subscribe("market.*", on_msg)
    await asyncio.sleep(3600)  # 持续接收

asyncio.run(main())
```

---

## 架构

```
                      ┌──────────────────┐
                      │    Server        │
                      │  ┌─ 数据面 ─────┐│──── DEALER → 消费者
                      │  │ ROUTER :5555 ││
                      │  └──────────────┘│
                      │  ┌─ 控制面 ─────┐│
                      │  │ ROUTER :5556 ││←── DEALER  ← 生产者
                      │  └──────────────┘│
                      │  ┌─ Admin ──────┐│
                      │  │ HTTP :9090   ││─── REST / SSE / Web UI
                      │  └──────────────┘│
                      │  ZAP (bcrypt)    │
                      └──────────────────┘
```

| 端口 | 协议 | 用途 |
|------|------|------|
| `5555` | ROUTER (数据面) | 消息发布/接收 |
| `5556` | ROUTER (控制面) | REGISTER / HEARTBEAT / SUBSCRIBE / DISCONNECT |
| `9090` | HTTP | 监控 Web UI + REST API + SSE |

### 数据流

```
生产者 DEALER ──encode→  ROUTER  decode_header  match topic ──→ 消费者 DEALER
                              ↓
                         TrafficStats.record  LatencyStats.sample
```

服务端**不解压/不反序列化** payload（`decode_header` 仅提取头部），转发后由消费者完整 `decode` 还原。

### 控制面

| 命令 | 方向 | 作用 |
|------|------|------|
| `REGISTER` | Client → Server | 注册上线（含用户名、角色、订阅列表） |
| `HEARTBEAT` | Client → Server | 保活（每秒 1 次，6 秒超时自动下线） |
| `SUBSCRIBE` | Client → Server | 订阅 topic 模式 |
| `UNSUBSCRIBE` | Client → Server | 取消订阅 |
| `DISCONNECT` | Client → Server | 优雅下线 |

---

## 数据类型与序列化

### 支持的数据类型

| Python 类型 | 可用序列化器 | record_count |
|-------------|-------------|-------------|
| `pd.DataFrame` | msgpack, json, **pyarrow** | 行数 |
| `dict` | msgpack, json, **pyarrow** | 1 |
| `str` | **str**（仅此一种） | 1 |
| `bytes` | **bytes**（仅此一种） | 1 |

> pyarrow 对 DataFrame/dict 均可直接序列化；json/msgpack 下 DataFrame 先转 `list[dict]`
> 再以 `data_type=DATAFRAME` 标记，接收端 `_restore_type` 还原回 DataFrame。

### 序列化格式

| 格式 | 后端 | 适合 | 批处理场景 |
|------|------|------|-----------|
| `msgpack` | msgspec | 结构化小消息 ✅ | 批量 DataFrame 需先 to_dict |
| `json` | msgspec | 人类可读、跨语言 | 同上 |
| `pyarrow` | pyarrow IPC | 列存/分析 ✅ | **直接序列化 DataFrame，最快** |
| `str` | UTF-8 | 纯文本 | ❌ |
| `bytes` | 透传 | 二进制 | ❌ |

### 压缩算法

| 算法 | 适用场景 |
|------|---------|
| `none` | 小消息，极速 |
| `snappy` | 速度优先 |
| `lz4` | 批数据，平衡 |
| `zstd` | 压缩比优先（带宽受限） |

小消息场景压缩是**负收益**（计算开销 > 传输节省）；批量 DataFrame 场景 `lz4`/`zstd` 有明显效果。

---

## 客户端生命周期

### 首次启动

| 场景 | 异常 | exit code |
|------|------|-----------|
| 密码错误 | `AuthenticationError` | 3 |
| 服务器不可达 | `ClientStartupError` | 4 |

### 运行期重连

断线后 `Client._reconnect_loop` 按指数退避自动重连：

```
断线 → disconnected → cancel bg tasks → 新 Transport → PLAIN 认证
  → REGISTER（同 client_id）
    ├─ ALREADY_ONLINE → 退避重试（等待心跳超时释放）
    ├─ auth_failed → _reconnect_fatal → exit 3
    └─ OK → 恢复订阅 → 重启 recv/heartbeat
```

**业务无感**：订阅自动恢复，消息继续接收。

---

## 监控

### Web UI

浏览器打开 `http://localhost:9090/?token=<admin_token>` 查看实时面板：

- 4 个指标卡片：活跃主题、消息量/秒、流量/秒、运行时间
- 4 个客户端卡片：在线用户、生产者数、消费者数、订阅数
- ECharts 流量趋势折线图（分钟级，1H/6H 切换，最多 5 topic 叠加）
- 延迟 P50/P95/P99 柱状图
- 实时事件流（认证/连接/断线）
- 在线 Client 详情弹窗

### REST API

```bash
# 实时指标
curl 'http://localhost:9090/api/v1/stats/realtime?token=<token>'

# 主题列表
curl 'http://localhost:9090/api/v1/topics?token=<token>'

# 主题分钟级历史
curl 'http://localhost:9090/api/v1/topics/market.tick/history?minutes=60&token=<token>'

# 在线客户端明细
curl 'http://localhost:9090/api/v1/clients?token=<token>'

# 生命周期事件
curl 'http://localhost:9090/api/v1/events?token=<token>'

# 健康检查（无需 token）
curl http://localhost:9090/healthz
```

### SSE 实时流

```bash
curl -N 'http://localhost:9090/api/v1/stats/stream?token=<token>'
```

每 1 秒推送一帧 JSON，包含 topics / latency / online_users / sse_events 等。

---

## 协议帧格式

**单 bytes 帧**（非 ZMQ 多帧，通过 DEALER/ROUTER 传输）：

```
magic(2) ver(1) msg_type(1) flags(1) data_type(1) topic_len(2 BE)
topic(N) ts(8 BE ns) record_count(4 BE) payload(变长) [CRC32?(4)]
```

- `magic` = `"PM"`
- `msg_type` = DATA(0x01) / CONTROL(0x02)
- `flags` = 编码序列化器(3bit) + 压缩算法(2bit) + CRC(1bit)
- `data_type` = UNKNOWN(0x00) / DICT(0x01) / DATAFRAME(0x02) / STR(0x03) / BYTES(0x04)
- `record_count` 上限 **1,000,000**
- CRC 可选（由 flags 指示）

---

## 凭据管理

### CLI

```bash
# 添加用户（自动 bcrypt 哈希）
pulsemq users add trader1 secret123 --role publisher --role subscriber

# 启动服务端
pulsemq server

# 列出所有用户
pulsemq users list

# 禁用/启用用户
pulsemq users disable trader1
pulsemq users enable  trader1

# 修改密码
pulsemq users passwd trader1 new_secret

# 热加载凭据（SIGHUP 或 CLI）
pulsemq users reload
```

### 文件格式 (`pulsemq_users.toml`)

```toml
[users.admin]
hashed_password = "$2b$12$..."
roles = ["admin"]
enabled = true
created_at = "2026-06-27T00:00:00Z"
```

### Admin Token

首次启动自动生成 32 字节随机 base64url token，写入 `pulsemq_admin.token`（0600 权限）。
Web UI 和 REST API 通过 `?token=...` 或 `Authorization: Bearer ...` 传递。

可通过环境变量 `PULSEMQ_ADMIN_TOKEN` 或配置文件覆盖。

---

## 日志

日志输出到 `logs/` 目录，每日滚动，保留 30 天：

```
logs/
├── pulsemq_2026-06-27.log
├── pulsemq_2026-06-28.log
└── ...
```

stderr 同步输出（容器/交互可见）。

---

## 性能基准

### 单条消息（dict，每帧 1 条）

| 序列化 | 压缩 | 吞吐量 | P50 延迟 |
|--------|------|--------|----------|
| msgpack | none | **18.9 K/s** | 183.9 ms |
| msgpack | lz4 | 18.4 K/s | 124.1 ms |
| json | none | 18.6 K/s | 151.4 ms |
| pyarrow | none | 2.1 K/s | 781.6 ms |

> pyarrow 适用于批处理，单条 dict 场景序列化开销过大（8-9x 慢于 msgpack）。
> 压缩对< 200B payload 为负收益。机器：Windows 11, Python 3.14, 单机 localhost。

### 批量 DataFrame（1000 行/帧）

| 序列化 | 压缩 | 记录/s | P50 延迟 |
|--------|------|--------|----------|
| **pyarrow** | **zstd** | **363.6 K/s** | 1677 ms |
| pyarrow | lz4 | 341.1 K/s | 1481 ms |
| msgpack | none | 292.5 K/s | 1955 ms |
| json | none | 290.5 K/s | 1983 ms |

> 批量场景 pyarrow + 压缩是**最优组合**（直接序列化 DataFrame，无需格式转换）。
> 延迟为压测下 ZMQ 缓冲区排队所致，真实场景以固定间隔发送时远低于此。

### 基准脚本

```bash
# 端到端压测（单进程 consumer + producer）
python bench_v2_e2e.py --duration 10

# 分离进程行情压测
python bench_v2_market.py producer --duration 10
python bench_v2_market.py consumer --duration 13

# 全矩阵 ser × comp × data_type 类型保真验证
python bench_v2_matrix.py

# 行情全矩阵性能（12 组合）
python bench_market_full.py --duration 5

# DataFrame 批量性能（1000 行/帧）
python bench_df_batch.py --duration 5
```

---

## 配置

### 环境变量

| 变量 | 说明 | 默认 |
|------|------|------|
| `PULSEMQ_DATA_ENDPOINT` | 数据面绑定地址 | `tcp://0.0.0.0:5555` |
| `PULSEMQ_CONTROL_ENDPOINT` | 控制面绑定地址 | `tcp://0.0.0.0:5556` |
| `PULSEMQ_ADMIN_ENDPOINT` | 管理 HTTP 绑定地址 | `0.0.0.0:9090` |
| `PULSEMQ_CREDENTIALS_FILE` | 凭据 TOML 路径 | `./pulsemq_users.toml` |
| `PULSEMQ_ADMIN_TOKEN` | 监控接口 token | 自动生成 |
| `PULSEMQ_ADMIN_TOKEN_FILE` | token 文件路径 | `./pulsemq_admin.token` |
| `PULSEMQ_ADMIN_PASSWORD` | 首次启动默认密码 | 随机 16 字符 |
| `PULSEMQ_STATS_DB` | SQLite 路径 | `./pulsemq_stats.sqlite` |
| `PULSEMQ_HEARTBEAT_TIMEOUT` | 心跳超时（秒） | 6.0 |
| `PULSEMQ_LATENCY_SAMPLE_RATE` | 延迟采样率 | 0.01 (1%) |
| `PULSEMQ_STATS_RETENTION` | 内存窗口（分钟） | 480 (8h) |
| `PULSEMQ_BCRYPT_COST` | bcrypt 代价因子 | 12 |

---

## 更新日志

### v2 (current)

完整重构：PUB/SUB → Client/Server (ROUTER/DEALER)

- **架构变更**：单 PUB socket → 双 ROUTER（数据面 + 控制面）+ HTTP admin
- **认证升级**：api_key 明文 → bcrypt CredentialStore + ZAP PLAIN + 用户 CLI
- **类型保真**：`_restore_type` 确保 DataFrame/dict str/bytes 端到端还原
- **监控增强**：延迟 P50/P95/P99、在线客户端、事件流、SSE 实时推送、独立 Admin 线程
- **性能优化**：`decode_header` 服务端零反序列化路由、`TrafficStats` 单 `dict.get`
- **自动重连**：运行期指数退避重连（1s→2s→4s→...→30s）
- **日志系统**：loguru 统一，每日滚动写入 `logs/`，30 天保留
- **安全性**：密码 bcrypt 哈希、admin token 随机生成（0600）、凭据文件原子写入
- **批次处理**：DataFrame 批量 1000 行/帧 达到 363 K records/s（pyarrow+zstd）

---

## 许可证

[MIT](LICENSE)
