跳过正文

RESP3 push frame torn apart when a pubsub message >= 16384 bytes is delivered to the publishing connection (9.0 regression)

作者
杨全烨
系统软件:操作系统、网络与分布式系统。
目录

binbin第一次直接assign给我的issue,我来尝试处理一下并且学习一下这里的代码。

原issue界面:https://github.com/valkey-io/valkey/issues/4231


1. 这个 bug 发生在什么业务场景
#

一条 TCP 连接同时干两件事:

  1. SUBSCRIBE wiretest —— 成为 channel 的订阅者
  2. 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.hclient 定义):

结构 作用
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);
}

注意顺序:

  1. pubsubPublishMessage... —— 给所有订阅者排队 push
  2. 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 段:

  1. >3\r\n(或 RESP2 的 *3
  2. message
  3. channel 名(bulk)
  4. 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
#

addReplyPubsubMessagepushing=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. 读代码时建议的顺序(当「导览」)
#

  1. publishCommand / addReplyPubsubMessage —— src/pubsub.c
  2. _addReplyToBufferOrList 的 defer 注释 —— src/networking.c ~744
  3. afterCommandlistJoin —— src/server.c ~4195
  4. addReplyBulk / isCopyAvoidPreferred / _addBulkStrRefToBufferOrList —— 同文件
  5. addBulkStringToReplyIOV —— 理解 CA 如何变成 $len\r\ndata

本地验证:起一个实例,跑 issue 里的 Python;或 CONFIG GET min-string-size-avoid-copy-reply 看阈值。把阈值设成 0 关掉 CA,大消息也会好——侧面证明根因是 CA 绕过 defer,不是 pubsub 逻辑本身。