binbin第一次直接assign给我的issue,我来尝试处理一下并且学习一下这里的代码。
原issue界面:https://github.com/valkey-io/valkey/issues/4231
1. 这个 bug 发生在什么业务场景 #
一条 TCP 连接同时干两件事:
SUBSCRIBE wiretest—— 成为 channel 的订阅者PUBLISH wiretest <大 payload>—— 又是发布者
Valkey 会给每一个订阅者投递一份消息,其中也包括发布者自己。
SignalR Redis backplane、某些 multiplexer 客户端就是「同连接 pub+sub」,所以会踩到。
另一条连接只订阅、不发布:投递是正常的。坏的只有「发布方自己那份」。
2. 协议层:客户端看到的字节流是什么 #
Valkey 用 RESP。客户端从 socket read() 出来的是一串类型标记 + 负载,必须按顺序解析,不能乱。
RESP2 的 pubsub 投递是普通 array:
*3\r\n
$7\r\nmessage\r\n
$8\r\nwiretest\r\n
$16384\r\nXXXX...\r\n
RESP3 用 push(>)表示「服务器主动推送,不是某条命令的回复」:
>3\r\n
$7\r\nmessage\r\n
$8\r\nwiretest\r\n
$16384\r\nXXXX...\r\n
同连接上再发 PUBLISH,命令本身还有一个整数回复(订阅者个数):
:1\r\n
客户端(如 StackExchange.Redis)会维护「这条命令期待什么 reply」。它期望:
:1\r\n ← PUBLISH 的回复
>3\r\n$message\r\n$channel\r\n$payload\r\n ← 推送
若先看到裸的 $16384\r\n...,就会当成 PUBLISH 的 reply → 报错并和后续字节流失步。
3. 服务器侧:一条连接的「待写出队列」长什么样 #
每个 client 大致有三块和写出相关的结构(可跳到 src/server.h 里 client 定义):
| 结构 | 作用 |
|---|---|
c->buf + c->bufpos |
小 reply 先堆在这块静态 buffer(约 16KB,PROTO_REPLY_CHUNK_BYTES) |
c->reply |
list,装不下就 spill 成一块块 clientReplyBlock |
server.pending_push_messages |
全局临时 list:当前命令执行期间,给「自己」的 push 先放这里 |
写出时(write / writev)顺序大致是:先 c->buf,再按顺序扫 c->reply list。
谁先进入 buf/reply,谁先上 socket。
事件循环里典型路径:
read(query) → 解析命令 → 执行 → 把回复塞进 c->buf / c->reply
→ 事件循环稍后 writev(fd, iov...) 把缓冲刷到内核
4. Pub/Sub 投递:从 PUBLISH 到字节 #
命令入口在 src/pubsub.c:
void publishCommand(client *c) {
if (server.sentinel_mode) {
sentinelPublishCommand(c);
return;
}
int receivers = pubsubPublishMessageAndPropagateToCluster(c->argv[1], c->argv[2], 0);
if (!server.cluster_enabled) forceCommandPropagation(c, PROPAGATE_REPL);
addReplyLongLong(c, receivers);
}
注意顺序:
- 先
pubsubPublishMessage...—— 给所有订阅者排队 push - 再
addReplyLongLong(c, receivers)—— 给当前客户端写:N
投递核心:
int pubsubPublishMessageInternal(robj *channel, robj *message, pubsubtype type) {
...
while (hashtableNext(&iter, &c)) {
addReplyPubsubMessage(c, channel, message, *type.messageBulk);
...
receivers++;
}
组装一帧 push:
void addReplyPubsubMessage(client *c, robj *channel, robj *msg, robj *message_bulk) {
struct ClientFlags old_flags = c->flag;
c->flag.pushing = 1;
if (c->resp == 2)
addReply(c, shared.mbulkhdr[3]);
else
addReplyPushLen(c, 3);
addReply(c, message_bulk);
addReplyBulk(c, channel);
if (msg) addReplyBulk(c, msg);
if (!old_flags.pushing) c->flag.pushing = 0;
}
语义上连续写 4 段:
>3\r\n(或 RESP2 的*3)message- channel 名(bulk)
- payload(bulk)
中间设了 c->flag.pushing = 1,告诉下层 reply API:「这是 push,不是普通命令回复」。
5. 为什么要有 pending_push_messages(defer)
#
若发布者就是当前连接,且在执行 PUBLISH 中途就把 push 写进 c->buf,会发生:
[push 半帧或整帧] [ :1 ] [或许更多]
更糟的是 MULTI/EXEC:一串命令回复中间插入 push,客户端协议状态机很难处理。
所以 Valkey 约定:
给正在执行命令的这个 client 投递 push 时,先不进正式 reply 队列,放进
pending_push_messages;命令跑完、命令回复写好之后,再listJoin拼到c->reply末尾。
逻辑在 src/networking.c 的 _addReplyToBufferOrList:
/* If we're processing a push message into the current client (i.e. executing PUBLISH
* to a channel which we are subscribed to, then we wanna postpone that message ... */
int defer_push_message = c->flag.pushing && c == server.current_client && server.executing_client &&
!cmdHasPushAsReply(server.executing_client->cmd);
...
if (defer_push_message) {
_addReplyProtoToList(c, server.pending_push_messages, s, len);
return;
}
条件拆开:
c->flag.pushing:正在组 push(addReplyPubsubMessage设的)c == server.current_client:接收 push 的人就是发命令的人(self-publish)- 当前命令不是 SUBSCRIBE 族(那些命令本身用 push 当「回复」)
命令结束后:
if (!server.execution_nesting) listJoin(c->reply, server.pending_push_messages);
理想时序:
执行 PUBLISH
├─ 给自己的 push 各段 → pending_push_messages
└─ :1 → c->buf / c->reply
afterCommand
└─ pending 接到 c->reply 尾部
写出顺序: :1 → 完整 >3 push
投递给别人时 c != current_client,不 defer,整帧直接进对方 buf/reply,顺序天然正确。
6. Reply Copy Avoidance:9.0 引入的优化 #
大字符串 bulk 回复如果 memcpy 进 reply buffer,CPU/内存带宽浪费大。
#2078 做了「copy avoidance」:
- reply buffer 里不存字符串本体,只存一个小结构
bulkStrRef { robj *obj; char *str; } - 真正
writev时,把 iovec 直接指到对象内存里的 sds
写出展开(可跳到 src/networking.c):
static void addBulkStringToReplyIOV(char *buf, size_t buf_len, replyIOV *reply, bufWriteMetadata *metadata) {
bulkStrRef *str_ref = (bulkStrRef *)buf;
while (buf_len > 0 && !reply->limit_reached) {
size_t str_len = sdslen(str_ref->str);
/* RESP encodes bulk strings as $<length>\r\n<data>\r\n */
...
addPlainBufferToReplyIOV(... prefix "$len\r\n" ...);
addPlainBufferToReplyIOV(str_ref->str, str_len, ...); /* 直接指对象内存 */
addPlainBufferToReplyIOV(reply->crlf, 2, ...);
入口在 addReplyBulk:
void addReplyBulk(client *c, robj *obj) {
if (tryAvoidBulkStrCopyToReply(c, obj) == C_OK) {
...
return; /* 走了 copy-avoid,不再走下面的普通路径 */
}
addReplyBulkLen(c, obj);
addReply(c, obj); /* → _addReplyToBufferOrList(会 defer) */
addReplyProto(c, "\r\n", 2);
}
是否启用由 isCopyAvoidPreferred 决定;单线程默认阈值 16384:
return server.min_string_size_copy_avoid && sdslen(objectGetVal(obj)) >= (size_t)server.min_string_size_copy_avoid;
所以:
- payload 16383:不走 CA →
addReply→_addReplyToBufferOrList→ 会 defer → 正确 - payload ≥16384:走 CA →
_addBulkStrRefToBufferOrList→ 不管 defer → bug
这和「16KB chunk」表象一致,本质是 CA 默认阈值碰巧等于 16KB。
7. 两条路径合在一起:逐步跟一次 bug #
场景:同连接 HELLO 3 + SUBSCRIBE + PUBLISH 16384 字节。
步骤 A:开始组自己的 push #
addReplyPubsubMessage 设 pushing=1,然后:
| 调用 | 走哪条 API | 结果 |
|---|---|---|
addReplyPushLen(>3) |
_addReplyToBufferOrList |
defer → pending |
addReply("message") |
同上 | defer → pending |
addReplyBulk(channel) |
小字符串,CA 不启用 → 普通路径 | defer → pending |
addReplyBulk(msg) 16384 |
CA 启用 | 见下一步 |
此时 pending_push_messages 里大致是:
>3\r\n $7\r\nmessage\r\n $8\r\nwiretest\r\n
还缺最后的 payload。
步骤 B:大 payload 走了错误出口 #
tryAvoidBulkStrCopyToReply → _addBulkStrRefToBufferOrList:
static void _addBulkStrRefToBufferOrList(client *c, robj *obj) {
...
/* 没有 defer_push_message 判断! */
if (!_addBulkStrRefToBuffer(c, ...)) {
_addBulkStrRefToToList(c, ...); /* 硬编码进 c->reply */
}
}
bulkStrRef 进了 c->buf(或 c->reply),不进 pending。
内存视图:
c->buf: [payloadHeader][bulkStrRef{ptr→argv[2]}]
pending: >3 ... message ... channel (缺 payload)
步骤 C:PUBLISH 自己的回复 #
addReplyLongLong(c, 1) → :1\r\n 也进 c->buf,紧挨在 bulkStrRef 后面:
c->buf: [bulkStrRef] [:1\r\n]
pending: >3 ... message ... channel
步骤 D:afterCommand 拼接 #
listJoin(c->reply, pending):
写出顺序:
1) 展开 bulkStrRef → $16384\r\nXXXX...\r\n
2) :1\r\n
3) >3\r\n$message\r\n$channel\r\n ← 没有 payload 的残缺 push
这就是 issue 里的 wire capture。客户端以为 $16384... 是 PUBLISH reply → 协议失步。
8. 你可以用一张图记住 #
addReplyPubsubMessage (pushing=1)
│
┌───────────────────┼───────────────────┐
│ │ │
小片段 (header/ channel payload ≥16KB
message) (小 bulk) │
│ │ │
▼ ▼ ▼
_addReplyToBufferOrList addReplyBulk
│ │
│ defer? ├─ <阈值: 普通路径 → defer ✓
│ (self+pushing) └─ ≥阈值: copy-avoid
▼ │
pending_push_messages ▼
_addBulkStrRef*
→ c->buf / c->reply ✗
│
publishCommand 再写 :1 ──────────────────────────┘
│
afterCommand: listJoin(pending → reply 尾) │
▼
撕碎的 socket 字节流
9. 读代码时建议的顺序(当「导览」) #
publishCommand/addReplyPubsubMessage——src/pubsub.c_addReplyToBufferOrList的 defer 注释 ——src/networking.c~744afterCommand的listJoin——src/server.c~4195addReplyBulk/isCopyAvoidPreferred/_addBulkStrRefToBufferOrList—— 同文件addBulkStringToReplyIOV—— 理解 CA 如何变成$len\r\ndata
本地验证:起一个实例,跑 issue 里的 Python;或 CONFIG GET min-string-size-avoid-copy-reply 看阈值。把阈值设成 0 关掉 CA,大消息也会好——侧面证明根因是 CA 绕过 defer,不是 pubsub 逻辑本身。