ITADN
paddor/omq.rs
paddor/omq.rs · 文件 下载 ZIP
文件最后提交记录最后更新时间
README.md
以下内容由 AI 翻译,如有问题请点此提交 issue 反馈

ØMQ.rs

纯 Rust ZeroMQ:用于分布式和并发应用的无代理消息传递。在进程内、进程间以及网络上表现一致的套接字级消息传递模式。

  • 适用于 Linux、macOS 和 Windows 的 Tokio 后端
  • 20 种套接字类型:稳定的 ZMQ 模式以及草案 CLIENT/SERVER、RADIO/DISH、SCATTER/GATHER、CHANNEL/PEER 和 STREAM
  • 8 种稳定传输协议:TCP、IPC、inproc、UDP、WS、WSS、lz4+tcp://、 以及 lz4+ws://;实验性的 zstd+tcp:// 压缩传输协议位于 zstd 特性之后
  • 3 种安全机制:NULL、PLAIN、CURVE
  • 无需 C 编译器,无需 libzmq,无需 libsodium
  • Python 绑定(pyomq),C API(omq-libzmq

PUSH/PULL throughput: TCP implementations

REQ/REP 延迟

REQ/REP latency: TCP implementations

PUB/SUB 吞吐量

PUB/SUB throughput: TCP implementations

LZ4 PUSH/PULL 吞吐量

LZ4 PUSH/PULL throughput over TCP

完整对比图表

难点

OMQ 旨在实现真实的 ZMQ 行为,而不仅仅是理想路径下的 PUSH/PULL 吞吐量。您将获得:

  • 无需额外调优的 ZeroMQ 语义:没有拓扑特定的 socket 类型,没有用户可见的批处理 API,没有手动重连循环。
  • 传输故障是常态:重连、先连接后绑定、对端变动以及绑定端重启都是设计的一部分。
  • 对端故障不会变成用户错误:send()recv() 在断开连接、重连、慢消费者以及绑定端重启期间保持正常工作。
  • 在负载下的 HWM 背压和路由公平性,而不仅仅是在空队列示例中。
  • 关于无对端发送、linger 和 HWM 的已记录 libzmq 兼容性边缘情况: doc/libzmq/semantics.md.
  • 热路径具有大小感知和低延迟意识:小消息保持内联且无分配,inproc 按值传递消息,大负载在关键处使用零拷贝缓冲区。
  • 唯一遵循 libzmq 架构的 Rust ZeroMQ 实现:应用线程与专用的后台 IO 线程分离,IO 工作在这些线程上线性扩展,PUB 对端自动分配给 IO 通道。
  • 公共 crate 使用内存安全的 Rust。unsafe 被隔离并使用 Miri 进行检查。
  • 基准测试覆盖真实形态:CPU 核算、扇入/扇出、公平性、传输差异。

用法

[!NOTE] API 仍在演变,可能在次要版本之间发生变化。欢迎提交 bug 报告并在真实工作负载中进行测试。

Rust 后端是 omq-tokio:在 Linux、 macOS 和 Windows 上使用 tokio + mio。它可以在 OMQ 拥有的运行时线程上 或现有的 tokio 运行时内运行套接字 IO。

API运行时位置扩展模型
Context::new().socket(...)异步套接字,OMQ 拥有的 IO 线程线性 IO 通道扩展;PUB 对等节点跨通道分片
Context::current().socket(...)异步套接字,调用者的活动 tokio 运行时Tokio 调度器/工作窃取;PUB 扇出保留在一个 OMQ 通道上
Context::new().blocking_socket(...)同步套接字,OMQ 拥有的 IO 线程线性 IO 通道扩展;调用者线程不介入 IO

支持的 32 位 Linux 目标是 i686-unknown-linux-gnuarmv7-unknown-linux-gnueabihf。它们需要原生 64 位原子操作。ZMTP 线 长度字段保持 64 位,但实际的帧/消息大小受 平台分配限制约束(在 32 位上低于 4 GiB)。

如果你了解 ZeroMQ,你就了解 OMQ。相同的套接字类型,相同的 connect/bind/send/recv:

use omq_tokio::{Context, Message, Options, SocketType};

let ctx = Context::new();

let push = ctx.socket(SocketType::Push, Options::default());
push.connect("tcp://127.0.0.1:5555".parse()?).await?;
push.send(Message::single("hello")).await?;

let pull = ctx.socket(SocketType::Pull, Options::default());
pull.bind("tcp://127.0.0.1:5555".parse()?).await?;
let msg = pull.recv().await?;
assert_eq!(&msg[0], b"hello");

更多示例见 examples/zguide-tokio/, 这是 ZeroMQ Guide 模式到 OMQ 的移植。

Cargo 特性

均为可选。默认构建为最小部署:NULL 机制 + TCP / IPC / inproc / UDP,无需 C 编译器。可启用以下任意一项:

feature添加内容额外依赖
plainPLAIN 用户名/密码认证(RFC 24)-
curveCURVE 加密握手机制(RFC 26)crypto_box, crypto_secretbox
lz4lz4+tcp:// 压缩传输(RFC)lz4rip
zstd实验性 zstd+tcp:// 压缩传输zrip
wsWebSocket(ws://)和 安全 WebSocket(wss://)传输rustls, rustls-native-certs

设计亮点

特性详情
Sans-I/O ZMTP 编解码器 (omq-proto)字节输入 / 事件输出;热路径上无异步,无 trait。镜像 rustls::ConnectionCommon
消息计数 HWMsend_hwm/recv_hwm 计数完整消息,而非字节。发送 HWM 针对每个出站管道/环形缓冲区,因此当存在多个管道或发送槽时,原生缓冲的总消息数可能超过一个 send_hwm
连续帧负载&msg[0] 直接提供 &[u8];无可能失败的借用,无合并步骤。
零拷贝发送和接收发送:大型 Bytes 负载无需任何数据拷贝即可到达内核 writev。接收:大型帧直接读取到预分配的缓冲区中,绕过中间队列。
Patricia 树订阅匹配器复杂度为 O(M)(M 为主题长度),而非 O(NxM)。
LZ4 字典自动训练默认关闭。启用后,从前 100 条消息中进行训练,并向对端发送一次;将有效压缩阈值从 512 B 降低至 64 B。
监控事件类套接字的 Stream,在每个连接 / 断开连接 / 握手事件上拥有 PeerInfo

工作区

五个 Cargo 工作区 crate 加上 Python 绑定。

Crate功能不安全策略
omq-proto无 I/O 的 ZMTP 3.x 核心:编解码器、消息、机制、订阅#![forbid(unsafe_code)]
omq-tokio多线程 tokio 后端(Linux/macOS/Windows)#![forbid(unsafe_code)]
omq-libzmq兼容 libzmq 的 C 接口(libomq_zmq 动态/静态库)不安全的 C ABI 边界
yring有界 SPSC 环形缓冲区,采用 ypipe 风格的批量刷新/预取不安全的环形核心,经 Miri 测试
omq-bench基准测试运行器和 SVG 图表生成器仅限基准测试的进程控制和 CPU 核算
pyomqPython 绑定(基于 omq-tokio 的 PyO3,同步 + asyncio)PyO3 FFI 边界

测试

每种套接字类型、传输、机制和功能组合均 由集成测试覆盖。测试套件是分层的:

  • 700+ Rust tests 涵盖 socket 类型、传输层、机制以及 libzmq 兼容的 C API 行为。
  • Feature-gated coverage 针对 PLAIN、CURVE、LZ4 以及 pyzmq/libzmq 互操作性。WebSocket 拥有专门的测试和 soak 覆盖。
  • Protocol fuzzing(默认 opt-in 运行约 100 万次迭代, 可配置更长的运行时间):针对 wire 解析器和 socket-action 状态机的手工编写模糊测试。
  • 20+ soak scenarios 涵盖 Rust 和 pyomq:peer churn、重连 风暴、PUB/SUB churn、ROUTER/DEALER churn、HWM 重连、取消 安全性、压缩(lz4)、PLAIN / CURVE 认证、机制重连、 大消息吞吐量、多 socket、inproc 跨线程、 WebSocket 吞吐量和重连。Soak 运行会采样 RSS 和 FD 计数。
  • Loom 覆盖 yring SPSC 内存排序、异步唤醒以及 omq-tokio 信号竞争窗口。
  • Miri 针对 yring
  • Release semver review 通过 release-plz
./scripts/test-all.sh              # standard sweep with local perf gate
OMQ_FUZZ=1 ./scripts/test-all.sh   # include fuzz suites
OMQ_SKIP_PYOMQ=1 ./scripts/test-all.sh
OMQ_SKIP_PERF=1 ./scripts/test-all.sh

浸泡测试有意与完整扫描分离:

FEATURES="soak lz4 plain curve ws"
OMQ_SOAK_DURATION_SECS=600 cargo test -p omq-tokio \
  --features "$FEATURES" --release --test omq_soak_peer_churn -- --nocapture

延伸阅读

平台与要求

Linux 是主要的开发和基准测试平台。 PR CI 必需检查涵盖 Linux x86_64、macOS ARM64 和 Windows Rust 测试, 以及 32 位 Linux 交叉检查。Fmt 和 clippy 也在 macOS Intel 上运行。 扩展 CI 涵盖 Ubuntu ARM64。macOS 作业串行运行 Rust 测试, 因为在托管运行器上 socket/定时器时序更为敏感。

macOS 在 CI 中同时覆盖 Intel 和 ARM64 运行器。 omq-tokio 使用 mio / kqueue。omq-libzmq 使用基于管道的 通知 fd 用于 zmq_poll/ZMQ_FD 就绪状态。

Windows 在 CI 中覆盖。omq-tokio 支持 TCP、IPC 命名 管道、inproc、UDP 和 WebSocket 传输。omq-libzmq 在 Windows 上构建并 测试受支持的 C API 表面。

pyomq 发布 Linux、macOS 和 Windows 的 wheels 以及 sdist。

要求:

  • Rust 1.93 或更高版本(edition 2024)。

贡献

请参阅 CONTRIBUTING.md 获取指南,以及 DEVELOPMENT.md 获取构建、测试和基准测试命令。

AI 披露

本项目在架构、实现、测试、基准测试基础设施和文档方面均获得了大量 LLM 辅助。它是一次关于 LLM 辅助开发能做什么和不能做什么的实验。设计决策和方向由我负责。

License

ISC.