NNG协议
NNG 最吸引人的地方,是几行 C 就能把线程、进程或机器连起来。
本文固定在 NNG 1.7.2 讨论,源码对应 tag v1.7.2、commit daa55213d1bc7367f9eff29b010b563c71fb42cf。文中的程序在 Ubuntu 24.04、GCC 13.3.0 上实际编译运行,使用 -Wall -Wextra -Werror,没有拿新版本手册去解释旧版本二进制。
附件放在文末,也可以直接下载:
先分清六个名词
我现在看 NNG 程序,会先问六个问题。
| 对象 | 它回答的问题 | 例子 |
|---|---|---|
| protocol | 消息按什么规则流动 | REQ/REP、PUB/SUB、PUSH/PULL |
| transport | 字节借什么通道移动 | inproc://、ipc://、tcp:// |
| socket | 哪个协议实例共享端点与连接 | 一个 cooked REQ socket |
| dialer/listener | 谁发起、谁接受连接 | dial / listen |
| pipe | 当前已经建立的哪一条连接 | 某个 TCP 会话对应的 pipe |
| context | 哪一份独立协议状态 | 某个 REQ 的 request ID、重发计时器 |
NNG 的 协议与 transport 索引把前两层分得很清楚:协议定义通信模式,transport 负责搬运消息。
一个 socket 在创建时就固定一个协议,却可以同时拥有多个 dialer、listener 和 pipe。nng_ctx 不会再建 TCP 连接,它共享 socket 的端点和 pipe,只另外保存协议状态。nng_ctx(5)列出的典型状态正是 REQ 的 request ID 与 retry timer。
context != thread:context 是协议状态容器,线程、事件循环或 AIO callback 才是应用的执行方式。pipe != socket:pipe 对应 socket 上的一条已建立连接;断线后旧 pipe 消失,重连会产生新 pipe。nng_pipe(5)给出了对象定义。
dial/listen 与 REQ/REP 没有主从绑定
REQ 常被写成 dialer,REP 常被写成 listener,只是部署习惯。
连接方向与协议方向是两套关系。REP 完全可以主动连接正在监听的 REQ;连接建好后,仍然必须由 REQ 先发请求,REP 收到后才能回复。nng_dialer(5)也明确说 endpoint 的 client/server 关系与协议角色正交。
附件 nng_roles.c 故意反着部署:
1 | |
编译后可以开两个终端,也可以用附件脚本自动跑:
1 | |
我的实际输出如下:
1 | |
NNG_FLAG_ALLOC 返回的是指定长度的二进制消息,不承诺末尾带 NUL。所以例子用 %.*s 按长度打印,并用 nng_free(ptr, size) 配对释放。
cooked REQ/REP 到底替应用做了什么
nng_req0_open() 与 nng_rep0_open() 默认打开 cooked socket。
cooked 模式不只是给消息加一个类型。它维护一套有限状态机:
1 | |
REQ 在没有 outstanding request 时接收,会得到 NNG_ESTATE。REP 在没有先接到请求时发送,也会得到 NNG_ESTATE。同一个 context 也不能同时挂两次接收。REQ 手册与 REP 手册把这些顺序当成协议契约。
发送第二个 REQ 会取消第一个请求在本地的等待状态,旧回复随后会被丢弃。已经到达远端的旧请求不会因此停止执行。
这条边界很重要:协议能取消“我还等不等”,不能撤销“对方已经做了什么”。
raw socket 则把路由头与状态责任交给应用,适合写 device/forwarder。raw 模式不支持 context;为了绕开 NNG_ESTATE 就改 raw,通常只会把匹配、重试、TTL 与坏帧处理一起揽到自己身上。
request ID 与 backtrace
REQ/REP v0 的协议头是一组 32 位大端整数。
最后一项是 request ID,最高位必须为 1。它前面可以有若干 peer ID,最高位为 0。raw forwarder 收到请求时把入站 peer ID 压到前面,回复回来再弹出,以此选择下一跳。
例如下面只是一个便于阅读的示意值:
1 | |
应用在 cooked REP 收到的 body 里看不到这些字节。库把 backtrace 留在内部 context,发送回复时再补回去。REQ 协议头说明给出了同样的栈模型。
它解决的是“回复属于哪个请求、沿哪条路返回”,不是业务去重 ID。
request ID 在一次运行中用于协议匹配;订单号、扣款流水号这类幂等键必须由应用放进 payload,并在服务端保存处理结果。
重发带来的语义:至少可能执行一次
REQ 默认会保留请求副本,直到收到匹配回复、被新请求替换、超时/取消,或 context 关闭。
NNG_OPT_REQ_RESENDTIME 到期会再次发送请求;peer 断开或等待期间出现可用 peer,也可能触发再次发送。NNG_OPT_REQ_RESENDTICK 是 socket 共享的扫描粒度,而 resend time 可以按 context 配置。REQ 手册的 options 部分对此有完整说明。
考虑这条时间线:
1 | |
如果服务端只看 order=42 然后直接扣款,协议重试就可能变成重复副作用。
比较稳妥的处理是给操作分配稳定的 op_id,在同一个事务中记录 op_id -> result。重复请求命中记录时,返回原结果,不再执行副作用。
接收超时也不能证明远端没有执行。它只证明本地在截止时间前没有拿到匹配回复。
context 才是 cooked REQ/REP 的并发单位
socket 自带一个默认 context,所以直接对 REQ socket 调 nng_sendmsg/nng_recvmsg 时,一次只能有一个 outstanding request。
要并发,不需要为每个请求重新建 socket 和 TCP 连接。可以在同一个 socket 上开多个 context:
1 | |
每个 REQ context 独立保存 request ID、待收 reply 和 resend time;它们共享 socket 的 pipe 池。每个 context 内仍然是一问一答,context 之间可以并行。
REP 同理。一个 REP context 在收到请求后保存 backtrace 与入站 pipe,之后的 send 才知道回复该回到哪里。开多个 REP context,才能同时持有并处理多个请求。
我的实验使用六个 REQ context、三个 REP context。每个服务端 context 绑定一个 AIO worker;不同请求延迟 10–180 ms,实际观察到三个服务任务同时处于处理阶段。
固定 1.7.2 读一遍源码
下面不是猜调用链,而是对 v1.7.2 的定点走读。内部 nni_* 不是应用 ABI,升级版本后应重新核对。
REQ 的 context 与 socket 各保存什么
req.c:32-67里,req0_ctx 有这些关键字段:
1 | |
socket 侧则有 ready_pipes、busy_pipes、send_queue、retry_queue 和 requests ID map。也就是说,context 保存每笔请求的状态,socket 负责在共享 pipe 间调度。
req0_sock_init把 ID map 范围初始化为 0x80000000..0xffffffff,随机选择起点;同时把默认 retry 设为 60 秒、retry tick 设为 1 秒。
这也能从代码层确认 request ID 的最高位为 1。
一次 cooked send 如何进入队列
req0_ctx_send先取消同 context 的旧 send/recv,再重置旧状态,分配新的 request ID,并把 ID 写进 NNG 消息 header。
随后它做三件事:
ctx->req_msg = msg,协议接管原消息;- 把 context 放入 retry queue;
- 把 context 放入 send queue,调用
req0_run_send_queue()。
req0_run_send_queue取一个 ready pipe,把用户的 send AIO 完成为成功,然后 clone req_msg 交给 pipe。
原消息仍留在 context 中,正是为了可能的重发。这里的 send 成功表示协议已经接受请求,不表示 REP 已经收到,更不表示业务已完成。
reply 如何回到正确 context
req0_recv_cb从消息前面取出 4 字节 ID,在 requests map 中查 context。
找不到 ID、请求还没真正发出,或者该 context 已经有 reply 时,消息会被释放。匹配成功后,它移除 ID map 项和重发节点,释放保留的 request,再完成等待中的 receive AIO;如果应用还没调用 receive,就暂存在 rep_msg。
这就是发送新请求后,旧回复不会串给新请求的原因。
重发扫描做了什么
req0_retry_cb在 socket 锁内扫描 retry queue。
到期且仍有 req_msg 的 context 会重新进入 send queue。只有 retry queue 非空时才继续启动 timer,避免空闲 socket 一直唤醒。
pipe 断开时,req0_pipe_close把仍有请求副本的 context 放回 send queue;若禁用 retry,则等待中的 receive 以连接重置结束。
REP 为什么必须先 receive
REP context 在 rep.c:28-38保存 pipe_id、btrace、send AIO 与 receive AIO。
rep0_pipe_recv_cb每次从 body 前面取 4 字节,搬进内部 header,直到看到最高位为 1 的 request ID。超过 TTL 的帧被丢弃,连 4 字节都不够的坏帧会让 pipe 关闭。
它把完整 backtrace 和 pipe ID 保存进拿到该请求的 context,清掉应用可见 header,再完成 receive。
rep0_ctx_send发现 btrace_len == 0 时直接返回 NNG_ESTATE。正常路径会把 backtrace 放回 header,并只向原 pipe 发送。
还有一个容易忽略的细节:若回复时原 pipe 已消失,1.7.2 会释放消息并把 send 完成为成功。这是有意保持 REP 状态机前进,再次说明 send 成功不能作为送达凭据。
消息所有权:只看完成结果
同步 nng_sendmsg() / nng_ctx_sendmsg() 的规则很简单:
- 返回 0,socket 接管
msg,应用不能再读、改或 free; - 返回非 0,应用仍拥有
msg,可以重试或 free。
nng_ctx_sendmsg(3)还说明,成功只代表 context 接受了消息,队列和物理链路上仍可能有延迟或丢失。
异步发送要等 callback 才能分支。
安全的回调骨架是:
1 | |
nng_send_aio(3)的措辞很精确:成功时是 socket accepted/queued;失败、取消或超时时,callback 必须从 AIO 取回消息并处置。
提交 nng_ctx_send(ctx, aio) 之后,不能立即 free msg,也不能立即拿同一个 AIO 发起另一个操作。提交函数本身没有返回异步结果。
一个完整的 REP AIO worker
实验中的 worker 有四个显式状态:
1 | |
processing 与 aio 上的 message 分开记录,是为了让取消路径知道消息究竟停在哪一段。
核心 callback 如下,完整错误检查和实验驱动见附件:
1 | |
这里用 nng_sleep_aio() 模拟业务处理,让同一个 AIO 依次驱动 receive、timer 与 send。生产代码也可以把 CPU 工作交给线程池,再把结果送回事件状态机;关键是同一 worker 不同时复用一个 busy AIO。
每个 REP context 只提交一个 receive。三个 worker 就有三个独立 receive,可以同时拿走三条请求,各自保存自己的 backtrace。
stop、cancel 与 free 的顺序
关停最危险的写法,是先 free worker,再 close socket,希望 pending callback 自己消失。callback 很可能仍持有 worker *,这就是 use-after-free。
实验使用的顺序是:
1 | |
先设置 stopping,callback 就不会再次提交 receive。随后 stop 全部 AIO,最后才逐个释放 context、AIO 和 owner。
nng_aio_stop(3)会以 NNG_ECANCELED 中止操作,并等待操作及 callback 完全结束;它还会阻止这个 AIO 再次 begin。官方特别建议多个 AIO 先全部 stop,再 free 其中任何一个。
不要在该 AIO 自己的 callback 里调用阻塞式 nng_aio_free()。确实要从 callback 丢弃它时,接口提供了后台回收的 nng_aio_reap();更容易审计的做法仍是由 owner 统一停机。
这套顺序没有在持有应用 mutex 时调用 nng_aio_stop()。否则主线程拿着锁等 callback 退出,而 callback 又等同一把锁,会形成死锁。
实测:并发、非法状态、超时与所有权
运行环境与命令:
1 | |
脚本编译 nng_aio_context_lab.c、nng_flow_lab.c 和 nng_roles.c。主实验只绑定 tcp://127.0.0.1:18892,其余检查使用 inproc://。
开头两条是专门设计的错误路径:
1 | |
第一条在 REP context 尚未 receive 时异步 send,结果为 NNG_ESTATE;失败消息仍挂在 AIO 上,应用取回并 free。
第二条让 REP 收到请求后故意不回复,REQ 的 NNG_OPT_RECVTIMEO=120ms 最终得到 NNG_ETIMEDOUT。
六个 context 的关键输出是:
1 | |
短请求没有被 180 ms 的第一个请求挡住,max concurrent server jobs: 3 也由程序断言。每次成功异步 send 后 msg_on_aio=no,符合所有权转移规则。
关闭时三个 pending receive 都收到取消:
1 | |
这不是“忽略错误”,而是 stop 路径的预期完成结果。
PUSH 与 PUB 的“发送成功”完全不同
PUSH/PULL 和 PUB/SUB 都能一对多连接,但发送策略相反。
| 情况 | PUSH | PUB |
|---|---|---|
| 分发对象 | 一个可接收的 PULL | 每条现存 pipe 各一份 |
| 无可用 peer | 等待/排队,最终可能超时 | 直接成功并丢弃 |
| 慢 peer | 产生背压,不选暂不可接收的 pipe | 每 pipe 队列满时丢旧消息 |
| 应用级确认 | 没有 | 没有 |
| 典型用途 | 工作分配 | 瞬时状态/事件广播 |
nng_push(7)规定 PUSH 在可用 PULL 中轮转,并遵守 flow control。没有 peer 能接收时,send 会等到出现可用 peer或超时。
NNG_OPT_SENDBUF 默认是 0,表示无 socket 中间缓冲;设为正数后,容量单位是 消息条数,范围 0–8192。transport 还可能另有缓冲,所以 send 入队成功依然不是远端确认。
1.7.2 的 push0_sock_send依次尝试 ready pipe、socket queue,最后才把用户 AIO 挂进等待队列。pipe 再次 ready 时,push0_pipe_ready优先排空缓冲消息。
我的测试把 SENDBUF=0、SENDTIMEO=80ms,又不创建 PULL:
1 | |
PUSH 尽量不丢并提供背压,但协议没有 ack。消息被 socket 接受后 pipe 随即损坏,仍可能丢失;需要端到端确认时应设计 REQ/REP 或应用层 receipt。
PUB 则是 best effort。1.7.2 的 pub0_sock_send遍历当前 pipes,为每条 pipe clone 消息;pipe 忙且队列满时,先移除最旧消息,再放入新消息。若一个 pipe 都没有,循环为空,原消息被 free,用户 send 仍以 0 完成。
所以这段程序没有竞态解释空间:它先发布一个 能匹配订阅前缀 的消息,再建立 SUB pipe,接收必然超时。
1 | |
这类启动阶段丢失不是把 sleep 调大就能成为协议保证。若第一条状态不能丢,需要应用握手、快照接口,或者选带持久会话语义的系统。
SUB 过滤的是原始字节前缀
SUB topic 不是 MQTT 风格的层级主题,也没有 +、# 通配符。
它是消息 body 开头的一段任意字节。多个订阅之间是 OR;空前缀匹配所有消息。nng_sub(7)还指出,PUB 会把消息发给所有 subscriber,过滤发生在 SUB 本地,因此订阅不能节省链路带宽。
二进制前缀必须用长度明确的接口:
1 | |
实验发送两条三字节消息:
1 | |
实际输出:
1 | |
源码 sub0_matches就是长度检查加 memcmp(topic, body, topic_len)。
字符串接口有一个坑:1.7.2 手册说明 nng_socket_set_string/旧名 nng_setopt_string 会把结尾 NUL 算进 topic。订阅字符串 "sensor/" 实际前缀可能是:
1 | |
而 body "sensor/temp" 的第八字节是 t,自然不匹配。想订阅可见字符前缀,应使用通用 byte-array setter 和 strlen(topic)。
SUB 自己也有队列丢弃策略。NNG_OPT_SUB_PREFNEW=true 是默认值,队列满时移除最旧消息以容纳新消息;设为 false 则拒绝新消息。对应实现可见 sub0_recv_cb。
取消一个订阅后,1.7.2 还会重新扫描该 context 已排队的消息,清掉不再匹配剩余 topic 的项;见 sub0_ctx_unsubscribe。
队列与背压要画到具体一层
“NNG 有队列”信息量很低。至少要说明是哪一层:
1 | |
每层的“成功”只说明它把责任交给了下一层。
队列深度通常按消息数计算,不按总字节数。若消息从 100 B 到 5 MB 都可能出现,SENDBUF=100 不能直接推导出内存上限。生产环境还要限制单消息大小 NNG_OPT_RECVMAXSZ,并把应用队列计入预算。
PUSH 的背压能把压力传回发送调用,但 PUB 为了不让慢订阅者拖住发布者,会选择丢消息。二者没有一个统一的“可靠等级”数字,必须从协议路径判断。
四类时间不能混成一个 timeout
1. send timeout
NNG_OPT_SENDTIMEO 限制本地 send 等待 socket/context 接受消息的时间。它不等待远端应用处理。
对 SENDBUF=0 的 PUSH,它常表现为等待可用 PULL;对 best-effort PUB,send 本来就不会等待订阅者确认。
2. receive timeout
NNG_OPT_RECVTIMEO 限制一次本地 receive。REQ receive 超时会结束本次等待;下一次业务尝试应重新 send,而不是在已经被取消的状态上无限 receive。
源码 req0_ctx_cancel_recv明确把取消 pending receive 当成终止整个 context 状态机,并清理 request。
3. REQ resend time
NNG_OPT_REQ_RESENDTIME 是协议内部在等待同一逻辑 reply 时的重发周期。它不是调用方的总 deadline。
总 deadline 为 2 秒、resend 为 500 ms 时,一个调用可能在截止前发出多份请求。业务必须按同一幂等键识别它们。
4. dialer reconnect time
NNG_OPT_RECONNMINT 与 NNG_OPT_RECONNMAXT 控制失败连接的重试间隔。max 非 0 时,间隔从 min 指数增长直到 max;max 为 0 时固定使用 min。nng_options(5)给出了这两个选项的范围。
默认同步 nng_dial() 的首次连接若立即失败,会把错误返回给调用者,不自动安排首次失败后的重试。使用 NNG_FLAG_NONBLOCK 时,首次拨号异步进行,失败会在后台重试。
一旦 dialer 曾经连通,pipe 后来断开,则即使最初是同步 dial,dialer 也会异步重连。nng_dialer_start(3)区分了这两个阶段。
因此,启动成功、发送被接受、协议重发和业务 deadline 应分别记录。只打一条“timeout”日志,很难知道是哪层计时器触发。
排查时留下能关联的证据
每笔业务请求至少记录稳定的 op_id、context/worker、endpoint、操作名、消息长度、NNG 数值错误码与 nng_strerror()、耗时和绝对 deadline。能取得 pipe ID 时一并记录,断线重建后就能分辨新旧连接。
看到 ETIMEDOUT,先确认来自 send、receive 还是业务 deadline;看到 send 0,再看协议是 PUSH、PUB 还是 REQ。异步崩溃则检查 callback 是否重提操作、AIO 是否仍挂着 message,以及 owner 是否早于全部 AIO stop 被释放。
复现实验
Ubuntu/Debian 已安装 1.7.2 开发包时:
1 | |
自定义 NNG_PREFIX 需要包含 include/nng 与 lib;可用 NNG_BUILD_DIR="$PWD/build" 改变产物目录。脚本也会把前缀的 lib 目录加入 LD_LIBRARY_PATH,供独立包树中的 NNG 与 TLS 依赖加载。
端口 18891 用于角色反转例子,18892 用于并发例子,都只监听 127.0.0.1。若端口被占用,修改源码中的常量后重新编译即可。
