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 → rdmaProcessPendingData(rdma.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的位置,就是当前的线程停止在了
语义层面的重构 #
我们分析一下有什么重构的手法?
引入这个问题的原因是,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
我们在这个写路径上,
之后再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,但是我们分析的比较清晰了。