跳过正文

多线程下的RDMA的崩溃问题(无锁队列IO模型)

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

If it’s broken, I’ll fix it!

现在重新引入了一个有意思的问题(背景阐述)
#

描述: 本次问题解决的是RDMA的多线程导致的崩溃问题. 在我们解决这个问题的同时,实际上主分支正在重构I/O的模型,我们不知道新的I/O模型是否还会出现这样的问题. (之后再进行验证)

同时我们还需要理解新的I/O模型(之后看PR历史进行学习,这个可以在特性:RedisIO模型中查看问题所在).

大致是这样的:

flowchart LR
    Main["主线程"]
    IO["IO 线程"]
    Inbox["io_shared_inbox<br/>SPMC"]
    Outbox["io_shared_outbox<br/>MPSC"]

    Main -->|"JOB_REQ_READ/WRITE"| Inbox
    Inbox --> IO
    IO --> Outbox
    Outbox -->|"JOB_RES_*"| Main

因为现在的主分支的I/O模型已经发生了变化,我想先验证一下在新的模型下,原来的RDMA多线程是否还会存在问题. 我想开一个分支验证一下. 是会存在问题的,这很好.

现在想在新的I/O模型下测试
#

如果按照原有的力度进行测试,实际上会出问题,但是该概率非常之小,比如如下这样的指令就不会出现问题.

./src/valkey-benchmark -h 10.0.0.1 -p 6379 -d 256 --threads 16 -c 100 -P 32 -n 20000000 -t get --rdma

通过复现可以看出来实际上还是会存在问题,但是概率很小?→扩大压测的力度?

# -c 500 增加并发,-P 256 极限流水线,极容易打满 RDMA 的网卡队列
./src/valkey-benchmark -h 10.0.0.1 -p 6379 -d 256 --threads 16 -c 500 -P 256 -n 50000000 -t set,get --rdma

这个混合测试真的很猛,对于新的I/O模型应该是直接可以报错,错误就还是
c->cmd_queue.len == 0' is not true 断言的问题.

这个命令太猛了,就是直接在我们当前的代码测试也会直接I/O hang. 证明还是存在问题的. (这里的思考就是我们压测的力度是否是合理的,就是因为我们的压测命令的参数是没有上限的,难道对于任意的压测指令,我们的程序都应该正常的运转?但是无论如何断言错误是不应该出现的对么?是的,可以处理的慢,但是不能出现段错误) 1.新模型的I/O 会直接assertion报错 2.原来的模型,这个也是会死锁.(或者是出现I/O hang的问题?这个问题可以暂时先放着.)

# -d 8192 甚至更高,测试大块数据在 RDMA 传输时是否会导致内存踩踏或分配器崩溃
./src/valkey-benchmark -h 10.0.0.1 -p 6379 -d 8192 --threads 16 -c 200 -P 64 -n 10000000 -t set --rdma

这个就是很慢,但是最终还是可以完成,感觉也是存在一些问题. 这个我们就暂时认为是没有问题的. 现在也不会卡死了,我们真的在解决问题。

继续尝试解决这个问题
#

server进行编译:

sudo prlimit --memlock=unlimited --pid $$
sudo modprobe dummy
sudo ip link add dummy0 type dummy
sudo ip link set dummy0 up
sudo ip addr add 10.0.0.1/24 dev dummy0
sudo rdma link add rxe_dummy type rxe netdev dummy0
sudo systemctl stop valkey
make -j$(nproc) BUILD_RDMA=yes USE_FAST_FLOAT=yes # 注意这里开启多线程的编译
sudo ./src/valkey-server valkey.conf --rdma-bind 10.0.0.1 --rdma-port 6379

每当端口真的被占用的时候:

sudo lsof -i :6379

lsof (List Open Files) 用于列出当前系统打开的文件。 在 Linux 中“一切皆文件”,网络 socket 也被视为文件。-i :6379 表示过滤出占用 6379 端口的网络连接。

~/Project/valkey fix-rdma-assertion-new-io* ❯ sudo lsof -i :6379
COMMAND     PID USER FD   TYPE DEVICE SIZE/OFF NODE NAME
valkey-se 26355 root 8u  IPv4 148550      0t0  TCP *:redis (LISTEN)
~/Project/valkey fix-rdma-assertion-new-io* ❯ sudo kill 26355   

client测试:

sudo prlimit --memlock=unlimited --pid $$
./src/valkey-benchmark -h 10.0.0.1 -p 6379 -d 256 --threads 16 -c 500 -P 256 -n 50000000 -t set,get --rdma

然后我们的valkey.conf是这样. 确实,这个单线程确实没有问题,只会在存在多个io线程的时候,才会报错.

io-threads 4
save ""
protected-mode no
bind 0.0.0.0

然后我们还是能稳定复现之前的错误.

细节讨论
#

你还是首先需要深刻理解IO对于客户端任务的流转特性:RedisIO模型 一定要深刻理解我们需要如何解决这个问题,这是我们第一次和pizhenwei进行比较深刻的技术讨论。 画了一张架构图深刻理解问题在哪里,这是我的得意之作。(虽然看起来很垃圾)

首先声明和梳理几个变量:

状态机 含义
io_read_state / io_write_state offload 生命周期:IDLE → PENDING_IO → COMPLETED_IO → IDLE
cmd_queue IO 已解析、主线程尚未执行的 pipeline 命令->这里就是解析命令的作用

这里的断言就是你不能还没执行之前的命令,就又来解析新的命令。

RDMA 特殊之处:没有独立 POLLOUT事件; connUpdateState() → connRdmaEventHandler() 会同步调 read_handler / write_handler

解决问题的时候发现GET居然还会出现access nil的错误
#

这种多线程的问题真的太复杂了。難しい過ぎます。

错误日志太长,我们把这个日志放在一个文件内部进行分析, 有的时候大多数错误没有意义,但是你看不到最前面的错误了。

sudo ./src/valkey-server valkey.conf --rdma-bind 10.0.0.1 --rdma-port 6379 > valkey-run.log 2>&1

我觉得可能还是之前在多线程下RDMA的崩溃问题老的线程模型中引入过的问题,我们可以继续探究一下,这个印度人确实很厉害,但是实际上看起来不是一个问题。 使用:

sudo prlimit --memlock=unlimited --pid $$
./src/valkey-benchmark -h 10.0.0.1 -p 6379 -d 256 --threads 16 -c 500 -P 256 -n 50000000 -t set,get --rdma

进行压测的时候,进入GET阶段的时候,如果ctrl + c则server必然会crash导致错误,这同样还是一个严重的问题。 就是遇到这种崩溃,你具体应该怎么抓,就是谁的调用才导致了本次的崩溃? 先检查server二进制文件是否存在调试符号:

~/Project/valkey fix-rdma-assertion-new-io* ❯ file ./src/valkey-server
./src/valkey-server: ELF 64-bit LSB pie executable, x86-64, version 1 (GNU/Linux), dynamically linked, interpreter /lib64/ld-linux-x86-64.so.2, BuildID[sha1]=afbcceb012f03c869136fe49f6f71261a65185af, for GNU/Linux 4.4.0, with debug_info, not stripped

然后在gdb内部运行:

~/Project/valkey fix-rdma-assertion-new-io* ❯ sudo gdb --args ./src/valkey-server valkey.conf --rdma-bind 10.0.0.1 --rdma-port 6379

先跑server:

set pagination off
handle SIGPIPE nostop noprint
run

等到崩溃之后:

thread apply all bt
bt
frame 1
info frame

这样可以来查看thread frame来找到底是谁调用到了access nil的函数。 这个就非常爽了,直接就能查看到底是哪里炸掉了:

Thread 1 (Thread 0x7ffff7c90740 (LWP 42401) "valkey-server"):
#0  listUnlinkNode (list=0x7ffff743d420, node=0x0) at /home/ada/Project/valkey/src/adlist.c:194
#1  0x000055555570ef56 in listDelNode (list=0x7ffff743d420, node=0x0) at /home/ada/Project/valkey/src/adlist.c:183
#2  rdmaProcessPendingData () at /home/ada/Project/valkey/src/rdma.c:1785
#3  0x0000555555747f93 in connTypeProcessPendingData () at /home/ada/Project/valkey/src/connection.c:147
#4  beforeSleep (eventLoop=<optimized out>) at /home/ada/Project/valkey/src/server.c:1878
#5  beforeSleep (eventLoop=<optimized out>) at /home/ada/Project/valkey/src/server.c:1842
#6  0x00005555555fa1df in aeProcessEvents (flags=27, eventLoop=0x7ffff744d000) at /home/ada/Project/valkey/src/ae.c:426
#7  aeMain (eventLoop=0x7ffff744d000) at /home/ada/Project/valkey/src/ae.c:543
#8  0x00005555555e94fb in main (argc=6, argv=0x7fffffffe298) at /home/ada/Project/valkey/src/server.c:7795

从下往上可以清晰地看到调用链,真是分析的利器。 可以看到上面的node=0x0的时候就是出问题的时候。 beforeSleep → connTypeProcessPendingData → rdmaProcessPendingDatardma.c:1785)→ listDelNode(pending_list, rdma_conn->pending_list_node) → listUnlinkNode,其中 node == NULL。 这里就是进行了多次释放导致问题的出现。

闹了半天又绕回来了,实际上这个问题就是我们在Valkey Over RDMA测试中server UAF崩溃的问题中解决的单独的ctrl + c导致server crash的问题 此时已经闭环了,理论上只要合入最新的代码就不会存在这个问题了link 去看看修复的原理。 测试不存在这个问题了,那真的很爽了。

细节打磨
#

目前看起来好像是解决了问题,但是细节还需要进行进一步打磨. 经过初步的修改我们发现了很有趣的地方,就是对于这个混合指令集:

./src/valkey-benchmark -h 10.0.0.1 -p 6379 -d 256 --threads 16 -c 500 -P 256 -n 50000000 -t set,get --rdma

实际上这里应该是先进行很多次set,然后紧接着进行很多次get. 原来没有调整顺序之前,set没有问题,但是之后紧接着get的时候会报这个错误,就是我认领的另外一个issue: https://github.com/valkey-io/valkey/issues/3356 那我们就在这个基础上继续去解决问题.

我现在不太清楚的一个点就是,我们在新的I/O模型下的改动是否能直接应用到旧的分支上, 这样我们就不需要开启两个PR来尝试解决这个问题了,就是还是只需要在原来的基础上修改即可. 我们现在清楚手法,但是需要一个更干净的手段来处理问题

这个应该是不行的,针对最新的unstable分支和原来的老版本,我们需要使用不同的手法来进行处理. 还是先尝试实验性的修改, 难度真的很大,等到有评论之后再进行推进. 目前先提交一版代码,之后存在问题的情况再说吧.

大概看看新的IO模型
#

特性:RedisIO模型 相当于在这里IO线程由被动变成主动的了 架构对比图 旧架构 (轮询模式): 主线程 IO线程1 IO线程2 IO线程3

–[任务]–>[队列1]——>
–[任务]–>[队列2]——>———–>
–[任务]–>[队列3]——>———————>
<-[轮询检查所有客户端]<-
(低效,浪费CPU)

新架构 (事件驱动): 主线程 IO线程1 IO线程2 IO线程3

↓ ↓ ↓
–[任务]–>[共享SPMC队列]←-竞争拉取-
(自动负载均衡)
<—[MPSC响应队列]<—–+—-+—-+
(快速通知,无需轮询)

目前还是存在IO hang的问题
#

不能work around,能卡住证明这个改动本身就是存在问题的. htop排查问题 使用这个来看内存,CPU的使用率,问题可能出在哪里? perf查看热点函数 找到了某个进程之后,你就可以查看热点函数在哪里,分析为什么停在了这里. 我还是想再研究一下,旧的逻辑是否能够直接应用上去. 现在找不到卡死的原因.

我们记录一下解决卡死问题的过程.
#

目前就是偶发的会卡住. 现象: 1.一般两个CPU核心空转. 2.结束之后,再起benchmark的时候还是会直接卡死.

先看看卡在client还是server
#

这里利用strace命令来进行分析,就是先查看benchmark上的线程都在干什么?

pid="$(pgrep -nf valkey-benchmark)"
sudo timeout 15 strace -c -f -p "$pid"

timeout 15 15s自动停止,进行所有线程的统计,可以观察到现象:

~/Project/valkey fix-rdma-assertion-new-io* ❯ pid="$(pgrep -nf valkey-benchmark)"
sudo timeout 15 strace -c -f -p "$pid"
strace: Process 64770 attached with 17 threads
strace: Process 64791 detached
strace: Process 64793 detached
strace: Process 64792 detached
strace: Process 64790 detached
strace: Process 64789 detached
strace: Process 64788 detached
strace: Process 64787 detached
strace: Process 64786 detached
strace: Process 64785 detached
strace: Process 64784 detached
strace: Process 64783 detached
strace: Process 64782 detached
strace: Process 64781 detached
strace: Process 64780 detached
strace: Process 64779 detached
strace: Process 64778 detached
strace: Process 64770 detached
% time     seconds  usecs/call     calls    errors syscall
------ ----------- ----------- --------- --------- ----------------
 99.89    0.233949         244       956           epoll_wait
  0.11    0.000252           4        59           write
------ ----------- ----------- --------- --------- ----------------
100.00    0.234201         230      1015           total

benchmark的程序都在epoll_wait,这意味着server侧的RDMA路径应该卡住,并且不再回包导致的.

接下来我们看看server到底在干什么
#

PID=64539
top - 23:06:29 up  9:08,  1 user,  load average: 2.75, 2.84, 2.81
Threads: 13 total, 2 running, 11 sleep, 0 d-sleep, 0 stopped, 0 zombie
%Cpu(s):  7.7 us,  0.5 sy,  0.0 ni, 91.1 id,  0.5 wa,  0.2 hi,  0.2 si,  0.0 st
MiB Mem :  15673.7 total,   1471.2 free,  12377.8 used,   4286.4 buff/cache
MiB Swap:   4096.0 total,   2244.7 free,   1851.3 used.   3295.9 avail Mem

    PID USER      PR  NI    VIRT    RES    SHR S  %CPU  %MEM     TIME+ COMMAND
  64539 root      20   0 2946372   2.6g   1.5g R  99.3  16.9  40:12.97 valkey-+
  64545 root      20   0 2946372   2.6g   1.5g R  99.3  16.9  40:16.35 io_thd_1
  64540 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 bio_clo+
  64541 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 bio_aof
  64542 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 bio_laz+
  64543 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 bio_rdb+
  64544 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 bio_tls+
  64546 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:06.25 io_thd_2
  64547 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:05.07 io_thd_3
  64548 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 jemallo+
  64549 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 jemallo+
  64550 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 jemallo+
  64551 root      20   0 2946372   2.6g   1.5g S   0.0  16.9   0:00.00 jemallo+                                                      

你应该能看到是一个主线程CPU和一个IO线程在忙等?主要就是这两个线程.

不使用GDB查看内核栈的快照
#

sudo cat /proc/$pid/stack
for tid in /proc/$pid/task/*; do
  echo "=== $(basename $tid) ==="
  sudo cat $tid/stack 2>/dev/null | head -20
done | less

遍历所有的pid下的tid来查看栈. 但是上面的64539和64545没有抓到?

接着我们使用GDB来查看快照
#

~/Project/valkey fix-rdma-assertion-new-io* ❯ sudo gdb -p "$pid" -batch \
  -ex "set pagination off" \
  -ex "thread apply all bt" \
  -ex "detach" \
  -ex "quit" \
  2>&1 | tee /tmp/valkey-hang-bt.txt

# lightweight process 轻量级进程:这12个是子进程 最后面的64539是主进程
[New LWP 64551]
[New LWP 64550]
[New LWP 64549]
[New LWP 64548]
[New LWP 64547]
[New LWP 64546]
[New LWP 64545]
[New LWP 64544]
[New LWP 64543]
[New LWP 64542]
[New LWP 64541]
[New LWP 64540]
[Thread debugging using libthread_db enabled]
Using host libthread_db library "/usr/lib/libthread_db.so.1".
postponeClientRead (c=0x7f00c6b68e80) at /home/ada/Project/valkey/src/networking.c:6409
6409        if (ProcessingEventsWhileBlocked) return 0;

Thread 13 (Thread 0x7f00ff7326c0 (LWP 64540) "bio_close_file"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f21b00b in mutexQueuePop (theQueue=0x7f0100852380, blocking=<optimized out>) at /home/ada/Project/valkey/src/mutexqueue.c:123                                                                            
#5  0x0000558c5f15ecee in bioProcessBackgroundJobs (arg=0x558c5f4f1280 <bio_workers.lto_priv>) at /home/ada/Project/valkey/src/bio.c:269                                                                                
#6  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#7  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 12 (Thread 0x7f00fef316c0 (LWP 64541) "bio_aof"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f21b00b in mutexQueuePop (theQueue=0x7f01008523f0, blocking=<optimized out>) at /home/ada/Project/valkey/src/mutexqueue.c:123                                                                            
#5  0x0000558c5f15ecee in bioProcessBackgroundJobs (arg=0x558c5f4f1298 <bio_workers.lto_priv+24>) at /home/ada/Project/valkey/src/bio.c:269                                                                             
#6  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#7  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 11 (Thread 0x7f00fe7306c0 (LWP 64542) "bio_lazy_free"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f21b00b in mutexQueuePop (theQueue=0x7f0100852460, blocking=<optimized out>) at /home/ada/Project/valkey/src/mutexqueue.c:123                                                                            
#5  0x0000558c5f15ecee in bioProcessBackgroundJobs (arg=0x558c5f4f12b0 <bio_workers.lto_priv+48>) at /home/ada/Project/valkey/src/bio.c:269                                                                             
#6  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#7  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 10 (Thread 0x7f00fdf2f6c0 (LWP 64543) "bio_rdb_save"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f21b00b in mutexQueuePop (theQueue=0x7f01008524d0, blocking=<optimized out>) at /home/ada/Project/valkey/src/mutexqueue.c:123                                                                            
#5  0x0000558c5f15ecee in bioProcessBackgroundJobs (arg=0x558c5f4f12c8 <bio_workers.lto_priv+72>) at /home/ada/Project/valkey/src/bio.c:269                                                                             
#6  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#7  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 9 (Thread 0x7f00fd72e6c0 (LWP 64544) "bio_tls_reload"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f21b00b in mutexQueuePop (theQueue=0x7f0100852540, blocking=<optimized out>) at /home/ada/Project/valkey/src/mutexqueue.c:123                                                                            
#5  0x0000558c5f15ecee in bioProcessBackgroundJobs (arg=0x558c5f4f12e0 <bio_workers.lto_priv+96>) at /home/ada/Project/valkey/src/bio.c:269                                                                             
#6  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#7  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 8 (Thread 0x7f00fcf2d6c0 (LWP 64545) "io_thd_1"):
#0  0x0000558c5f21a612 in __rdtsc () at /usr/lib/gcc/x86_64-pc-linux-gnu/15.2.1/include/ia32intrin.h:114
#1  getMonotonicUs_x86 () at /home/ada/Project/valkey/src/monotonic.c:46
#2  0x0000558c5f1dfa66 in IOThreadMain (myid=<optimized out>) at /home/ada/Project/valkey/src/io_threads.c:314                                                                                                          
#3  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#4  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 7 (Thread 0x7f00fc72c6c0 (LWP 64546) "io_thd_2"):
#0  0x00007f0100c938d0 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9a0c4 in pthread_mutex_lock () from /usr/lib/libc.so.6
#2  0x0000558c5f1dfbf3 in IOThreadMain (myid=<optimized out>) at /home/ada/Project/valkey/src/io_threads.c:384                                                                                                          
#3  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#4  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 6 (Thread 0x7f00fbf2b6c0 (LWP 64547) "io_thd_3"):
#0  0x00007f0100c938d0 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9a0c4 in pthread_mutex_lock () from /usr/lib/libc.so.6
#2  0x0000558c5f1dfbf3 in IOThreadMain (myid=<optimized out>) at /home/ada/Project/valkey/src/io_threads.c:384                                                                                                          
#3  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#4  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 5 (Thread 0x7f00fb72a6c0 (LWP 64548) "jemalloc_bg_thd"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f34a8f5 in background_thread_sleep (tsdn=<optimized out>, info=<optimized out>, interval=<optimized out>) at src/background_thread.c:137                                                                 
#5  background_work_sleep_once (tsdn=<optimized out>, info=<optimized out>, ind=0) at src/background_thread.c:229                                                                                                       
#6  background_thread0_work (tsd=<optimized out>) at src/background_thread.c:374
#7  background_work (tsd=<optimized out>, ind=<optimized out>) at src/background_thread.c:412
#8  background_thread_entry (ind_arg=<optimized out>) at src/background_thread.c:444
#9  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#10 0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 4 (Thread 0x7f00fa7ff6c0 (LWP 64549) "jemalloc_bg_thd"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f34a239 in background_thread_sleep (tsdn=<optimized out>, info=0x7f0100a16950, interval=<optimized out>) at src/background_thread.c:137                                                                  
#5  background_work_sleep_once (tsdn=<optimized out>, info=<optimized out>, ind=<optimized out>) at src/background_thread.c:229                                                                                         
#6  background_work (tsd=<optimized out>, ind=<optimized out>) at src/background_thread.c:419
#7  background_thread_entry (ind_arg=<optimized out>) at src/background_thread.c:444
#8  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#9  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 3 (Thread 0x7f00f9ffe6c0 (LWP 64550) "jemalloc_bg_thd"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f34a239 in background_thread_sleep (tsdn=<optimized out>, info=0x7f0100a16a20, interval=<optimized out>) at src/background_thread.c:137                                                                  
#5  background_work_sleep_once (tsdn=<optimized out>, info=<optimized out>, ind=<optimized out>) at src/background_thread.c:229                                                                                         
#6  background_work (tsd=<optimized out>, ind=<optimized out>) at src/background_thread.c:419
#7  background_thread_entry (ind_arg=<optimized out>) at src/background_thread.c:444
#8  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#9  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 2 (Thread 0x7f00f93ff6c0 (LWP 64551) "jemalloc_bg_thd"):
#0  0x00007f0100c9ef32 in ?? () from /usr/lib/libc.so.6
#1  0x00007f0100c9339c in ?? () from /usr/lib/libc.so.6
#2  0x00007f0100c9368c in ?? () from /usr/lib/libc.so.6
#3  0x00007f0100c95e5e in pthread_cond_wait () from /usr/lib/libc.so.6
#4  0x0000558c5f34a239 in background_thread_sleep (tsdn=<optimized out>, info=0x7f0100a16af0, interval=<optimized out>) at src/background_thread.c:137                                                                  
#5  background_work_sleep_once (tsdn=<optimized out>, info=<optimized out>, ind=<optimized out>) at src/background_thread.c:229                                                                                         
#6  background_work (tsd=<optimized out>, ind=<optimized out>) at src/background_thread.c:419
#7  background_thread_entry (ind_arg=<optimized out>) at src/background_thread.c:444
#8  0x00007f0100c9697a in ?? () from /usr/lib/libc.so.6
#9  0x00007f0100d1a2bc in ?? () from /usr/lib/libc.so.6

Thread 1 (Thread 0x7f0100e32740 (LWP 64539) "valkey-server"):
#0  postponeClientRead (c=0x7f00c6b68e80) at /home/ada/Project/valkey/src/networking.c:6409
#1  readQueryFromClient (conn=0x7f00f3035f40) at /home/ada/Project/valkey/src/networking.c:4256
#2  0x0000558c5f260a9a in callHandler (conn=0x7f00f3035f40, handler=<optimized out>) at /home/ada/Project/valkey/src/connhelpers.h:79                                                                                   
#3  connRdmaEventHandler (el=<optimized out>, fd=<optimized out>, clientData=0x7f00f3035f40, mask=<optimized out>) at /home/ada/Project/valkey/src/rdma.c:735                                                           
#4  0x0000558c5f22d3b5 in connUpdateState (conn=<optimized out>) at /home/ada/Project/valkey/src/connection.h:524                                                                                                       
#5  processClientIOWriteDone (c=0x7f00c6b68e80) at /home/ada/Project/valkey/src/networking.c:3215
#6  processClientIOWriteDone (c=0x7f00c6b68e80) at /home/ada/Project/valkey/src/networking.c:3204
#7  0x0000558c5f1dfec7 in handleWriteJobs (write_jobs=0x7fff24438e80, write_count=<optimized out>) at /home/ada/Project/valkey/src/io_threads.c:884                                                                     
#8  processIOThreadsResponses () at /home/ada/Project/valkey/src/io_threads.c:934
#9  0x0000558c5f29d228 in processIOThreadsResponses () at /home/ada/Project/valkey/src/io_threads.c:73
#10 beforeSleep (eventLoop=<optimized out>) at /home/ada/Project/valkey/src/server.c:1975
#11 beforeSleep (eventLoop=<optimized out>) at /home/ada/Project/valkey/src/server.c:1842
#12 0x0000558c5f14f1df in aeProcessEvents (flags=27, eventLoop=0x7f010084d000) at /home/ada/Project/valkey/src/ae.c:426                                                                                                 
#13 aeMain (eventLoop=0x7f010084d000) at /home/ada/Project/valkey/src/ae.c:543
#14 0x0000558c5f13e4fb in main (argc=6, argv=0x7fff24439388) at /home/ada/Project/valkey/src/server.c:7795
[Inferior 1 (process 64539) detached]

从下往上读(从外层函数向内层核心深挖)。 看到这个#0的位置,就是当前的线程停止在了

我们尝试去解决root cause: 就是一个状态机错位的问题,都已经写到了这个修复CPU空转的问题commit message的内部,可以复习. 这个也是经典的状态机错乱导致的错误。

语义层面的重构
#

我们分析一下有什么重构的手法?

引入这个问题的原因是,offload之前,需要把postpone这个东西置起来,结束之后回回来,然后再updateEvent.

有个逻辑没有用到,就是postpone_update_state的时候,有个参数没有使用。从两个不同的路径上触发的,有可能可以使用。 正常的event只有pollin才能进来(因为这里只是注册了pollin事件?)。 mask 只有pollin event才能进来。 我们从上面调用的那个handler里触发下来的话,可以使用其他参数来控制一下的,控制它是不是调用readhandler和writehandler的操作, 使用mask,我们知道我们的触发路径是什么,要么就是靠系统(RDMA的 completion queue)的pollin事件从aeEvent进来的,这个里面mask肯定只有pollin一种可能, postpone置成0,updatestate的时候,这条路径上我们可以使用一个特定的mask,看下这个error是怎么触发的? 或者提前把某个mask reserve起来,internal的一种用法,就是我们自己使用,通过这个来决定是不是要唤醒某些特定的handler,有可能能解决这个问题。 就是不引入我们加入的int标志位,postpone这个东西应该是TLS引入的,RDMA不希望让networking层的代码变脏了,但是可以尝试一下。

理一遍Root Cause
#

我们走读一边代码重新理一下。 同时在解决问题的时候建立这样的思想,主线程和IO线程之间就是利用无锁队列来进行通信的。 目前来说root cause已经找的非常清晰了,剩下的只是我们解决问题的手法是否漂亮的问题了.

1.主线程状态分析
#

我们从触发读回调开始分析(因为本身重入问题就出在读回调函数里面),读回调函数是readQueryFromClient

void readQueryFromClient(connection *conn) {
	client *c = connGetPrivateData(conn);
	/* Check if we can send the client to be handled by the IO-thread */
	// 尝试能否推迟ClientRead
	if (postponeClientRead(c)) return;
	if (c->io_write_state != CLIENT_IDLE || c->io_read_state != CLIENT_IDLE) return;
	......
	}

调到postponeClientRead:(所谓推迟,就是看能否暂时交给IO线程来处理)

/* Return 1 if the client read is handled using threaded I/O.
* 0 otherwise. */
int postponeClientRead(client * c) {
    if (ProcessingEventsWhileBlocked) return 0;
    // 尝试看这个client read是否能够派发给IO线程来进行处理。
    return (trySendReadToIOThreads(c) == C_OK);
}

继续调用trySendReadToIOThreads.同注释的逻辑:

int trySendReadToIOThreads(client * c) {
	// 这里有自己的计算方法
    if (server.active_io_threads_num <= 1) return C_ERR;
    /* Fake/teardown clients may have no connection; never offload those. */
    if (!c - > conn) return C_ERR;
    /* If IO thread is still reading, return C_OK so the main thread does not race it. */
    if (c - > io_read_state == CLIENT_PENDING_IO) return C_OK;
    /* A completed read must be finished by processClientIOReadsDone on the main thread
     * before we try to offload another read; do not treat it like PENDING_IO. */
       
    // 这里返回C_ERR来解决状态机紊乱的问题。 目前的状态是C
    
    if (c - > io_read_state == CLIENT_COMPLETED_IO) return C_ERR;
    if (c - > io_write_state == CLIENT_PENDING_IO) return C_OK;
    /* For simplicity, don't offload replica clients reads as read traffic from replica is negligible */
    if (getClientType(c) == CLIENT_TYPE_REPLICA) return C_ERR;
    /* With Lua debug client we may call connWrite directly in the main thread */
    if (c - > flag.lua_debug) return C_ERR;
    /* For simplicity let the main-thread handle the blocked clients */
    if (c - > flag.blocked || c - > flag.unblocked) return C_ERR;
    if (c - > flag.close_asap) return C_ERR;

    c - > read_flags = canParseCommand(c) ? 0 : READ_FLAGS_DONT_PARSE;
    c - > read_flags |= authRequired(c) ? READ_FLAGS_AUTH_REQUIRED : 0;
    c - > read_flags |= isReplicatedClient(c) ? READ_FLAGS_REPLICATED : 0;
    c - > io_read_state = CLIENT_PENDING_IO;
    connSetPostponeUpdateState(c - > conn, 1);
	
	// 单生产者多消费者队列
    if (unlikely(spmcEnqueue( & io_shared_inbox, tagJob(c, JOB_REQ_READ_CLIENT)) == false)) {
        c - > read_flags = 0;
        c - > io_read_state = CLIENT_IDLE;
        connSetPostponeUpdateState(c - > conn, 0);
        return C_ERR;
    }

	// 此时真的offload当前这个readjob
    io_jobs_submitted++;
    server.stat_io_reads_pending++;
    c - > flag.pending_read = 1;
    return C_OK;
}

上面返回成功之后,我们本次的clientRead就应该已经offload给了IO线程来进行处理。 注意这里的两个点(这条线上我们继续分析主线程,你也可以暂时先走读下面client的代码):

c -> io_read_state = CLIENT_PENDING_IO;
connSetPostponeUpdateState(c -> conn, 1);

我们设置read状态的第一个状态机CLIENT_PENDING_IO,同时调用了connSetPostponeUpdateState, 如果这个函数指针是被注册过的话,我们就回调这个函数

static inline void connSetPostponeUpdateState(connection *conn, int on) {
	if (conn->type->postpone_update_state) {
		conn->type->postpone_update_state(conn, on);
	}
}

这个回调指针是在connection.c中进行定义的:

typedef struct connectionType {
	/* Postpone update state - with IO threads & TLS we don't want the IO threads to update the event loop events - let
* the main-thread do it */
  // 这个回调指针是TLS引入的,就是说在TLS连接层和开启IO线程的情况下,我们不希望IO线程来更新 事件循环事件(这是什么?)而是使用主线程来进行更新。
	void (*postpone_update_state)(struct connection *conn, int);
}

我们来分析这个postpone_update_state到底干了什么,有什么作用,实际上只有RDMA和TLS使用到了这个回调指针:

static ConnectionType CT_RDMA = {
......
.postpone_update_state = postPoneUpdateRdmaState, 
......
}

上面是具体的函数,我们来看看这个函数postPoneUpdateRdmaState:

static void postPoneUpdateRdmaState(struct connection *conn, int postpone) {
rdma_connection *rdma_conn = (rdma_connection *)conn;
if (postpone) {
	rdma_conn->flags |= RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE;
} else {
	rdma_conn->flags &= ~RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE;
}
}

就是如果需要postpone(就是传进来的是非零的时候,我们就加上RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE这个标志位,不然的话一定要去掉这个标志位),那我们看看哪里使用到了这个标志位?那么理论上在当前的连接之内,我们都是可以使用这个标志位来进行操作的。(实际上就是在接收到IO线程的sendToMainThread之后会使用到这些代码,所以我们可以继续从_main开始接收IO线程的回复结果之后开始分析。)

#define RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE (1 << 0)

上面的定义就是1。

============ 分割线 ============

_main线程等到IO线程返回结果之后的代码逻辑
#

beforeSleep 在每轮「准备阻塞等待 I/O」之前执行: 上一轮 poll 触发的可读/可写回调已经跑完, beforeSleep 负责把这一阶段该收尾的活集中做完,然后再进 poll(或等价地 custompoll/IO 线程上的 poll)。

我们需要注意的是beforeSleep中调用processIOThreadsResponses()中: 注意这里是在干什么? 把 IO 线程完成的读/写/accept 结果收回主线程:对应 networking/client 状态的推进。

int processIOThreadsResponses(void) {

    /* We don't check for threads number since some threads may return jobs then deactivate/shut-down */
    /* Quick check if any pending operations exist */
    if (getPendingIOResponsesCount() == 0) return 0;
    int total_processed = 0;
    void * jobs[JOB_BATCH_SIZE];
    client * read_jobs[JOB_BATCH_SIZE];
    client * write_jobs[JOB_BATCH_SIZE];

    /* Loop until we consume all pending jobs */
    // 循环遍历处理之前IO线程返回的每一个JOB
    while (1) {
        int received_responses = 0;
        int dequeued_count = 0;
        int read_count = 0;
        int write_count = 0;
        /* Try to dequeue JOB_BATCH_SIZE */
        while (received_responses < JOB_BATCH_SIZE) {

            // 从多生产者单消费者队列中取jobs进行处理
            // 批量进行出队的操作
            dequeued_count = mpscDequeueBatch( & io_shared_outbox, jobs, JOB_BATCH_SIZE - received_responses);
            /* Stop if we can't get more jobs from the queue. */

            if (dequeued_count == 0) break;
            received_responses += dequeued_count;
            total_processed += dequeued_count;
            
            for (int i = 0; i < dequeued_count; i++) {
                void * data;
                int job_type;
                // 这里把client给解出来
                untagJob(jobs[i], & data, & job_type);
                client * c = (client * ) data;
                
                if (job_type == JOB_RES_READ_CLIENT) {
                
                    // 此时的状态必须是IO处理结束的
                    serverAssert(c - > io_read_state == CLIENT_COMPLETED_IO);
                    read_jobs[read_count++] = c;
                } else if (job_type == JOB_RES_WRITE_CLIENT) {
                
                    serverAssert(c - > io_write_state == CLIENT_COMPLETED_IO);
                    write_jobs[write_count++] = c;
                } else {
                    serverPanic("Unknown job type %d", job_type);
                }
            }
        }
        // 对于要读和要写的操作调用不同的函数。
        if (read_count) handleReadJobs(read_jobs, read_count);
        if (write_count) handleWriteJobs(write_jobs, write_count);
        /* If the queue was empty at the last try - don't try again */
        if (dequeued_count == 0) return total_processed;
    }
}

看上面遍历了很多jobs,IO入队的时候存放的就是read jobs(我们目前整条链路都在分析读任务), 那么这里我们就对于整个read jobs收集的数组调用handleReadJobs这个函数, 我们来看看这个函数(这里是我们没有修改时候的函数,因为我们现在正在梳理root cause):

/* Function to handle read jobs */
static void handleReadJobs(client ** read_jobs, int read_count) {
	// 我们这次处理了这么多读任务
    server.stat_io_reads_pending -= read_count;
    serverAssert(server.stat_io_reads_pending >= 0);

    /* process each client */
    // 遍历每个client客户端
    for (int i = 0; i < read_count; i++) {
        client * c = read_jobs[i];
        processClientIOReadsDone(c);
    }

    /* Process commands in batch if we processed any reads */
    if (read_count) {
        server.stat_io_reads_processed += read_count;
        processClientsCommandsBatch();
    }
}

这里原来的逻辑就是,遍历处理每一个收集到的client, 然后调用processClientIOReadsDone,我们看看这个函数又干了什么?

这里相当于应该已经是收尾的处理了。

void processClientIOReadsDone(client * c) {
    serverAssert(c - > io_read_state == CLIENT_COMPLETED_IO);

    if (ProcessingEventsWhileBlocked) {
        /* When ProcessingEventsWhileBlocked we may call processIOThreadsReadDone recursively.
         * In this case, there may be some clients left in the batch waiting to be processed. */
        processClientsCommandsBatch();
    }

    c - > flag.pending_read = 0;
    c - > io_read_state = CLIENT_IDLE;
    /* Don't post-process-reads from clients that are going to be closed anyway. */
    if (c - > flag.close_asap) return;
    /* If a client is protected, don't do anything,
     * that may trigger read/write error or recreate handler. */
    if (c - > flag.protected) return;

    /* Save the current conn state, as connUpdateState may modify it */
    int in_accept_state = (connGetState(c - > conn) == CONN_STATE_ACCEPTING);
    connSetPostponeUpdateState(c - > conn, 0);
    connUpdateState(c - > conn);

    /* In accept state, no client's data was read - stop here. */
    if (in_accept_state) return;
    /* On read error - stop here. */
    if (handleReadResult(c) == C_ERR) {
        return;
    }

    if (!(c - > read_flags & READ_FLAGS_DONT_PARSE)) {
        parseResult res = handleParseResults(c);
        /* On parse error - stop here. */
        if (res == PARSE_ERR) {
            return;
        } else if (res == PARSE_NEEDMORE) {
            beforeNextClient(c);
            return;
        }
    }
    if (c - > argc > 0) {
        c - > flag.pending_command = 1;
    }
    /* try to add the command to the batch */
    int ret = addCommandToBatchAndProcessIfFull(c);
    /* If the command was not added to the commands batch, process it immediately */
    if (ret == C_ERR) {
        if (processPendingCommandAndInputBuffer(c) == C_OK) beforeNextClient(c);
    }
}

注意,这里更改了client的状态为IDLE,就是目前是处于空闲的状态,被重置了。

  c - > flag.pending_read = 0;
  c - > io_read_state = CLIENT_IDLE;
  /* Don't post-process-reads from clients that are going to be closed anyway. */
  if (c - > flag.close_asap) return;

进行状态更新,这里就是会出问题的地方:

    /* Save the current conn state, as connUpdateState may modify it */
    int in_accept_state = (connGetState(c - > conn) == CONN_STATE_ACCEPTING);
    connSetPostponeUpdateState(c - > conn, 0);
    connUpdateState(c - > conn);

在处理结束之后,我们又把这个postponeUpdateState函数的指针置回来了。processClientIOReadsDone中,我们进行了状态的更新,就是调用了RDMA的connUpdateState。来看看:

static inline void connUpdateState(connection *conn) {
	if (conn->type->update_state) {
		conn->type->update_state(conn);
	}
}

这里就是调用了当时注册的update_state的回调指针:

static ConnectionType CT_RDMA = {
......
.postpone_update_state = postPoneUpdateRdmaState,
.update_state = updateRdmaState,
......
};

我们来看看这个updateRdmaState:

static void updateRdmaState(struct connection * conn) {
    rdma_connection * rdma_conn = (rdma_connection * ) conn;
    connRdmaSetRwHandler(conn);
    connRdmaEventHandler(NULL, -1, rdma_conn, 0);
}

上面的connRdmaSetRwHandler,只是用来注册可读的事件handler,因为IB肯定只能注册POLLIN的事件。 之后的connRdmaEventHandler中,进行调用:

static void connRdmaEventHandler(struct aeEventLoop * el, int fd, void * clientData, int mask) {
    rdma_connection * rdma_conn = (rdma_connection * ) clientData;
    connection * conn = & rdma_conn - > c;
    struct rdma_cm_id * cm_id = rdma_conn - > cm_id;
    RdmaContext * ctx = cm_id - > context;

    int ret = 0;

    UNUSED(el);
    UNUSED(fd);
    UNUSED(mask);

    ret = connRdmaHandleCq(rdma_conn);
    if (ret == C_ERR) {
        conn - > state = CONN_STATE_ERROR;
        return;
    }

    /* uplayer should read all */
    while (!(rdma_conn - > flags & RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE) && ctx - > rx.pos < ctx - > rx.offset) {
        if (conn - > read_handler && (callHandler(conn, conn - > read_handler) == C_ERR)) {
            return;
        }
    }

    /* recv buf is full, register a new RX buffer */
    if (ctx - > rx.pos == ctx - > rx.length) {
        connRdmaRegisterRx(ctx, cm_id);
    }

    /* RDMA comp channel has no POLLOUT event, try to send remaining buffer */
    if (!(rdma_conn - > flags & RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE) && ctx - > tx.offset < ctx - > tx.length && conn - > write_handler) {
        callHandler(conn, conn - > write_handler);
    }
}

你会发现在这里更新状态的时候,还会调用callHandler,就是调到read_handler(之前注册的)。 read_hander是谁?readQueryFromClient,此时就导致了重入,从readQueryFromClient再走一遍上面的逻辑,一直调用到下面的:

void parseInputBuffer(client *c) {
	/* The command queue must be emptied before parsing. */
	// 指令队列在解析之前必须为空
	serverAssert(c->cmd_queue.len == 0);
	......
}

就是此时重入了,但是之前可能没有处理结束导致的断言错误(我的个人推断。) 实际上在processClientIOWriteDone中因为也会调用更新函数,导致进入读路径,造成另一处的:

/* If io_last_written_data_len is nonzero it must relate to c->buf */
serverAssert(c->io_last_written.data_len == 0 || c->io_last_written.buf == c->buf);

断言错误,而这个PR在尝试同时解决这两个问题。

IO线程状态分析
#

因为实际上二者的状态流转是交叉的,我们可以交替进行分析。我现在就是IO线程。

上面在主线程把client read offload给了我们之后,就是 我们从IO的主线程往下走:

static void * IOThreadMain(void * myid) {
        /* The ID is the thread ID number (from 1 to server.io_threads_num-1). ID 0 is the main thread. */
        long id = (long) myid;
        char thdname[32];
        ......
        while (1) {
            ......
            /* PRIORITY 2: Shared Global Queue (SPMC)
             * Only checked after SPSC is drained. */
			 // 这里就是经典的从环形缓冲区尝试进行出队的操作,单生产者多消费者。
            void * tagged_job = spmcDequeue( & io_shared_inbox);
            // 看需要消费的job类型
            if (tagged_job) {
                void * data;
                int type;
                untagJob(tagged_job, & data, & type);
                switch (type) {
                // 
                    case JOB_REQ_READ_CLIENT:
                        ioThreadReadQueryFromClient((client * ) data);
                        break;
                    case JOB_REQ_WRITE_CLIENT:
                        ioThreadWriteToClient((client * ) data);
                        break;
                    case JOB_REQ_FREE_OBJ:
                        zfree(data);
                        break;
                    case JOB_REQ_ACCEPT:
                        ioThreadAccept((client * ) data);
                        break;
                    case JOB_REQ_POLL:
                        ioThreadPoll((aeEventLoop * ) data);
                        break;
                    default:
                        serverPanic("Invalid SPMC job type: %d", type);
                }
                processed++;
                ......
            } {

_main线程入队的时候我们可以看到,那个时候我们给的job tag就是JOB_REQ_READ_CLIENT(里面的代码逻辑是入队失败的情况):

    if (unlikely(spmcEnqueue( & io_shared_inbox, tagJob(c, JOB_REQ_READ_CLIENT)) == false)) {
        c - > read_flags = 0;
        c - > io_read_state = CLIENT_IDLE;
        connSetPostponeUpdateState(c - > conn, 0);
        return C_ERR;
    }

OK,我们继续分析,在ioMain的switch语句中我们继续调用ioThreadReadQueryFromClient:

void ioThreadReadQueryFromClient(client * c) {
	// 断言,刚进这个函数的时候,我们的状态肯定是 `CLIENT_PENDING_IO` 就是等待处理的状态。
    serverAssert(c-> io_read_state == CLIENT_PENDING_IO);

    /* Read */
    readToQueryBuf(c);
    if (c - > flag.close_asap) {
        goto done;
    }

    /* Check for read errors. */
    if (c - > nread <= 0) {
        goto done;
    }

    /* Skip command parsing if the READ_FLAGS_DONT_PARSE flag is set. */
    if (c - > read_flags & READ_FLAGS_DONT_PARSE) {
        goto done;
    }

    /* Handle QB limit */
    if (c - > read_flags & READ_FLAGS_QB_LIMIT_REACHED) {
        goto done;
    }
	
	// ======== 这个解析就是重入点,之前的指令主线程还没处理结束 ========
    parseInputBuffer(c);
    
    trimCommandQueue(c);
    prepareCommandQueue(c);

    /* Parsing was not completed - let the main-thread handle it. */
    if (!(c - > read_flags & READ_FLAGS_PARSING_COMPLETED)) {
        goto done;
    }

    /* Empty command - Multibulk processing could see a <= 0 length. */
    if (c - > argc == 0) {
        goto done;
    }

    done:
        /* Only trim query buffer for non-primary clients
         * Primary client's buffer is handled by main thread using repl_applied position */
        if (!(c - > read_flags & READ_FLAGS_REPLICATED)) {
            trimClientQueryBuffer(c);
        }

    c - > io_read_state = CLIENT_COMPLETED_IO;
    c - > cur_tid = getCurTid();
    sendToMainThread(c, JOB_RES_READ_CLIENT);
}

上面这段代码中,我们读取到缓冲区并且进行相应的解析(这里的逻辑之后可以进一步研究,此处不是重点), 同时在处理完read client之后,注意这里我们更改状态机为CLIENT_COMPLETED_IO, 就是此时已经被IO处理完毕了,然后我们把处理的结果返回给主线程:

c - > io_read_state = CLIENT_COMPLETED_IO;
c - > cur_tid = getCurTid();
sendToMainThread(c, JOB_RES_READ_CLIENT);

我们继续看sendToMainThread这个函数都干了什么?

void sendToMainThread(void * data, int type) {
    if (unlikely(pending_io_responses)) {
        flushPendingIOResponses(0);
    }

	// 拿到我们IO当前完成的任务
    void * job = tagJob(data, type);
    // 把我们当前完成的job放进一个队列中,这是所有IO都放到一个队列内部,然后主线程进行对应的处理。
    // 所以就是多生产者(IO线程)单消费者(_main线程)的链表
    if (unlikely(pending_io_responses || !mpscEnqueue( & io_shared_outbox, job, & io_thread_ticket))) {
        /* Failed to push new job: initialize list if needed and save job */
        // 实在没有放进去的话我们就开启一个全新的队列(这个处理机制是什么???)
        if (pending_io_responses == NULL) {
            pending_io_responses = listCreate();
        }
        listAddNodeTail(pending_io_responses, job);
    }
}

OK,此时IO线程的使命已经结束,处理完read client之后,然后把结果进行相应的入队操作, 我们可以回到_main线程继续看做了什么?back to main

更优雅的解决方式?
#

在分析结束了root cause之后,我们想使用更好的方式尝试来解决问题,因为现在上层的代码看起来很脏。
就是你不能写 if (连接类型 = RDMA){ 处理逻辑 },这很不好看,有没有什么更优雅的处理手法。

思路,就是原来是直接在上层判断类型,然后修改逻辑,现在就是想利用postpone_update_state来使用更客观的手法来处理, 但是改动很多,非常容易出错感觉,这里还需要再仔细进行研究,这里的逻辑过于复杂了, 我们的想法就是逐渐给出一个最小的改动来处理,而不是一下子进行大改。 思考一下,这种修改本来就不可能一下就OK的,肯定非常复杂。

有的时候你想给AI看,或者要不分页的git diff来进行查询,就需要这样来进行处理:

$ git --no-pager -diff # 这个是你可以展示目前工作区的diff

但是有的时候已经commit,我们应该如何查看:

$ git --no-pager show <commit_id> # 不分页展示某个已经提交的commit

以上两个都非常有用。 目前来说看起来差不多了,从二月份,那个时候IO线程的job模型还没有被修改成MPMC,我处理的旧版的IO存在的问题, 之后IO模型修改,我又处理了新的版本,现在竟然到了五月底,时间真的好快,中间还解决了两个小问题,希望能顺利进行!

我们修改之后的时序类似:

sequenceDiagram
    participant Epoll as RDMA/epoll
    participant Main as 主线程
    participant IO as IO 线程
    participant RDMA as connRdmaEventHandler

    Epoll->>Main: readQueryFromClient
    Main->>IO: trySendReadToIOThreads → inbox
    Note over Main: PENDING_IO + postpone READ(这里推迟了读路径)

    IO->>IO: ioThreadReadQueryFromClient<br/>parseInputBuffer → cmd_queue
    IO->>Main: COMPLETED_IO → outbox

    Main->>Main: handleReadJobs ① processClientIOReadsDone
    Note over Main: IDLE, mask=READ(+WRITE?), connUpdateState
    Note over RDMA: READ 被挡,不进 read_handler

    Main->>Main: ② processClientsCommandsBatch
    Main->>Main: ③ drain + clientConnPostponeMask + connUpdateState
    Note over RDMA: queue 空 → 可排空 RX

这个 PR 的本质,就是在 client/connection 生命周期的特定节点, 设置 postpone_mask 里的 READ/WRITE 位; 在 RDMA 的 connRdmaEventHandler 里检查这些位, 从而决定「这一刻能不能同步调用 read/write handler」。

CI失败
#

目前我们发现了在当前这个分支上会出现一个CI测试的稳定失败test-ubuntu-latest-cmake-tls中的单个io_threads的测试。 那么到底是我的问题还是测试的问题? 可以尝试来解决一下这个问题,先尝试复现test吧:

# 使用 CMake 编译带 TLS 支持的 Valkey (与 CI 环境保持一致) 
cmake -B build -DUSE_TLS=ON 
cmake --build build -j$(nproc)

现在看来的话,还是在postpone mask的时候出现了一点问题需要处理,就是对于TLS连接,究竟哪里出现了不正常的现象。

基本测试流程
#

pkill -9 -f valkey-server; pkill -9 -f runtest # 首先杀掉原来的进程
make BUILD_TLS=yes -j$(nproc) # 带TLS进行编译
./utils/gen-test-certs.sh     # 生成相应的证书
./runtest --single tests/unit/io-threads.tcl --tls # 跑这个报错的单测

为了保证和test-ubuntu-latest-cmake-tls是完全一样的,我们的测试流程是:

rm -rf build-release
mkdir -p build-release
cd build-release
cmake -DCMAKE_BUILD_TYPE=Release .. -DBUILD_TLS=yes -DBUILD_UNIT_GTESTS=yes
make -j$(nproc)
cd ..

其实就是在TLS在accept阶段的时候,你不应该封死读路径(if判断修正一下即可),这样的话TLS的状态机没有办法往后更新了。 这样处理之后CI就全部正常通过了。

之后还收到了review:

processClientIOWriteDone(c)
  → connUpdateState(c->conn)
    → updateRdmaState()
      → connRdmaEventHandler()
        → read handler
          → readQueryFromClient()
            → ... parse/execute ...
              → beforeNextClient(c)
                → if (c->flag.close_asap) freeClient(c);  // c is freed
  → come back to processClientIOWriteDone,try to access c  // use-after-free

我们在这个写路径上,

也应当先查询一下ID,然后我们再判断是否存在,防止产生相同的UAF。

之后再ping一把viktor:

@zuiderkwast Gentle ping on this one when you have bandwidth — it's in the Valkey 9.1 board under **Needs Review**.

Quick recap since the last round of feedback:
- RDMA + IO threads re-entrancy fix via directional `CONN_POSTPONE_READ` / `CONN_POSTPONE_WRITE` (approved by @pizhenwei).
- `handleReadJobs` drains `cmd_queue` before resuming transport.
- Latest follow-ups: removed blanket RDMA postpone skips (950c32b), inlined the read-done postpone mask (5586e73), and addressed the write-path UAF with post-`connUpdateState()` re-lookup (f842d16).

CI is green on the latest commit. No rush at all — just bubbling it up in case it got buried. Thanks!

截至文章写到现在,还是没有merge,但是我们分析的比较清晰了。