内存与文件的统一写状态机
1. 它是什么,原理是什么?
输出队列中的元素统一表示为:
using OutputSegment = std::variant<MemorySegment, FileSegment>;
MemorySegment保存用户态缓冲区、offset、remaining。FileSegment保存UniqueFd、文件偏移、剩余长度。- 队列只允许队尾追加、队首发送。
状态机的核心规则是:
- 只处理队首,后面的 segment 不能越过队首。
- 内存段调用
send,文件段调用sendfile。 - 成功发送
n字节,就推进 offset,减少 remaining。 EINTR:原地重试同一个 segment。EAGAIN:保留当前 offset 和 remaining,打开EPOLLOUT,等待下一次可写事件。- remaining 变成 0:segment 退役、出队、归还内存或文件 FD 资源。
- 队列为空:关闭
EPOLLOUT。
因此,Memory → File → Memory 仍然严格按顺序输出。
需要注意,项目中“统一”的是队列模型、顺序、状态推进、错误处理和生命周期;发送系统调用本身仍然不同。
2. 项目里怎么用,为什么用?
静态文件请求的路径是:
FileHandler
→ openat2 打开文件
→ fstat 获取固定长度
→ HttpResponse::SetFile
→ ResponseBatch(Header + File)
→ Reservation 统一准入
→ OutputQueue
→ send / sendfile
动态响应的 body 会构造成 MemorySegment;静态文件则构造成 FileSegment。
使用它主要有四个原因:
- 保证 HTTP pipeline 顺序:响应头、文件正文、下一个动态响应不能乱序。
- 减少内存占用:文件内容不整体读入用户态缓冲,直接使用
sendfile。 - 统一处理慢客户端:文件和内存都能通过短写、
EAGAIN、EPOLLOUT续传。 - 统一资源回收:
UniqueFd在文件发送完成、连接关闭或异常时自动关闭 FD。
项目还特意把两个概念分开计数:
flow_backlog_bytes:内存和文件都计入,用于 HIGH/LOW 背压。accounted_output_payload_bytes:只计算用户态内存,用于内存预算。
文件虽然不占用户态 payload,但仍然占文件 FD 和队列 segment,所以还有单连接 FD 数和 segment 数上限。
pipeline 请求与背压解除后的异步续处理
1. 它是什么?原理是什么?
它是什么?原理是什么?
这里包含三个概念:
- HTTP pipeline(流水线请求):客户端在同一条 HTTP/1.1 连接上,不必等前一个响应回来,就连续发送多个请求。服务器需要正确区分请求边界,并按请求顺序返回响应。
- 背压:客户端读取响应较慢,服务器输出队列不断积压,就暂时减少上游输入和处理,避免继续产生大量响应。
- 解除背压后的异步续处理:输出积压降下来后,把“继续处理剩余输入”投递到事件循环,稍后执行。
注意,Keep-Alive 表示连接可以复用;pipeline 进一步允许前一个响应还没收到,就发送后续请求。它也不意味着服务器必须并行处理这些请求。
最关键的是区分两个位置:
内核 socket 接收缓冲区 → 用户态 _in_buffer → HTTP 解析与业务处理
关闭 EPOLLIN 只能停止继续从 socket 读取,已经进入 _in_buffer 的请求仍然存在。所以暂停时必须同时停止解析,恢复时也必须考虑这部分输入。
2. 我的项目里怎么用?为什么用?
我的服务器支持 HTTP pipeline,
Onmessage()会循环解析输入缓冲区里的完整请求。为了避免慢客户端导致响应持续积压,我用输出队列的待发送字节数控制背压,同时控制 socket 读取和 HTTP 解析。积压下降后,恢复读事件,并异步续处理已经缓冲的请求。
具体流程如下。
① 正常情况下,连续处理完整请求。
HttpServer::Onmessage() 在 while (buf->ReadAbleSize()) 中解析请求,只有请求完整后才进入路由、生成响应,再重置上下文处理下一个请求。响应通过 ResponseBatch 提交到统一 FIFO 输出队列,保持响应顺序。
② 输出积压达到 HIGH,同时暂停读取和解析。
项目当前设置:
| 条件 | 行为 |
|---|---|
| 积压达到或超过 4 MiB | 进入背压,关闭读事件,停止继续解析 |
| 已处于背压,积压仍大于 1 MiB | 保持暂停 |
| 积压降到 1 MiB 或以下 | 解除背压,恢复读取,安排输入续处理 |
这里统计的是 _flow_backlog_bytes,即输出队列里所有 segment 的剩余发送量,文件待发送部分也计入。
暂停时设置 _backpressured = true,关闭读事件;HTTP 解析循环通过 CanProcessInput() 检查这个状态,在当前响应触发 HIGH 后停止处理后续请求。
③ 写出数据后,检查是否可以恢复。
HandleWrite() 随实际发送量扣减积压,随后调用 ResumeReadInLoop()。满足 LOW 条件后:
- 清除背压状态。
- 在没有输入越限的情况下恢复读事件。
- 如果输入缓冲区还有数据,通过
QueueInLoop()安排续处理任务。
举个例子:
一次读入请求 A、B
↓
处理 A,产生大响应,输出积压达到 HIGH
↓
停止解析,B 留在用户态缓冲区
↓
客户端逐渐读取 A 的响应,输出积压下降到 LOW
↓
恢复读事件,并投递续处理任务
↓
任务执行,继续解析 B
这里“异步”指延后到同一个 owner EventLoop 的任务阶段执行,不是另开线程。 如果继续产生大响应,也可以再次进入背压。
为什么这样做:既控制慢客户端造成的积压,又保证已经收到的请求能够继续推进;同时避免写回调直接嵌套调用 HTTP 解析和业务处理。
EventLoop 与连接状态的线程归属
1. 它是什么,原理是什么?
EventLoop 本质是一个事件分发循环:
- 用
epoll_wait等待 socket、eventfd、timerfd的事件。 Poller返回活跃的Channel。- EventLoop 在所属线程中依次调用 Channel 的读、写、关闭、错误回调。
- 再执行跨线程投递到任务队列中的任务。
项目中的主循环是:
epoll_wait
-> Channel::HandEvent()
-> RunAllTask()
每个 EventLoop 在创建时记录线程 ID,并通过 AssertInLoop() 强制要求:
Channel/Poller/Timer/Connection状态
只能由该线程操作。
跨线程任务通过:
- 互斥锁保护任务队列;
eventfd唤醒阻塞中的epoll_wait;- owner 线程取出任务并执行。
因此它类似一个“单线程 Actor”:
外部线程可以发消息,但不能直接修改 Actor 内部状态。
2. 项目里怎么用,为什么用?
连接如何分配线程
TcpServer 有一个 BaseLoop:
- 负责监听 socket;
- 接收新连接;
- 管理连接表;
- 处理全局停止流程。
同时可以创建多个 Worker EventLoop。新连接到来时,通过线程池轮询选择一个 worker:
EventLoop *owner_loop = _pool.NextLoop();
conn.reset(new Connection(owner_loop, ...));
所以一个 Connection 从创建开始就固定绑定一个 _loop:
EventLoop *_loop;
连接的以下状态全部归 owner loop 管理:
_statu_channel_socekt_in_buffer_output_queue- 定时器
- 背压状态
- 请求 deadline
- 关闭状态
跨线程如何操作连接
例如业务线程调用:
conn->Send(...)
conn->Shutdown()
conn->EnableInactiveRelease(...)
这些接口不会直接改连接,而是进入 [`DispatchToOwner()` (line 2343)](source/server.hpp:2343):
- 如果当前已经是 owner 线程,直接执行;
- 如果是其他线程,封装成任务并
QueueInLoop(); - owner loop 再执行真正的
SendInLoop()、ShutdownInLoop()等逻辑。
Send() 还会先复制数据,保证调用方随后修改或释放原始缓冲区不会影响异步发送。
但是 SendResponseBatch() 明确要求在 owner loop 调用,因为它要保证整批响应在 owner loop 中一次性提交到 FIFO。
关闭为什么也必须回 owner loop
连接关闭不是简单 close(fd),还要完成:
- 从 epoll 删除 Channel;
- 关闭 socket;
- 取消 timer;
- 清空输出队列;
- 归还内存和 fd 预算;
- 调用关闭回调;
- 通知 BaseLoop 从连接表删除。
项目还特别把 Release() 延迟到下一轮任务执行:
_loop->QueueInLoop([self]() {
self->ReleaseInLoop();
});
因为当前 epoll_wait 返回的 active Channel 快照中,可能还保存着这个连接的 Channel*。如果当前回调里立即析构连接,后面的事件分发可能访问悬空指针。
为什么使用这种模型
主要原因有三个:
-
减少锁竞争
连接内部状态不需要每个字段都加锁,状态机由 owner loop 串行推进。 -
保证操作顺序
读、写、关闭、超时、发送任务都在同一条队列中线性化,输出队列天然保持 FIFO。 -
避免生命周期竞态
关闭、移除 epoll、关闭 fd、取消定时器都在同一个线程完成,降低 use-after-free 和重复 teardown 风险。
跨线程任务、shared_ptr 与连接所有权
1. 它是什么,原理是什么?
跨线程任务是把“不能在当前线程直接执行的操作”封装成闭包,放进目标线程的任务队列,由目标线程执行。
项目中的 EventLoop::QueueInLoop():
- 用 mutex 保护任务队列;
- 把任务放入
_tasks; - 写
eventfd唤醒可能阻塞在epoll_wait的线程; - 目标 EventLoop 被唤醒后执行任务。
RunInLoop() 如果发现当前已经是 owner 线程,就直接执行,否则进入队列。
shared_ptr 使用引用计数管理对象生命周期。任务闭包按值捕获 shared_ptr 后,即使外部的最后一个引用释放,任务执行前对象也不会析构。
这里要注意:
shared_ptr的引用计数管理是线程安全的;Connection内部字段并不会因此自动线程安全;- 项目通过 owner loop 保证连接状态只在一个线程修改;
weak_ptr只观察对象,不延长生命周期,适合 Timer 回调。
连接所有权可以分成三层:
TcpServer的连接表持有连接的强引用;Connection自己持有 socket、Channel、输入缓冲和输出队列;- 排队任务临时持有
shared_ptr<Connection>,保证任务执行时对象有效。
2. 项目里怎么用,为什么这样用?
连接继承了 std::enable_shared_from_this<Connection>,并定义了:
using PtrConnection = std::shared_ptr<Connection>;
连接建立后,服务器先配置回调和资源预算,再把连接放进连接表,最后异步执行 EstablishedInLoop()。
Established()、Send()、Shutdown()、EnableInactiveRelease() 等公共接口都可以被非 owner 线程调用,但它们最终都会经过 DispatchToOwner():
PtrConnection self = shared_from_this();
loop->QueueInLoop([self, command]() {
command(self);
});
例如 Send() 有两个生命周期保护:
- 先把调用方传入的数据复制到
std::string; - 再把
shared_ptr<Connection>和数据一起放进任务。
所以调用方即使马上复用或释放原始缓冲区,也不会影响异步发送。
关闭连接时,Release() 不直接释放,而是继续投递一个强持有连接的任务:
PtrConnection self = shared_from_this();
_loop->QueueInLoop([self]() {
self->ReleaseInLoop();
});
ReleaseInLoop() 在 owner loop 中完成:
- 设置
DISCONNECTED; - 关闭跨线程投递闸门;
- 从 Poller 移除 Channel;
- 关闭 socket;
- 取消 Timer;
- 清空输出队列;
- 归还内存和文件描述符配额;
- 执行 close callback。
服务器停止时,也不会直接从 BaseLoop 操作所有连接,而是按 owner loop 分组,给每个 Worker 投递强引用快照,等待所有 Worker 完成清理后再停止线程。
这样做的原因是:
- 避免多个线程同时修改同一个连接;
- 避免 Poller、Channel、Timer 被错误线程操作;
- 避免排队任务访问已经析构的对象;
- 避免停止过程中连接仍被 Worker 使用。
延迟回收与幂等关闭
1. 它是什么,原理是什么?
“延迟回收”解决的是事件循环中的生命周期问题。
项目的 Poller::Poll() 会先生成一个 active Channel* 快照,然后统一执行回调,最后才执行 RunAllTask()。如果在某个回调里直接析构 Connection,后续快照中仍可能存在该连接的 Channel*,就会出现悬空指针和 use-after-free。
项目采用两阶段关闭:
收到 EOF / 错误 / 超时 / Shutdown
↓
CONNECTED → DISCONNECTING
↓
立即停止读事件、关闭新的命令投递
↓
Release() 投递到 EventLoop 任务队列
↓
当前 active 事件分发结束
↓
ReleaseInLoop() 真正 teardown
Release() 会捕获 shared_ptr<Connection>,保证延迟任务执行前对象还活着;ReleaseInLoop() 才真正执行:
- 从 Poller 移除 Channel;
- 关闭 socket;
- 取消请求 deadline 和 inactive timer;
- 清空输出队列;
- 关闭 FileSegment 的 fd;
- 归还内存预算、文件 fd 计数和连接指标;
- 最后触发 close callback。
“幂等关闭”指的是:关闭请求可以来自多个路径,但最终 teardown 只能执行一次。
项目通过以下状态保证幂等:
CONNECTING → CONNECTED → DISCONNECTING → DISCONNECTED
ShutdownInLoop() 遇到 DISCONNECTING 或 DISCONNECTED 直接返回;ReleaseInLoop() 遇到 DISCONNECTED 也直接返回。因此 EOF、超时、写错误、业务主动关闭、服务停止同时发生时,只有第一次关闭真正生效。
2. 项目里怎么用,为什么用?
连接级别:
- HTTP 请求处理完且响应需要关闭时,调用
conn->Shutdown()。 - HTTP 错误响应先写入输出队列,再调用
Shutdown()。 - 对端 EOF 走优雅关闭,允许已经排队的响应继续发送。
- 读错误、写错误、请求 deadline 超时走硬关闭,直接丢弃未发送输出并延迟 teardown。
- inactive timer 和 request deadline 都只捕获
weak_ptr,避免 timer 反过来延长 Connection 生命周期。 - 跨线程调用
Shutdown()、Send()时,先通过DispatchToOwner()投递到连接所属的 EventLoop。
服务级别:
TcpServer::StopNow() 也采用同样思想:
- 用生命周期状态
Running → Stopping → Stopped抢占停止权; - 停止 Acceptor,禁止新连接进入;
- 按 owner EventLoop 对连接做快照;
- 把每个 worker 的连接 teardown 投递回各自线程;
- 用
TeardownBarrier等待所有 worker 完成; - 清空连接表,停止 worker,最后退出 BaseLoop。
这样做的原因有三个:
- Poller、TimerWheel、Channel、socket 都要求在 owner loop 操作;
- 当前事件快照中保存的是裸
Channel*; - 外部线程可能在连接关闭期间继续调用
Send()或Shutdown()。
所以项目的核心原则是:
状态变化先同步生效,资源回收延迟到安全点;所有真实 I/O 和 Poller 操作只在 owner loop 执行。
如果不延迟释放只采用智能指针,那么状态就不能及时修改
时间轮的绝对超时时间是为了防止固定长度的时间轮返回回来之后判断不出新旧时间事件
请求的绝对时间是由链接内变量控制,请求行和header是一个阶段,body一个阶段,每个阶段有每个阶段超时时间,这样防止一个链接不断发送短信息不到一个请求但始终占着链接。
停机屏障与断连 / 超时 / 停机竞态
1. 它是什么,原理是什么?
这个项目里,“停机屏障、断连、超时、停机竞态”可以统一理解为:先阻止新的工作进入,再让已有工作在正确的 owner loop 中完成清理,最后回收线程和对象。
这里的“停机屏障”不是 CPU 内存屏障,而是一个 teardown barrier,类似 CountDownLatch:
CONNECTED -> DISCONNECTING -> DISCONNECTED
CONNECTED:正常收发。DISCONNECTING:停止接收新请求,必要时继续冲刷已经提交的响应。DISCONNECTED:Channel、socket、timer、输出队列都已清理。
连接的所有状态只由所属的 EventLoop 修改。跨线程调用 Send()、Shutdown() 时,只是把任务投递到 owner loop。
服务器停机时:
- 将生命周期改为
Stopping,不再接受新连接; - 停止 Acceptor;
- 对
ConnectionMap做强引用快照,并按 owner loop 分组; - 每个 worker 在自己的 loop 中执行所有连接的
ReleaseInLoop(); - worker 完成后调用
barrier->Done(); - BaseLoop 等待所有 worker 完成;
- 清空连接表,停止并 Join worker;
- 设置为
Stopped,退出 BaseLoop。
超时分两类:
- 空闲超时:没有请求处理时启用,I/O 活动会刷新;
- 请求超时:Header 和 Body 分别使用从阶段开始计算的绝对 deadline,收到零碎字节不会无限续命。
时间轮负责大致调度,真正关闭前还会用 steady_clock 再确认一次绝对时间,避免提前关闭。
2. 项目里怎么用,为什么这样用?
连接断开有三条主要路径:
- 对端 EOF:走优雅关闭,停止读事件,但把已经进入
OutputQueue的响应发送完; - 读写错误、请求 deadline 到期:停止读并排队释放,挂起输出直接丢弃;
- HTTP 错误或业务主动关闭:先生成错误响应并提交,再调用
Shutdown()。
ShutdownInLoop() 会:
- 清除请求 deadline;
- 把连接改为
DISCONNECTING; - 移除读事件,禁止新输入;
- 如果有待发送数据,打开
EPOLLOUT; - 输出队列清空后进入
ReleaseInLoop()。
- 关闭跨线程投递闸门;
- 状态改为
DISCONNECTED; - 从 Poller 移除 Channel;
- 关闭 socket;
- 取消 timer;
- 清空输出队列;
- 归还内存、文件 FD、segment 配额;
- 最后执行 close callback。
这样设计的原因是:
epoll的 active 列表中保存的是Channel*,不能在回调中直接析构对象;- 每个连接的 Channel、Timer 和队列都必须由 owner loop 操作;
- 停机时必须确保 worker 不再访问 Connection,
Start()返回才是安全析构服务器的完成点; - 外部即使还持有
shared_ptr<Connection>,teardown 也会立即关闭 fd 和释放资源,而不是等 shared_ptr 析构。
DispatchToOwner() 还使用 _dispatch_mutex 保护“检查闸门并入队”这一小段逻辑。teardown 先关闭闸门,再允许 worker 被 Join,避免外部线程向已经销毁的 EventLoop 投递任务。
故障注入 / 真实 socket 回归 / ASan·UBSan / TSan
1️⃣ 是什么 / 原理
- 故障注入:主动制造异常(EINTR/EAGAIN/EMFILE、文件截短、慢消费者),专测错误路径——网络服务器 80% 的 bug 在这里。
- 真实 socket 测试:用 loopback TCP/socketpair 走真实内核栈测试。EAGAIN 时机、partial write、RST、EPOLLRDHUP 顺序这些内核行为 mock 不出来。
- 回归验证:分层 gate(make step34),同一套测试在多种构建下重复执行,失败立刻定位到哪个安全层。
- ASan:编译插桩 + 影子内存 + redzone/quarantine,查越界、UAF、double-free、泄漏(LSan)。
- UBSan:查未定义行为(有符号溢出、错误对齐、非法枚举转换等)。
- TSan:基于 happens-before 关系追踪锁/原子操作,报数据竞争。
- 分工:ASan/UBSan 查「内存/语义错误」(单线程也能查),TSan 查「并发访问模式」,互补不互替。
2️⃣ 项目里怎么用 / 为什么
构建矩阵(test/makefile,make step34 一键全跑):
Gate 配置 跑什么
step34-plain 无插桩 全部功能回归
step34-asan -O1 -g3 -fsanitize=address,undefined,halt_on_error=1 与 plain 完全相同的功能集
step34-tsan -fsanitize=thread 精选并发敏感子集:LoopThread、owner-loop 释放、时间轮、StopNow、跨线程 Send/Shutdown、全局 CAS 预算、并发 sendfile、metrics
step34-release -O2 -DNDEBUG 真实起服务 curl 验证动态/静态响应(防“只在 Debug 正确”)
典型故障注入(全部真实存在):
- EINTR: socketpair() 建的一对本地互连 fd,receiver 端等价于服务端读端,peer 不写就等于“客户端不发包”,recv() 阻塞。,辅助线程向主线程发送信号打断, Socket::Recv() 把 errno==EINTR 翻译成 RecvStatus::Interrupted 返回,测试再用 EXPECT 断言这个 status 和 error_code==EINTR 是否成立——成立就证明封装层正确识别了“被信号打断”,且当普通错误处理、也没死循环重试。
- EMFILE: 先获取下一个可用fd用 setrlimit 压低到这个fd, 强制 fd 耗尽,监听套接字返回fd耗尽的状态然后验证无 busy loop(1.2秒内重试的次数大约是五次不是无限,进程实际消耗cpu时间不会是1.2秒)、listenfd 不关、250ms 冷却后恢复,解除fd限制后能正常获取fd。
- 文件中途截短:传输中 ftruncate,验证 sendfile 提前 EOF 判硬错关闭、且尾部排队的动态数据绝不越过文件错发,线上行为(字节/EOF)由 peer 验证;服务端内部状态(队列/FD/不变量)由测试以进程内白盒方式直接验证,不依赖客户端“知道”什么。
- 慢消费者/背压:调小 SO_SNDBUF + 对端不读 → 强制 EAGAIN、partial write,触发 HIGH/LOW 水位(backpressure_test.cc:129)。
ASan:内存错误:越界、Double-Free、栈/全局/堆溢出、内存泄漏
ASan 是怎么发现越界和 UAF 的?
核心是影子内存:每 8 字节应用内存对应 1 字节影子字节,映射公式为 shadow(A) = ShadowBase + (A >> 3)。影子字节取值含义:0 表示这 8 字节全部可访问;1~7 表示前 N 字节有效(用于处理非 8 对齐的尾部);负值表示中毒。编译器为每一次 load/store 插入形如 __asan_load4(addr) 的调用,运行时先查影子字节,发现中毒立即报错
free 之后内存不会立即被复用,进入中毒状态
TSan 判定 data race 的条件是什么?
四个条件同时成立才报:① 同一内存位置;② 来自不同线程;③ 至少一个是写;④ 这两次访问之间没有 happens-before 关系。注意第四条——结果"碰巧对了"也照样报,因为它看的是逻辑顺序而非实际值。sqglobe.com+1
Q2:什么是 happens-before?哪些操作会建立 HB 边?
HB 是并发操作之间的偏序关系,由同步原语建立。TSan 能识别的 HB 边包括:mutex 的 unlock→lock、条件变量的 wait/notify、线程 create/join、正确的 acquire/release 原子操作(含 fence)。
UBSan 检测哪些 UB?
典型包括:有符号整数溢出(INT_MAX + 1)、越界移位(1 << 31 对 int 是 UB)、空指针解引用(通过 null 检查项)、未对齐访问、数组下标越界(静态可确定时)、INT_MIN / -1、浮点转换溢出、无效枚举值等。llvm.org+1
Q2:UBSan 和 ASan 的本质区别?
UBSan 处理的是编译器在生成代码时假设"永远不发生"的 UB,检查逻辑在编译期就能确定并插入条件判断;ASan 处理的是运行时才能确定的内存访问合法性,需要影子内存这种运行时数据结构支撑。所以 UBSan 开销极小,而 ASan 开销大。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/qq_55640460/article/details/166781868



