云计算百科
云计算领域专业知识百科平台

11.升级简单的Epoll,从 select 到 epoll:C++17 TCP 服务器的六层渐进式重构

文章目录

  • 从 select 到 epoll:C++17 TCP 服务器的六层渐进式重构
    • 一、为什么要采用“分层升级”
    • 二、先建立正确的 epoll 心智模型
      • 2.1 epoll 不是异步 I/O
      • 2.2 interest list 与 ready list
      • 2.3 epoll 并不是“完全零拷贝”
      • 2.4 不要把 epoll 的复杂度概括成“所有操作都是 O(1)”
  • 第一部分:第 0 层——原始 select 版本
    • 三、第 0 层的目的
    • 3.1 完整服务端:`00_select_original/server.cpp`
    • 3.2 完整客户端:`00_select_original/client.cpp`
    • 3.3 原始版本的具体问题
      • 问题一:每轮复制并修改 fd_set
      • 问题二:返回后仍然遍历全部客户端
      • 问题三:受到 fd_set 表示能力限制
      • 问题四:把 EOF 和错误混在一起
      • 问题五:默认阻塞 I/O
      • 问题六:没有处理部分发送
      • 问题七:TCP 没有消息边界
  • 第二部分:第 1 层——最小 epoll LT 迁移
    • 四、本层目标
    • 4.1 从 select 到 epoll 的最小映射
    • 4.2 完整服务端:`01_epoll_lt_minimal/server.cpp`
    • 4.3 完整客户端:`01_epoll_lt_minimal/client.cpp`
    • 4.4 本层修复了什么
      • 修复一:不再每轮复制完整监控集合
      • 修复二:不再遍历所有空闲客户端
      • 修复三:不再依赖 `max_fd + 1`
    • 4.5 本层仍然存在的问题
  • 第三部分:第 2 层——非阻塞 epoll LT
    • 五、为什么 LT 也建议使用非阻塞套接字
    • 5.1 本层新增的规则
      • 规则一:监听 fd 非阻塞
      • 规则二:通信 fd 也设为非阻塞
      • 规则三:循环读取到 EAGAIN
      • 规则四:本层遇到发送缓冲区满时直接关闭连接
    • 5.2 完整服务端:`02_epoll_lt_nonblocking/server.cpp`
    • 5.3 完整客户端:`02_epoll_lt_nonblocking/client.cpp`
    • 5.4 本层修复了什么
    • 5.5 为什么不能把 EAGAIN 当成错误
    • 5.6 本层仍然存在的问题
  • 第四部分:LT 与 ET 的真正区别
    • 六、水平触发 LT
    • 七、边缘触发 ET
  • 第五部分:ET 的两个典型错误
    • 八、错误一:EPOLLET + 单次 read
      • 8.1 错误逻辑
    • 8.2 完整错误复现代码
    • 九、错误二:EPOLLET + 阻塞套接字 + 循环 read
    • 9.1 完整错误复现代码
    • 十、ET 正确读法
  • 第六部分:第 3 层——正确的 epoll ET 读取
    • 十一、本层相对第 2 层改了什么
    • 11.1 完整服务端:`03_epoll_et_correct_read/server.cpp`
    • 11.2 完整客户端:`03_epoll_et_correct_read/client.cpp`
    • 11.3 本层修复了什么
    • 11.4 本层仍然没有解决写路径
  • 第七部分:为什么不能永久监听 EPOLLOUT
    • 十二、套接字“可写”意味着什么
    • 十三、为什么每个连接都需要独立发送状态
  • 第八部分:第 4 层——ET + 发送队列 + EPOLLOUT
    • 十四、本层新增的状态机
    • 14.1 完整服务端:`04_epoll_et_write_queue/server.cpp`
    • 14.2 完整客户端:`04_epoll_et_write_queue/client.cpp`
    • 14.3 `BuildClientEvents()` 的意义
    • 14.4 为什么收到数据后先主动发送
    • 14.5 `EPOLLRDHUP` 不能直接等同于“所有数据已经读完”
    • 14.6 本层修复了什么
    • 14.7 本层的新风险:内存无限增长
  • 第九部分:背压是什么
    • 十五、背压的核心含义
    • 十六、高低水位
    • 十七、为什么还需要硬上限
  • 第十部分:ET、公平预算和 EPOLLONESHOT
    • 十八、为什么不能无限排空一个非常活跃的连接
    • 十九、ET 下提前停止读取为什么危险
    • 二十、EPOLLONESHOT 的语义
  • 第十一部分:第 5 层——背压、预算与资源边界
    • 二十一、最终层的状态字段
    • 21.1 完整服务端:`05_epoll_et_backpressure/server.cpp`
    • 21.2 完整客户端:`05_epoll_et_backpressure/client.cpp`
    • 21.3 最终层的保护机制
      • 连接数上限
      • 单连接发送队列硬上限
      • 高低水位背压
      • 每事件读写预算
      • EPOLLONESHOT 重装
    • 21.4 最终层完整流程图
  • 第十二部分:epoll 三大 API 逐项解析
    • 二十二、`epoll_create1()`
    • 二十三、`epoll_ctl()`
      • ADD
      • MOD
      • DEL
    • 二十四、`epoll_wait()`
  • 第十三部分:常用事件标志
    • 二十五、事件位速查表
    • 二十六、为什么使用 SO_ERROR
  • 第十四部分:accept、recv、send 的完整状态表
    • 二十七、非阻塞 accept
    • 二十八、非阻塞 recv
    • 二十九、非阻塞 send
  • 第十五部分:TCP 半关闭
    • 三十、`shutdown(SHUT_WR)` 的意义
  • 第十六部分:LT 与 ET 应该如何选
    • 三十一、对比表
    • 三十二、ET 一定比 LT 快吗
  • 第十七部分:压力测试客户端
    • 三十三、完整大数据校验客户端
    • 33.1 测试 1 MiB 回显
    • 33.2 为什么不能用一次 write 作为压力测试
  • 第十八部分:编译与运行
    • 三十四、编译全部版本
    • 三十五、运行指定层次
      • 第 0 层
      • 第 1 层
      • 第 2 层
      • 第 3 层
      • 第 4 层
      • 第 5 层
  • 第十九部分:调试方法
    • 三十六、使用 strace 观察事件循环
    • 三十七、使用 ss 查看连接
    • 三十八、模拟慢客户端
  • 第二十部分:常见错误清单
    • 三十九、错误:ET 只加 EPOLLET,不改读取方式
    • 四十、错误:ET 使用阻塞 fd 循环读
    • 四十一、错误:把 EAGAIN 当成连接失败
    • 四十二、错误:永久监听 EPOLLOUT
    • 四十三、错误:忽略部分发送
    • 四十四、错误:发送队列没有上限
    • 四十五、错误:ET 设置处理预算但不重新激活
    • 四十六、错误:EPOLLRDHUP 立即 close
    • 四十七、错误:在 `epoll_event.data.ptr` 中保存裸指针却不管理生命周期
  • 第二十一部分:进一步工程化方向
    • 四十八、应用层协议拆包
    • 四十九、定时器
    • 五十、信号
    • 五十一、多线程
    • 五十二、跨线程唤醒
  • 第二十二部分:FAQ
    • 五十三、close 前必须 EPOLL_CTL_DEL 吗
    • 五十四、accepted socket 会继承 O_NONBLOCK 吗
    • 五十五、EAGAIN 和 EWOULDBLOCK 是否相同
    • 五十六、LT 可以使用阻塞套接字吗
    • 五十七、ET 是否每次都必须读到 EAGAIN
    • 五十八、epoll 可以监控普通文件吗
    • 五十九、epoll 能替代线程池吗
  • 第二十三部分:六层演进总结
    • 六十、每一层到底学到了什么
      • 第 0 层:select
      • 第 1 层:epoll LT 最小迁移
      • 第 2 层:非阻塞 LT
      • 第 3 层:ET 正确读取
      • 第 4 层:写事件状态机
      • 第 5 层:工程边界
    • 六十一、最终核心原则
  • 附录 A:源码目录
  • 附录 B:官方资料名称

从 select 到 epoll:C++17 TCP 服务器的六层渐进式重构

平台:Linux 语言标准:C++17 示例类型:单线程 TCP Echo 服务器 核心目标:不是一次性给出一份复杂的“最终代码”,而是从原始 select 程序开始,一层一层升级到具备非阻塞 I/O、ET 正确读写、发送队列、背压和公平性控制的 epoll 事件循环。


一、为什么要采用“分层升级”

很多 epoll 教程存在同一个问题:

  • 先展示一份非常简单的 select 代码;
  • 下一页直接出现 EPOLLET、非阻塞、unordered_map、发送缓冲区、EPOLLOUT、EPOLLONESHOT;
  • 读者能看懂每一行语法,却无法理解每一处状态为什么存在。
  • 本文不采用这种跳跃式讲法,而是固定为六个可以独立编译、独立运行的版本:

    层次版本本层只解决的核心问题
    第 0 层 原始 select 建立对照基线
    第 1 层 最小 epoll LT 把“扫描全部 fd”改成“只返回就绪事件”
    第 2 层 非阻塞 epoll LT 排空连接队列、区分 EAGAIN、避免单个 I/O 阻塞事件循环
    第 3 层 正确的 epoll ET 读取 非阻塞套接字 + 循环处理到 EAGAIN
    第 4 层 ET 发送队列 处理部分发送、发送缓冲区满和 EPOLLOUT
    第 5 层 背压与公平性 限制内存、暂停读取、处理预算、EPOLLONESHOT 重装

    每一层都包含:

    • 完整服务端代码;
    • 完整客户端代码;
    • 相对上一层改动了什么;
    • 修复了什么;
    • 本层仍然没有解决什么;
    • 编译与运行方法。

    二、先建立正确的 epoll 心智模型

    2.1 epoll 不是异步 I/O

    epoll 是一种 I/O 就绪事件通知机制。

    它告诉应用程序:

    某个文件描述符现在可能可以执行读取、写入或其他操作。

    它不会替应用程序执行 recv() 或 send(),也不会自动创建线程。

    典型执行过程仍然是同步的:

    应用线程

    │ epoll_wait()

    等待一个或多个 fd 就绪

    │ 返回就绪事件

    应用自己调用 recv()/send()/accept()


    处理业务状态

    └──────────────→ 再次 epoll_wait()

    所以更准确的描述是:

    epoll 负责“告诉你谁可能能做 I/O”,应用程序仍负责真正的 I/O 和状态管理。

    2.2 interest list 与 ready list

    从用户态视角,可以把一个 epoll 实例理解成两部分:

    epoll 实例
    ┌─────────────────────────────┐
    │ │
    │ interest list │
    │ 我关心哪些 fd、哪些事件 │
    │ │
    │ fd=3: EPOLLIN │
    │ fd=7: EPOLLIN|EPOLLRDHUP │
    │ fd=9: EPOLLIN|EPOLLOUT │
    │ │
    ├─────────────────────────────┤
    │ │
    │ ready list │
    │ 当前已经就绪的监控项 │
    │ │
    │ fd=7: EPOLLIN │
    │ fd=9: EPOLLOUT │
    │ │
    └─────────────────────────────┘

    三个核心 API 的角色分别是:

    epoll_create1()
    创建 epoll 实例

    epoll_ctl()
    修改 interest list
    ADD / MOD / DEL

    epoll_wait()
    从 ready list 取出当前就绪事件

    2.3 epoll 并不是“完全零拷贝”

    select() 每轮都要把待监控的 fd_set 传给内核,并在返回时把修改后的集合带回用户态。

    epoll 的优势是:

    • 监控集合通过 epoll_ctl() 持久保存在内核中;
    • epoll_wait() 不需要每轮重新提交整个监控集合;
    • 返回时只复制本轮就绪的 epoll_event。

    但 epoll_wait() 返回事件时仍然需要把就绪事件复制到用户缓冲区,所以不应该把 epoll 简化成“零拷贝”。

    2.4 不要把 epoll 的复杂度概括成“所有操作都是 O(1)”

    更稳妥的理解是:

    • epoll_ctl() 要查询并修改内核中的监控项;
    • epoll_wait() 的主要用户态处理量与返回的就绪事件数量相关;
    • 应用不再像 select() 那样线性检查每一个可能的 fd;
    • 内核实现细节可能随版本演化,不应把某一种内部数据结构当作稳定 ABI。

    第一部分:第 0 层——原始 select 版本

    三、第 0 层的目的

    这一层不追求健壮性,只保留最初代码的结构,用于回答:

    从 select 迁移到 epoll 时,最先发生变化的到底是什么?

    原始版本的主要结构是:

    master_fds 保存全部监控对象


    read_fds = master_fds


    select() 修改 read_fds


    FD_ISSET(listen_fd)

    ├─ accept()

    └─ 遍历 clients

    └─ FD_ISSET(client_fd)

    3.1 完整服务端:00_select_original/server.cpp

    #include <arpa/inet.h>
    #include <sys/select.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <algorithm>
    #include <cstring>
    #include <iostream>
    #include <string>
    #include <vector>

    using namespace std;

    int main() {
    // 1. 创建 TCP 监听套接字
    int fd = socket(AF_INET, SOCK_STREAM, 0);
    if (fd == 1) {
    cerr << "socket() failed" << endl;
    return 1;
    }

    // 2. 设置端口复用
    int opt = 1;
    setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt));

    // 3. 绑定地址并监听
    sockaddr_in addr{};
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = htonl(INADDR_ANY);
    addr.sin_port = htons(8888);

    if (bind(fd, reinterpret_cast<sockaddr*>(&addr), sizeof(addr)) == 1 ||
    listen(fd, 5) == 1) {
    cerr << "bind/listen failed" << endl;
    close(fd);
    return 1;
    }

    cout << "Server listening on port 8888…" << endl;

    // 4. 初始化 select 监控集合
    fd_set master_fds;
    fd_set read_fds;
    FD_ZERO(&master_fds);
    FD_SET(fd, &master_fds);
    int max_fd = fd;

    vector<int> clients;

    while (true) {
    // select 会修改 fd_set,因此每轮必须复制。
    read_fds = master_fds;

    int activity =
    select(max_fd + 1, &read_fds, nullptr, nullptr, nullptr);
    if (activity < 0) {
    perror("select");
    break;
    }

    // 5. 监听套接字可读:接收一个新连接
    if (FD_ISSET(fd, &read_fds)) {
    sockaddr_in client_addr{};
    socklen_t len = sizeof(client_addr);
    int cfd = accept(
    fd, reinterpret_cast<sockaddr*>(&client_addr), &len);

    if (cfd != 1) {
    char ip_str[INET_ADDRSTRLEN]{};
    inet_ntop(AF_INET,
    &client_addr.sin_addr,
    ip_str,
    sizeof(ip_str));

    cout << "New connection from "
    << ip_str << ':'
    << ntohs(client_addr.sin_port)
    << endl;

    FD_SET(cfd, &master_fds);
    max_fd = max(max_fd, cfd);
    clients.push_back(cfd);
    }
    }

    // 6. 线性遍历所有客户端,检查哪些 fd 可读
    for (auto it = clients.begin(); it != clients.end();) {
    int cfd = *it;

    if (FD_ISSET(cfd, &read_fds)) {
    string buf;
    buf.resize(64);

    ssize_t n = read(cfd, buf.data(), buf.size());

    if (n <= 0) {
    cout << "Client disconnected (fd="
    << cfd << ')' << endl;

    close(cfd);
    FD_CLR(cfd, &master_fds);
    it = clients.erase(it);
    } else {
    buf.resize(static_cast<size_t>(n));
    cout << "Received: " << buf << endl;

    const string reply = "Hello from server";
    write(cfd, reply.c_str(), reply.size());
    ++it;
    }
    } else {
    ++it;
    }
    }
    }

    close(fd);
    return 0;
    }

    3.2 完整客户端:00_select_original/client.cpp

    #include <arpa/inet.h>
    #include <cstring>
    #include <iostream>
    #include <sys/socket.h>
    #include <unistd.h>

    using namespace std;

    int main() {
    // 1. 创建 TCP 套接字
    int fd = socket(AF_INET, SOCK_STREAM, 0);

    // 2. 设置服务器地址
    sockaddr_in addr{};
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = inet_addr("127.0.0.1");
    addr.sin_port = htons(8888);

    // 3. 连接服务器
    connect(fd, reinterpret_cast<sockaddr*>(&addr), sizeof(addr));
    cout << "Connected to server" << endl;

    // 4. 发送一条消息
    const char* msg = "Hello from client";
    write(fd, msg, strlen(msg));
    cout << "Sent: " << msg << endl;

    // 5. 阻塞等待服务器回复
    char buf[64] = {};
    read(fd, buf, sizeof(buf));
    cout << "Received: " << buf << endl;

    close(fd);
    return 0;
    }

    3.3 原始版本的具体问题

    问题一:每轮复制并修改 fd_set

    read_fds = master_fds;
    select(max_fd + 1, &read_fds, nullptr, nullptr, nullptr);

    select() 会原地修改 read_fds,所以必须保留 master_fds,并在每轮重新复制。

    问题二:返回后仍然遍历全部客户端

    即使只有一个客户端有数据,应用仍然会检查 clients 中的每一个 fd:

    for (auto it = clients.begin(); it != clients.end(); ++it) {
    if (FD_ISSET(*it, &read_fds)) {
    // …
    }
    }

    问题三:受到 fd_set 表示能力限制

    在常见 Linux/glibc 环境中,fd_set 的容量由 FD_SETSIZE 约束。文件描述符数值一旦超出可表示范围,继续使用 FD_SET() 会导致未定义行为。

    问题四:把 EOF 和错误混在一起

    if (n <= 0) {
    // …
    }

    这里没有区分:

    n == 0 对端有序关闭写方向,接收端到达 EOF
    n == -1 系统调用失败,需要检查 errno

    问题五:默认阻塞 I/O

    即使 select() 报告套接字可读,后续操作仍可能因为竞态、错误的循环方式或不完整的状态判断而阻塞。

    问题六:没有处理部分发送

    write(cfd, reply.data(), reply.size());

    TCP write()/send() 成功返回时,返回值可能小于请求发送长度。不能把“一次调用成功”理解为“全部数据已经发送”。

    问题七:TCP 没有消息边界

    一次 write() 不保证对应对端的一次 read():

    发送端:
    write("ABC")
    write("DEF")

    接收端可能看到:
    read() -> "ABCDEF"

    也可能看到:
    read() -> "AB"
    read() -> "CDE"
    read() -> "F"


    第二部分:第 1 层——最小 epoll LT 迁移

    四、本层目标

    这一层只完成一件事:

    用 epoll_create1()、epoll_ctl()、epoll_wait() 替换 fd_set 和 select()。

    为了让差异足够清楚,本层刻意保留:

    • 阻塞套接字;
    • 每次只 accept() 一个连接;
    • 每次只 read() 一次;
    • 单次 write();
    • 没有发送队列。

    4.1 从 select 到 epoll 的最小映射

    selectepoll
    FD_ZERO() epoll_create1()
    FD_SET(fd) epoll_ctl(ADD, fd)
    修改监控事件 epoll_ctl(MOD, fd)
    FD_CLR(fd) epoll_ctl(DEL, fd)
    select() epoll_wait()
    遍历全部客户端 只遍历返回的就绪事件

    核心循环从:

    read_fds = master_fds;
    select(...);

    for (int fd : clients) {
    if (FD_ISSET(fd, &read_fds)) {
    // …
    }
    }

    变成:

    int count = epoll_wait(...);

    for (int i = 0; i < count; ++i) {
    int fd = events[i].data.fd;
    // fd 已经是本轮就绪对象
    }

    4.2 完整服务端:01_epoll_lt_minimal/server.cpp

    #include <arpa/inet.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <cstring>
    #include <iostream>
    #include <string>
    #include <unordered_set>
    #include <vector>

    namespace {
    constexpr uint16_t kPort = 8888;
    constexpr int kBacklog = 5;
    constexpr int kMaxEvents = 64;
    constexpr std::size_t kBufferSize = 64;

    bool AddReadEvent(int epoll_fd, int fd) {
    epoll_event event{};
    event.events = EPOLLIN; // 未指定 EPOLLET,默认是 LT
    event.data.fd = fd;

    if (epoll_ctl(epoll_fd, EPOLL_CTL_ADD, fd, &event) == 1) {
    std::perror("epoll_ctl(ADD)");
    return false;
    }
    return true;
    }

    void CloseClient(int epoll_fd,
    int client_fd,
    std::unordered_set<int>& clients) {
    // 对于单线程简单程序,close() 最终也会让内核移除监控项;
    // 显式 DEL 可以让用户态生命周期更加清晰。
    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_DEL,
    client_fd,
    nullptr) == 1 &&
    errno != ENOENT &&
    errno != EBADF) {
    std::perror("epoll_ctl(DEL)");
    }

    close(client_fd);
    clients.erase(client_fd);
    }
    }

    int main() {
    int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (listen_fd == 1) {
    std::perror("socket");
    return 1;
    }

    int reuse = 1;
    if (setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse)) == 1) {
    std::perror("setsockopt");
    close(listen_fd);
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
    server_addr.sin_port = htons(kPort);

    if (bind(listen_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("bind");
    close(listen_fd);
    return 1;
    }

    if (listen(listen_fd, kBacklog) == 1) {
    std::perror("listen");
    close(listen_fd);
    return 1;
    }

    // epoll 实例本身也是一个文件描述符。
    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
    if (epoll_fd == 1) {
    std::perror("epoll_create1");
    close(listen_fd);
    return 1;
    }

    if (!AddReadEvent(epoll_fd, listen_fd)) {
    close(epoll_fd);
    close(listen_fd);
    return 1;
    }

    std::vector<epoll_event> ready_events(kMaxEvents);
    std::unordered_set<int> clients;

    std::cout << "epoll LT minimal server listening on 0.0.0.0:"
    << kPort << '\\n';

    while (true) {
    int ready_count =
    epoll_wait(epoll_fd,
    ready_events.data(),
    static_cast<int>(ready_events.size()),
    1);

    if (ready_count == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("epoll_wait");
    break;
    }

    // epoll_wait 只返回本轮已经就绪的事件,不再遍历全部连接。
    for (int i = 0; i < ready_count; ++i) {
    const int current_fd = ready_events[i].data.fd;
    const uint32_t events = ready_events[i].events;

    if (current_fd == listen_fd) {
    // 为了保持“最小迁移”,本层仍然只 accept 一次。
    sockaddr_in client_addr{};
    socklen_t client_len = sizeof(client_addr);

    int client_fd =
    accept(listen_fd,
    reinterpret_cast<sockaddr*>(&client_addr),
    &client_len);

    if (client_fd == 1) {
    std::perror("accept");
    continue;
    }

    char ip[INET_ADDRSTRLEN]{};
    inet_ntop(AF_INET,
    &client_addr.sin_addr,
    ip,
    sizeof(ip));

    std::cout << "connected: "
    << ip << ':'
    << ntohs(client_addr.sin_port)
    << ", fd=" << client_fd << '\\n';

    if (!AddReadEvent(epoll_fd, client_fd)) {
    close(client_fd);
    continue;
    }

    clients.insert(client_fd);
    continue;
    }

    bool close_now = false;

    if ((events & EPOLLIN) != 0U) {
    char buffer[kBufferSize];
    ssize_t n = read(current_fd, buffer, sizeof(buffer));

    if (n > 0) {
    std::cout << "received from fd="
    << current_fd << ": ";
    std::cout.write(buffer, n);
    std::cout << '\\n';

    // 本层刻意保留“单次 write”的旧结构。
    ssize_t sent =
    write(current_fd,
    buffer,
    static_cast<std::size_t>(n));

    if (sent == 1) {
    std::perror("write");
    close_now = true;
    }
    } else if (n == 0) {
    close_now = true;
    } else {
    std::perror("read");
    close_now = true;
    }
    }

    if ((events & (EPOLLERR | EPOLLHUP)) != 0U) {
    close_now = true;
    }

    if (close_now) {
    std::cout << "closed fd=" << current_fd << '\\n';
    CloseClient(epoll_fd, current_fd, clients);
    }
    }
    }

    for (int client_fd : clients) {
    close(client_fd);
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    4.3 完整客户端:01_epoll_lt_minimal/client.cpp

    #include <arpa/inet.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <cstring>
    #include <iostream>
    #include <string>

    namespace {
    bool SendAll(int fd, const char* data, std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (socket_fd == 1) {
    std::perror("socket");
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(8888);

    if (inet_pton(AF_INET,
    "127.0.0.1",
    &server_addr.sin_addr) != 1) {
    std::cerr << "invalid server address\\n";
    close(socket_fd);
    return 1;
    }

    if (connect(socket_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("connect");
    close(socket_fd);
    return 1;
    }

    const std::string message = "Hello from epoll LT client";

    if (!SendAll(socket_fd, message.data(), message.size())) {
    close(socket_fd);
    return 1;
    }

    std::cout << "sent: " << message << '\\n';

    char buffer[256];
    ssize_t n;

    do {
    n = recv(socket_fd, buffer, sizeof(buffer), 0);
    } while (n == 1 && errno == EINTR);

    if (n > 0) {
    std::cout << "received: ";
    std::cout.write(buffer, n);
    std::cout << '\\n';
    } else if (n == 0) {
    std::cout << "server closed connection\\n";
    } else {
    std::perror("recv");
    }

    close(socket_fd);
    return n < 0 ? 1 : 0;
    }

    4.4 本层修复了什么

    修复一:不再每轮复制完整监控集合

    监控项通过 epoll_ctl() 加入 epoll 实例,之后由内核维护。

    修复二:不再遍历所有空闲客户端

    epoll_wait() 直接返回本轮就绪事件数组:

    for (int i = 0; i < ready_count; ++i) {
    int current_fd = ready_events[i].data.fd;
    }

    修复三:不再依赖 max_fd + 1

    epoll_wait() 不需要最大 fd 参数。

    4.5 本层仍然存在的问题

    • 监听套接字仍然阻塞;
    • 每次监听事件只 accept() 一个连接;
    • 通信套接字仍然阻塞;
    • 没有正确区分 EINTR、EAGAIN;
    • 单次 read() 可能无法读完当前数据;
    • 单次 write() 可能部分发送;
    • 没有处理慢客户端;
    • 没有背压。

    因此本层只是“API 迁移版”,不是最终版本。


    第三部分:第 2 层——非阻塞 epoll LT

    五、为什么 LT 也建议使用非阻塞套接字

    LT 模式的语义和 poll() 接近:只要条件仍然满足,就可能继续报告就绪。

    这并不意味着阻塞套接字一定安全。

    例如:

    epoll_wait 报告 fd 可读


    另一个执行单元先读走数据


    当前线程调用阻塞 recv()


    可能睡眠等待新数据

    在本文的单线程示例里这种竞态较少,但统一采用非阻塞 I/O 可以建立更稳定的事件循环不变量:

    任何可能等待的系统调用,都必须以 EAGAIN/EWOULDBLOCK 返回控制权,而不是阻塞整个事件循环。

    5.1 本层新增的规则

    规则一:监听 fd 非阻塞

    这样才能安全地循环 accept():

    while (true) {
    int cfd = accept(...);

    if (cfd >= 0) {
    // 继续 accept
    } else if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }
    }

    规则二:通信 fd 也设为非阻塞

    Linux 上,accept() 得到的新套接字不会因为监听套接字设置了 O_NONBLOCK 就自动继承该状态,因此必须单独设置,或者使用带 SOCK_NONBLOCK 的 accept4()。

    规则三:循环读取到 EAGAIN

    LT 并不强制必须这样做,但排空当前接收缓冲区可以减少重复进入 epoll_wait()。

    规则四:本层遇到发送缓冲区满时直接关闭连接

    这是一个有意保留的中间策略:

    好处:不会阻塞事件循环,也不会悄悄丢弃未发送数据
    代价:慢客户端会被断开

    真正保留慢客户端连接,需要在第 4 层引入发送队列。

    5.2 完整服务端:02_epoll_lt_nonblocking/server.cpp

    #include <arpa/inet.h>
    #include <fcntl.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>
    #include <unordered_set>
    #include <vector>

    namespace {
    constexpr uint16_t kPort = 8888;
    constexpr int kBacklog = 128;
    constexpr int kMaxEvents = 64;
    constexpr std::size_t kBufferSize = 4096;

    bool SetNonBlocking(int fd) {
    int old_flags = fcntl(fd, F_GETFL, 0);
    if (old_flags == 1) {
    return false;
    }

    return fcntl(fd,
    F_SETFL,
    old_flags | O_NONBLOCK) != 1;
    }

    bool AddClient(int epoll_fd, int client_fd) {
    epoll_event event{};
    event.events = EPOLLIN | EPOLLRDHUP; // 默认 LT
    event.data.fd = client_fd;

    return epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &event) != 1;
    }

    void CloseClient(int epoll_fd,
    int client_fd,
    std::unordered_set<int>& clients) {
    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_DEL,
    client_fd,
    nullptr) == 1 &&
    errno != ENOENT &&
    errno != EBADF) {
    std::perror("epoll_ctl(DEL)");
    }

    close(client_fd);
    clients.erase(client_fd);
    std::cout << "closed fd=" << client_fd << '\\n';
    }

    // 本层尚未引入发送队列。
    // 为避免非阻塞 send() 阻塞事件循环,如果不能一次发送完,直接关闭连接。
    bool SendImmediatelyOrFail(int fd,
    const char* data,
    std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    if (n == 1 &&
    (errno == EAGAIN || errno == EWOULDBLOCK)) {
    std::cerr
    << "send buffer full; layer 02 has no output queue\\n";
    return false;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int listen_fd = socket(AF_INET,
    SOCK_STREAM | SOCK_CLOEXEC,
    0);

    if (listen_fd == 1) {
    std::perror("socket");
    return 1;
    }

    if (!SetNonBlocking(listen_fd)) {
    std::perror("fcntl(listen_fd)");
    close(listen_fd);
    return 1;
    }

    int reuse = 1;
    if (setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse)) == 1) {
    std::perror("setsockopt");
    close(listen_fd);
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
    server_addr.sin_port = htons(kPort);

    if (bind(listen_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("bind");
    close(listen_fd);
    return 1;
    }

    if (listen(listen_fd, kBacklog) == 1) {
    std::perror("listen");
    close(listen_fd);
    return 1;
    }

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
    if (epoll_fd == 1) {
    std::perror("epoll_create1");
    close(listen_fd);
    return 1;
    }

    epoll_event listen_event{};
    listen_event.events = EPOLLIN; // 默认 LT
    listen_event.data.fd = listen_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    listen_fd,
    &listen_event) == 1) {
    std::perror("epoll_ctl(ADD listen_fd)");
    close(epoll_fd);
    close(listen_fd);
    return 1;
    }

    std::vector<epoll_event> ready_events(kMaxEvents);
    std::unordered_set<int> clients;

    std::cout << "epoll LT nonblocking server listening on 0.0.0.0:"
    << kPort << '\\n';

    while (true) {
    int ready_count =
    epoll_wait(epoll_fd,
    ready_events.data(),
    static_cast<int>(ready_events.size()),
    1);

    if (ready_count == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("epoll_wait");
    break;
    }

    for (int i = 0; i < ready_count; ++i) {
    int current_fd = ready_events[i].data.fd;
    uint32_t events = ready_events[i].events;

    if (current_fd == listen_fd) {
    // 监听 fd 已经非阻塞,因此可以一直 accept 到 EAGAIN。
    while (true) {
    sockaddr_in client_addr{};
    socklen_t client_len =
    sizeof(client_addr);

    int client_fd =
    accept(listen_fd,
    reinterpret_cast<sockaddr*>(
    &client_addr),
    &client_len);

    if (client_fd == 1) {
    if (errno == EINTR ||
    errno == ECONNABORTED) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("accept");
    break;
    }

    if (!SetNonBlocking(client_fd)) {
    std::perror("fcntl(client_fd)");
    close(client_fd);
    continue;
    }

    if (!AddClient(epoll_fd, client_fd)) {
    std::perror("epoll_ctl(ADD client)");
    close(client_fd);
    continue;
    }

    char ip[INET_ADDRSTRLEN]{};
    inet_ntop(AF_INET,
    &client_addr.sin_addr,
    ip,
    sizeof(ip));

    std::cout << "connected: "
    << ip << ':'
    << ntohs(client_addr.sin_port)
    << ", fd=" << client_fd << '\\n';

    clients.insert(client_fd);
    }

    continue;
    }

    bool close_now = false;

    if ((events & EPOLLERR) != 0U) {
    int socket_error = 0;
    socklen_t error_len = sizeof(socket_error);
    getsockopt(current_fd,
    SOL_SOCKET,
    SO_ERROR,
    &socket_error,
    &error_len);

    std::cerr << "socket error on fd="
    << current_fd
    << ": errno=" << socket_error
    << '\\n';
    close_now = true;
    }

    if (!close_now &&
    (events & EPOLLIN) != 0U) {
    // 即使 LT 不强制排空,非阻塞循环读取可以减少重复唤醒。
    while (true) {
    char buffer[kBufferSize];
    ssize_t n = recv(current_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    std::cout << "received "
    << n << " bytes from fd="
    << current_fd << '\\n';

    if (!SendImmediatelyOrFail(
    current_fd,
    buffer,
    static_cast<std::size_t>(n))) {
    close_now = true;
    break;
    }

    continue;
    }

    if (n == 0) {
    close_now = true;
    break;
    }

    if (errno == EINTR) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("recv");
    close_now = true;
    break;
    }
    }

    if ((events &
    (EPOLLRDHUP | EPOLLHUP)) != 0U) {
    close_now = true;
    }

    if (close_now) {
    CloseClient(epoll_fd,
    current_fd,
    clients);
    }
    }
    }

    for (int client_fd : clients) {
    close(client_fd);
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    5.3 完整客户端:02_epoll_lt_nonblocking/client.cpp

    客户端仍然使用 select() 同时监控:

    • 标准输入;
    • 服务器套接字。

    这样做是为了把本章的变量限制在服务端 epoll 重构上,而不是同时改变客户端模型。

    #include <arpa/inet.h>
    #include <sys/select.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>

    namespace {
    bool SendAll(int fd, const char* data, std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (socket_fd == 1) {
    std::perror("socket");
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(8888);

    if (inet_pton(AF_INET,
    "127.0.0.1",
    &server_addr.sin_addr) != 1) {
    std::cerr << "invalid IPv4 address\\n";
    close(socket_fd);
    return 1;
    }

    if (connect(socket_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("connect");
    close(socket_fd);
    return 1;
    }

    std::cout << "connected to 127.0.0.1:8888\\n"
    << "type text; Ctrl+D stops sending but keeps receiving\\n";

    bool stdin_open = true;

    while (true) {
    fd_set read_set;
    FD_ZERO(&read_set);

    if (stdin_open) {
    FD_SET(STDIN_FILENO, &read_set);
    }
    FD_SET(socket_fd, &read_set);

    int max_fd = socket_fd > STDIN_FILENO
    ? socket_fd
    : STDIN_FILENO;

    int ready =
    select(max_fd + 1,
    &read_set,
    nullptr,
    nullptr,
    nullptr);

    if (ready == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("select");
    break;
    }

    if (stdin_open &&
    FD_ISSET(STDIN_FILENO, &read_set)) {
    char buffer[4096];
    ssize_t n = read(STDIN_FILENO,
    buffer,
    sizeof(buffer));

    if (n > 0) {
    if (!SendAll(socket_fd,
    buffer,
    static_cast<std::size_t>(n))) {
    close(socket_fd);
    return 1;
    }
    } else if (n == 0) {
    stdin_open = false;

    if (shutdown(socket_fd, SHUT_WR) == 1 &&
    errno != ENOTCONN) {
    std::perror("shutdown");
    }
    } else if (errno != EINTR) {
    std::perror("read(stdin)");
    break;
    }
    }

    if (FD_ISSET(socket_fd, &read_set)) {
    char buffer[4096];
    ssize_t n = recv(socket_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    std::cout.write(buffer, n);
    std::cout.flush();
    } else if (n == 0) {
    std::cout << "\\nserver closed connection\\n";
    close(socket_fd);
    return 0;
    } else if (errno != EINTR) {
    std::perror("recv");
    break;
    }
    }
    }

    close(socket_fd);
    return 1;
    }

    5.4 本层修复了什么

    • 监听套接字不再因为多调用一次 accept() 而阻塞;
    • 一次监听事件可以排空当前已完成连接队列;
    • 通信套接字不再让事件循环长期阻塞;
    • 正确处理 EINTR;
    • 正确处理 EAGAIN/EWOULDBLOCK;
    • 使用 EPOLLRDHUP 感知对端关闭写方向;
    • 使用 MSG_NOSIGNAL 避免向已关闭连接发送数据时触发 SIGPIPE 终止进程;
    • 读取路径可以一次处理当前已经到达的多批数据。

    5.5 为什么不能把 EAGAIN 当成错误

    非阻塞套接字中:

    recv() == -1
    errno == EAGAIN / EWOULDBLOCK

    表示:

    当前没有更多数据可以立即读取,本轮处理正常结束。

    同理,非阻塞 send() 返回 EAGAIN 表示:

    当前发送缓冲区暂时没有足够空间,需要稍后等待可写事件。

    它们是事件驱动程序中的正常控制流,而不是异常崩溃条件。

    5.6 本层仍然存在的问题

    本层的读路径已经较稳健,但写路径仍然采用:

    立即写完

    ├─ 成功:继续
    └─ 写不完/EAGAIN:关闭连接

    下一步切换 ET 时,先保持这个限制不变,只专注解决 ET 读取规则。


    第四部分:LT 与 ET 的真正区别

    六、水平触发 LT

    LT 可以理解为:

    只要 fd 当前仍然满足就绪条件,就可能继续通知。

    假设接收缓冲区有 5000 字节,而应用每次只读 1024 字节:

    第 1 次 epoll_wait -> EPOLLIN
    read 1024,剩余 3976

    第 2 次 epoll_wait -> 仍可能 EPOLLIN
    read 1024,剩余 2952

    ……

    七、边缘触发 ET

    ET 更接近:

    当监控对象的就绪状态发生变化时通知;收到通知后,应用应把该 fd 当作可操作对象处理,直到非阻塞 I/O 返回 EAGAIN。

    注意,ET 不能被机械地简化为严格的“缓冲区从 0 字节变成非 0 字节只通知一次”。多个数据到达过程仍可能产生多个事件。真正可靠的编程契约是:

    收到 ET 事件


    认为 fd 当前可以继续处理


    非阻塞 read/write 循环

    └─ 直到 EAGAIN 才把控制权交回 epoll_wait


    第五部分:ET 的两个典型错误

    八、错误一:EPOLLET + 单次 read

    8.1 错误逻辑

    event.events = EPOLLIN | EPOLLET;

    // 收到事件后只读取一次
    read(fd, buffer, 1024);

    假设客户端发送 5000 字节:

    接收缓冲区
    ┌────────────────────────────────────────────┐
    │ 前 1024 字节 │ 剩余 3976 字节 │
    └────────────────────────────────────────────┘

    └─ 单次 read 只取走这里

    剩余数据并没有被内核“丢掉”,而是仍留在接收缓冲区中。

    问题在于:

    应用已经消费了这一轮边缘通知,却没有把当前可读状态处理完;之后可能长时间收不到新的通知。

    所以更准确的术语是:

    • 未读数据滞留;
    • 读取停滞;
    • 事件通知遗漏后的应用层饥饿。

    不应简单称为“TCP 数据丢失”。

    8.2 完整错误复现代码

    该代码只用于实验,不应作为正确实现使用。

    // 错误演示:EPOLLET + 非阻塞套接字 + 单次 read。
    // 结果不是“内核丢包”,而是未读数据可能长期留在接收缓冲区,
    // 事件循环却不一定再次收到通知。

    #include <arpa/inet.h>
    #include <fcntl.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>

    bool SetNonBlocking(int fd) {
    int flags = fcntl(fd, F_GETFL, 0);
    return flags != 1 &&
    fcntl(fd, F_SETFL, flags | O_NONBLOCK) != 1;
    }

    int main() {
    int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
    int reuse = 1;
    setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse));
    SetNonBlocking(listen_fd);

    sockaddr_in addr{};
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = htonl(INADDR_ANY);
    addr.sin_port = htons(8888);

    bind(listen_fd,
    reinterpret_cast<sockaddr*>(&addr),
    sizeof(addr));
    listen(listen_fd, 128);

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);

    epoll_event listen_event{};
    listen_event.events = EPOLLIN | EPOLLET;
    listen_event.data.fd = listen_fd;
    epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    listen_fd,
    &listen_event);

    epoll_event events[16]{};

    while (true) {
    int count = epoll_wait(epoll_fd,
    events,
    16,
    1);

    if (count == 1) {
    if (errno == EINTR) {
    continue;
    }
    std::perror("epoll_wait");
    break;
    }

    for (int i = 0; i < count; ++i) {
    int fd = events[i].data.fd;

    if (fd == listen_fd) {
    while (true) {
    int client_fd =
    accept(listen_fd,
    nullptr,
    nullptr);

    if (client_fd == 1) {
    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }
    break;
    }

    SetNonBlocking(client_fd);

    epoll_event client_event{};
    client_event.events =
    EPOLLIN | EPOLLET;
    client_event.data.fd = client_fd;

    epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &client_event);
    }
    } else {
    char buffer[1024];

    // 错误点:只读取一次,没有排空到 EAGAIN。
    ssize_t n = read(fd,
    buffer,
    sizeof(buffer));

    std::cout << "single read returned "
    << n << " bytes\\n";
    }
    }
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    九、错误二:EPOLLET + 阻塞套接字 + 循环 read

    有人发现单次读取不够,于是改成:

    while (true) {
    read(fd, buffer, sizeof(buffer));
    }

    但如果 fd 仍是阻塞模式:

    第 1 次 read -> 1024 字节
    第 2 次 read -> 1024 字节
    ……
    最后一批 read -> 成功
    下一次 read -> 接收缓冲区为空


    当前线程睡眠


    整个事件循环停止处理其他 fd

    这通常不是严格意义上的“死锁”,因为新数据到来后 read() 仍可能返回。

    更准确的描述是:

    单个阻塞套接字把整个单线程事件循环长期卡住,其他连接发生饥饿。

    9.1 完整错误复现代码

    该代码只用于实验,不应作为正确实现使用。

    // 错误演示:EPOLLET + 阻塞通信套接字 + while(read)。
    // 当接收缓冲区被读空后,下一次 read 会阻塞整个事件循环。
    // 这不是严格意义上的“死锁”,而是事件循环被单个连接长期阻塞。

    #include <arpa/inet.h>
    #include <fcntl.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>

    bool SetNonBlocking(int fd) {
    int flags = fcntl(fd, F_GETFL, 0);
    return flags != 1 &&
    fcntl(fd, F_SETFL, flags | O_NONBLOCK) != 1;
    }

    int main() {
    int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
    int reuse = 1;
    setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse));

    // 监听套接字非阻塞,方便循环 accept。
    SetNonBlocking(listen_fd);

    sockaddr_in addr{};
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = htonl(INADDR_ANY);
    addr.sin_port = htons(8888);

    bind(listen_fd,
    reinterpret_cast<sockaddr*>(&addr),
    sizeof(addr));
    listen(listen_fd, 128);

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);

    epoll_event listen_event{};
    listen_event.events = EPOLLIN | EPOLLET;
    listen_event.data.fd = listen_fd;

    epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    listen_fd,
    &listen_event);

    epoll_event events[16]{};

    while (true) {
    int count = epoll_wait(epoll_fd,
    events,
    16,
    1);

    if (count == 1) {
    if (errno == EINTR) {
    continue;
    }
    std::perror("epoll_wait");
    break;
    }

    for (int i = 0; i < count; ++i) {
    int fd = events[i].data.fd;

    if (fd == listen_fd) {
    while (true) {
    int client_fd =
    accept(listen_fd,
    nullptr,
    nullptr);

    if (client_fd == 1) {
    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }
    break;
    }

    // 错误点:没有把 client_fd 设置为非阻塞。
    epoll_event client_event{};
    client_event.events =
    EPOLLIN | EPOLLET;
    client_event.data.fd = client_fd;

    epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &client_event);
    }
    } else {
    char buffer[1024];

    while (true) {
    // 数据读空后,这里会睡眠等待下一批数据,
    // 事件循环无法处理其他 fd。
    ssize_t n = read(fd,
    buffer,
    sizeof(buffer));

    if (n <= 0) {
    break;
    }

    std::cout << "read "
    << n << " bytes\\n";
    }
    }
    }
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    十、ET 正确读法

    ET 读取路径应保持三个不变量:

    1. fd 必须是非阻塞模式
    2. 收到可读事件后循环 recv()
    3. 只有遇到 EAGAIN/EWOULDBLOCK 才结束本轮读取

    状态表:

    recv() 结果含义动作
    n > 0 读到 n 字节 处理数据,继续读
    n == 0 接收方向到达 EOF 标记对端写方向关闭
    n == -1, EINTR 被信号中断 重试
    n == -1, EAGAIN/EWOULDBLOCK 当前已经读空 正常结束本轮
    其他错误 连接异常 关闭连接

    第六部分:第 3 层——正确的 epoll ET 读取

    十一、本层相对第 2 层改了什么

    代码结构几乎不变,关键差异只有事件掩码:

    // LT
    EPOLLIN | EPOLLRDHUP

    // ET
    EPOLLIN | EPOLLRDHUP | EPOLLET

    因为第 2 层已经完成:

    • 非阻塞监听 fd;
    • 非阻塞通信 fd;
    • 循环 accept() 到 EAGAIN;
    • 循环 recv() 到 EAGAIN;

    所以切换 ET 不需要一次性重写整个事件循环。

    这正是分层升级的价值。

    11.1 完整服务端:03_epoll_et_correct_read/server.cpp

    #include <arpa/inet.h>
    #include <fcntl.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>
    #include <unordered_set>
    #include <vector>

    namespace {
    constexpr uint16_t kPort = 8888;
    constexpr int kBacklog = 128;
    constexpr int kMaxEvents = 64;
    constexpr std::size_t kBufferSize = 4096;

    bool SetNonBlocking(int fd) {
    int old_flags = fcntl(fd, F_GETFL, 0);
    if (old_flags == 1) {
    return false;
    }

    return fcntl(fd,
    F_SETFL,
    old_flags | O_NONBLOCK) != 1;
    }

    bool AddClient(int epoll_fd, int client_fd) {
    epoll_event event{};
    event.events =
    EPOLLIN | EPOLLRDHUP | EPOLLET;
    event.data.fd = client_fd;

    return epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &event) != 1;
    }

    void CloseClient(int epoll_fd,
    int client_fd,
    std::unordered_set<int>& clients) {
    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_DEL,
    client_fd,
    nullptr) == 1 &&
    errno != ENOENT &&
    errno != EBADF) {
    std::perror("epoll_ctl(DEL)");
    }

    close(client_fd);
    clients.erase(client_fd);
    std::cout << "closed fd=" << client_fd << '\\n';
    }

    // 本层重点是把 ET 的“读路径”做正确。
    // 写路径仍采用“必须立即全部写完,否则关闭连接”的简化策略。
    bool SendImmediatelyOrFail(int fd,
    const char* data,
    std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    if (n == 1 &&
    (errno == EAGAIN || errno == EWOULDBLOCK)) {
    std::cerr
    << "send buffer full; layer 03 has no EPOLLOUT queue\\n";
    return false;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int listen_fd = socket(AF_INET,
    SOCK_STREAM | SOCK_CLOEXEC,
    0);

    if (listen_fd == 1) {
    std::perror("socket");
    return 1;
    }

    if (!SetNonBlocking(listen_fd)) {
    std::perror("fcntl(listen_fd)");
    close(listen_fd);
    return 1;
    }

    int reuse = 1;
    if (setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse)) == 1) {
    std::perror("setsockopt");
    close(listen_fd);
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
    server_addr.sin_port = htons(kPort);

    if (bind(listen_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("bind");
    close(listen_fd);
    return 1;
    }

    if (listen(listen_fd, kBacklog) == 1) {
    std::perror("listen");
    close(listen_fd);
    return 1;
    }

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
    if (epoll_fd == 1) {
    std::perror("epoll_create1");
    close(listen_fd);
    return 1;
    }

    epoll_event listen_event{};
    listen_event.events = EPOLLIN | EPOLLET;
    listen_event.data.fd = listen_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    listen_fd,
    &listen_event) == 1) {
    std::perror("epoll_ctl(ADD listen_fd)");
    close(epoll_fd);
    close(listen_fd);
    return 1;
    }

    std::vector<epoll_event> ready_events(kMaxEvents);
    std::unordered_set<int> clients;

    std::cout << "epoll ET read-correct server listening on 0.0.0.0:"
    << kPort << '\\n';

    while (true) {
    int ready_count =
    epoll_wait(epoll_fd,
    ready_events.data(),
    static_cast<int>(ready_events.size()),
    1);

    if (ready_count == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("epoll_wait");
    break;
    }

    for (int i = 0; i < ready_count; ++i) {
    int current_fd = ready_events[i].data.fd;
    uint32_t events = ready_events[i].events;

    if (current_fd == listen_fd) {
    // ET:必须一直 accept 到 EAGAIN。
    while (true) {
    sockaddr_in client_addr{};
    socklen_t client_len =
    sizeof(client_addr);

    int client_fd =
    accept(listen_fd,
    reinterpret_cast<sockaddr*>(
    &client_addr),
    &client_len);

    if (client_fd == 1) {
    if (errno == EINTR ||
    errno == ECONNABORTED) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("accept");
    break;
    }

    if (!SetNonBlocking(client_fd)) {
    std::perror("fcntl(client_fd)");
    close(client_fd);
    continue;
    }

    if (!AddClient(epoll_fd, client_fd)) {
    std::perror("epoll_ctl(ADD client)");
    close(client_fd);
    continue;
    }

    char ip[INET_ADDRSTRLEN]{};
    inet_ntop(AF_INET,
    &client_addr.sin_addr,
    ip,
    sizeof(ip));

    std::cout << "connected: "
    << ip << ':'
    << ntohs(client_addr.sin_port)
    << ", fd=" << client_fd << '\\n';

    clients.insert(client_fd);
    }

    continue;
    }

    bool close_now = false;

    if ((events & EPOLLERR) != 0U) {
    int socket_error = 0;
    socklen_t error_len = sizeof(socket_error);
    getsockopt(current_fd,
    SOL_SOCKET,
    SO_ERROR,
    &socket_error,
    &error_len);

    std::cerr << "socket error on fd="
    << current_fd
    << ": errno=" << socket_error
    << '\\n';
    close_now = true;
    }

    if (!close_now &&
    (events & EPOLLIN) != 0U) {
    // ET 核心:非阻塞 + 循环 recv,直到 EAGAIN。
    while (true) {
    char buffer[kBufferSize];
    ssize_t n = recv(current_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    std::cout << "received "
    << n << " bytes from fd="
    << current_fd << '\\n';

    if (!SendImmediatelyOrFail(
    current_fd,
    buffer,
    static_cast<std::size_t>(n))) {
    close_now = true;
    break;
    }

    continue;
    }

    if (n == 0) {
    close_now = true;
    break;
    }

    if (errno == EINTR) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    // 当前接收缓冲区已经排空。
    break;
    }

    std::perror("recv");
    close_now = true;
    break;
    }
    }

    if ((events &
    (EPOLLRDHUP | EPOLLHUP)) != 0U) {
    close_now = true;
    }

    if (close_now) {
    CloseClient(epoll_fd,
    current_fd,
    clients);
    }
    }
    }

    for (int client_fd : clients) {
    close(client_fd);
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    11.2 完整客户端:03_epoll_et_correct_read/client.cpp

    #include <arpa/inet.h>
    #include <sys/select.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>

    namespace {
    bool SendAll(int fd, const char* data, std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (socket_fd == 1) {
    std::perror("socket");
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(8888);

    if (inet_pton(AF_INET,
    "127.0.0.1",
    &server_addr.sin_addr) != 1) {
    std::cerr << "invalid IPv4 address\\n";
    close(socket_fd);
    return 1;
    }

    if (connect(socket_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("connect");
    close(socket_fd);
    return 1;
    }

    std::cout << "connected to 127.0.0.1:8888\\n"
    << "type text; Ctrl+D stops sending but keeps receiving\\n";

    bool stdin_open = true;

    while (true) {
    fd_set read_set;
    FD_ZERO(&read_set);

    if (stdin_open) {
    FD_SET(STDIN_FILENO, &read_set);
    }
    FD_SET(socket_fd, &read_set);

    int max_fd = socket_fd > STDIN_FILENO
    ? socket_fd
    : STDIN_FILENO;

    int ready =
    select(max_fd + 1,
    &read_set,
    nullptr,
    nullptr,
    nullptr);

    if (ready == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("select");
    break;
    }

    if (stdin_open &&
    FD_ISSET(STDIN_FILENO, &read_set)) {
    char buffer[4096];
    ssize_t n = read(STDIN_FILENO,
    buffer,
    sizeof(buffer));

    if (n > 0) {
    if (!SendAll(socket_fd,
    buffer,
    static_cast<std::size_t>(n))) {
    close(socket_fd);
    return 1;
    }
    } else if (n == 0) {
    stdin_open = false;

    if (shutdown(socket_fd, SHUT_WR) == 1 &&
    errno != ENOTCONN) {
    std::perror("shutdown");
    }
    } else if (errno != EINTR) {
    std::perror("read(stdin)");
    break;
    }
    }

    if (FD_ISSET(socket_fd, &read_set)) {
    char buffer[4096];
    ssize_t n = recv(socket_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    std::cout.write(buffer, n);
    std::cout.flush();
    } else if (n == 0) {
    std::cout << "\\nserver closed connection\\n";
    close(socket_fd);
    return 0;
    } else if (errno != EINTR) {
    std::perror("recv");
    break;
    }
    }
    }

    close(socket_fd);
    return 1;
    }

    11.3 本层修复了什么

    • 监听 fd 在 ET 模式下排空连接队列;
    • 通信 fd 在 ET 模式下排空接收缓冲区;
    • 数据不会因为只读一次而长期滞留;
    • 事件循环不会因为阻塞 recv() 卡住;
    • EAGAIN 成为“本轮处理完成”的边界。

    11.4 本层仍然没有解决写路径

    当前逻辑仍然要求:

    收到一批数据


    立即 send

    ├─ 全部发送:正常
    └─ 部分发送/EAGAIN:关闭连接

    这不会悄悄丢数据,但会错误地拒绝慢客户端。

    要正确处理它,必须理解 EPOLLOUT。


    第七部分:为什么不能永久监听 EPOLLOUT

    十二、套接字“可写”意味着什么

    对 TCP 套接字而言,EPOLLOUT 通常表示:

    发送缓冲区当前有空间,执行非阻塞发送可能取得进展。

    新建立的 TCP 连接大多数时间都是可写的。

    如果无条件永久注册:

    EPOLLIN | EPOLLOUT

    事件循环可能不断得到 EPOLLOUT:

    epoll_wait -> fd 可写
    没有数据需要发送
    再次 epoll_wait -> fd 仍可写
    没有数据需要发送
    ……

    这会形成无意义的高频唤醒。

    因此常见策略是:

    没有待发送数据:
    只监听 EPOLLIN

    出现待发送数据:
    先主动 send

    全部发送完成:
    不监听 EPOLLOUT

    send 返回 EAGAIN 或只完成一部分:
    保存剩余数据
    注册 EPOLLOUT

    后续 EPOLLOUT 到来:
    继续发送

    十三、为什么每个连接都需要独立发送状态

    TCP 发送是每连接独立的,因此需要:

    struct ClientState {
    std::string output;
    std::size_t sent_offset;
    };

    假设:

    output = "ABCDEFGHIJ"
    sent_offset = 4

    表示:

    已经发送:ABCD
    等待发送:EFGHIJ

    下一次发送地址为:

    output.data() + sent_offset

    剩余长度为:

    output.size() sent_offset


    第八部分:第 4 层——ET + 发送队列 + EPOLLOUT

    十四、本层新增的状态机

    收到客户端数据


    追加 output


    主动尝试 send
    │ │
    │ └─ EAGAIN / 部分发送
    │ │
    ▼ ▼
    全部发完 保留 sent_offset
    │ │
    移除 EPOLLOUT 注册 EPOLLOUT


    下次可写时继续

    14.1 完整服务端:04_epoll_et_write_queue/server.cpp

    #include <arpa/inet.h>
    #include <fcntl.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>
    #include <string>
    #include <unordered_map>
    #include <vector>

    namespace {
    constexpr uint16_t kPort = 8888;
    constexpr int kBacklog = 128;
    constexpr int kMaxEvents = 64;
    constexpr std::size_t kBufferSize = 4096;

    struct ClientState {
    std::string peer;
    std::string output;
    std::size_t sent_offset = 0;
    bool peer_closed_write = false;
    };

    bool SetNonBlocking(int fd) {
    int old_flags = fcntl(fd, F_GETFL, 0);
    if (old_flags == 1) {
    return false;
    }

    return fcntl(fd,
    F_SETFL,
    old_flags | O_NONBLOCK) != 1;
    }

    std::size_t PendingBytes(const ClientState& client) {
    return client.output.size() client.sent_offset;
    }

    void CompactOutput(ClientState& client) {
    if (client.sent_offset == 0) {
    return;
    }

    if (client.sent_offset == client.output.size()) {
    client.output.clear();
    client.sent_offset = 0;
    return;
    }

    if (client.sent_offset >= 4096 &&
    client.sent_offset * 2 >= client.output.size()) {
    client.output.erase(0, client.sent_offset);
    client.sent_offset = 0;
    }
    }

    uint32_t BuildClientEvents(const ClientState& client) {
    uint32_t events = EPOLLET | EPOLLRDHUP;

    if (!client.peer_closed_write) {
    events |= EPOLLIN;
    }

    if (PendingBytes(client) > 0) {
    events |= EPOLLOUT;
    }

    return events;
    }

    bool ModifyClientEvents(int epoll_fd,
    int client_fd,
    const ClientState& client) {
    epoll_event event{};
    event.events = BuildClientEvents(client);
    event.data.fd = client_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_MOD,
    client_fd,
    &event) == 1) {
    std::perror("epoll_ctl(MOD client)");
    return false;
    }

    return true;
    }

    void CloseClient(
    int epoll_fd,
    int client_fd,
    std::unordered_map<int, ClientState>& clients) {
    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_DEL,
    client_fd,
    nullptr) == 1 &&
    errno != ENOENT &&
    errno != EBADF) {
    std::perror("epoll_ctl(DEL)");
    }

    close(client_fd);
    clients.erase(client_fd);
    std::cout << "closed fd=" << client_fd << '\\n';
    }

    bool FlushOutput(int client_fd,
    ClientState& client) {
    while (client.sent_offset < client.output.size()) {
    const char* data =
    client.output.data() + client.sent_offset;

    std::size_t remain =
    client.output.size() client.sent_offset;

    ssize_t n =
    send(client_fd,
    data,
    remain,
    MSG_NOSIGNAL);

    if (n > 0) {
    client.sent_offset +=
    static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    if (n == 1 &&
    (errno == EAGAIN || errno == EWOULDBLOCK)) {
    // 发送缓冲区暂时已满,保留未发送数据,
    // 等待下一次 EPOLLOUT。
    CompactOutput(client);
    return true;
    }

    std::perror("send");
    return false;
    }

    CompactOutput(client);
    return true;
    }
    }

    int main() {
    int listen_fd = socket(AF_INET,
    SOCK_STREAM | SOCK_CLOEXEC,
    0);

    if (listen_fd == 1) {
    std::perror("socket");
    return 1;
    }

    if (!SetNonBlocking(listen_fd)) {
    std::perror("fcntl(listen_fd)");
    close(listen_fd);
    return 1;
    }

    int reuse = 1;
    if (setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse)) == 1) {
    std::perror("setsockopt");
    close(listen_fd);
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
    server_addr.sin_port = htons(kPort);

    if (bind(listen_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("bind");
    close(listen_fd);
    return 1;
    }

    if (listen(listen_fd, kBacklog) == 1) {
    std::perror("listen");
    close(listen_fd);
    return 1;
    }

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
    if (epoll_fd == 1) {
    std::perror("epoll_create1");
    close(listen_fd);
    return 1;
    }

    epoll_event listen_event{};
    listen_event.events = EPOLLIN | EPOLLET;
    listen_event.data.fd = listen_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    listen_fd,
    &listen_event) == 1) {
    std::perror("epoll_ctl(ADD listen_fd)");
    close(epoll_fd);
    close(listen_fd);
    return 1;
    }

    std::vector<epoll_event> ready_events(kMaxEvents);
    std::unordered_map<int, ClientState> clients;

    std::cout << "epoll ET output-queue server listening on 0.0.0.0:"
    << kPort << '\\n';

    while (true) {
    int ready_count =
    epoll_wait(epoll_fd,
    ready_events.data(),
    static_cast<int>(ready_events.size()),
    1);

    if (ready_count == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("epoll_wait");
    break;
    }

    for (int i = 0; i < ready_count; ++i) {
    int current_fd = ready_events[i].data.fd;
    uint32_t events = ready_events[i].events;

    if (current_fd == listen_fd) {
    while (true) {
    sockaddr_in client_addr{};
    socklen_t client_len =
    sizeof(client_addr);

    int client_fd =
    accept(listen_fd,
    reinterpret_cast<sockaddr*>(
    &client_addr),
    &client_len);

    if (client_fd == 1) {
    if (errno == EINTR ||
    errno == ECONNABORTED) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("accept");
    break;
    }

    if (!SetNonBlocking(client_fd)) {
    std::perror("fcntl(client_fd)");
    close(client_fd);
    continue;
    }

    char ip[INET_ADDRSTRLEN]{};
    inet_ntop(AF_INET,
    &client_addr.sin_addr,
    ip,
    sizeof(ip));

    ClientState client;
    client.peer =
    std::string(ip) + ":" +
    std::to_string(
    ntohs(client_addr.sin_port));

    epoll_event client_event{};
    client_event.events =
    EPOLLIN | EPOLLRDHUP | EPOLLET;
    client_event.data.fd = client_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &client_event) == 1) {
    std::perror(
    "epoll_ctl(ADD client)");
    close(client_fd);
    continue;
    }

    std::cout << "connected: "
    << client.peer
    << ", fd=" << client_fd << '\\n';

    clients.emplace(client_fd,
    std::move(client));
    }

    continue;
    }

    auto it = clients.find(current_fd);
    if (it == clients.end()) {
    continue;
    }

    ClientState& client = it->second;
    bool close_now = false;

    if ((events & EPOLLERR) != 0U) {
    int socket_error = 0;
    socklen_t error_len = sizeof(socket_error);
    getsockopt(current_fd,
    SOL_SOCKET,
    SO_ERROR,
    &socket_error,
    &error_len);

    std::cerr << "socket error on fd="
    << current_fd
    << ": errno=" << socket_error
    << '\\n';
    close_now = true;
    }

    if (!close_now &&
    (events &
    (EPOLLIN | EPOLLRDHUP | EPOLLHUP)) != 0U) {
    while (true) {
    char buffer[kBufferSize];
    ssize_t n = recv(current_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    CompactOutput(client);
    client.output.append(
    buffer,
    static_cast<std::size_t>(n));

    std::cout << "queued "
    << n << " echo bytes for fd="
    << current_fd << '\\n';
    continue;
    }

    if (n == 0) {
    client.peer_closed_write = true;
    break;
    }

    if (errno == EINTR) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("recv");
    close_now = true;
    break;
    }
    }

    // EPOLLRDHUP/EPOLLHUP 只是“对端关闭方向发生变化”的提示。
    // 仍要通过 recv() 返回 0 确认接收方向已经到达 EOF,
    // 因为 FIN 之前可能还有尚未读取的数据。

    // 无论本轮是否带 EPOLLOUT,只要存在待发送数据,
    // 都先主动尝试一次发送;写不动再订阅 EPOLLOUT。
    if (!close_now &&
    PendingBytes(client) > 0) {
    if (!FlushOutput(current_fd, client)) {
    close_now = true;
    }
    }

    if (!close_now &&
    client.peer_closed_write &&
    PendingBytes(client) == 0) {
    close_now = true;
    }

    if (close_now) {
    CloseClient(epoll_fd,
    current_fd,
    clients);
    continue;
    }

    if (!ModifyClientEvents(epoll_fd,
    current_fd,
    client)) {
    CloseClient(epoll_fd,
    current_fd,
    clients);
    }
    }
    }

    for (const auto& [client_fd, client] : clients) {
    (void)client;
    close(client_fd);
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    14.2 完整客户端:04_epoll_et_write_queue/client.cpp

    #include <arpa/inet.h>
    #include <sys/select.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>

    namespace {
    bool SendAll(int fd, const char* data, std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (socket_fd == 1) {
    std::perror("socket");
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(8888);

    if (inet_pton(AF_INET,
    "127.0.0.1",
    &server_addr.sin_addr) != 1) {
    std::cerr << "invalid IPv4 address\\n";
    close(socket_fd);
    return 1;
    }

    if (connect(socket_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("connect");
    close(socket_fd);
    return 1;
    }

    std::cout << "connected to 127.0.0.1:8888\\n"
    << "type text; Ctrl+D stops sending but keeps receiving\\n";

    bool stdin_open = true;

    while (true) {
    fd_set read_set;
    FD_ZERO(&read_set);

    if (stdin_open) {
    FD_SET(STDIN_FILENO, &read_set);
    }
    FD_SET(socket_fd, &read_set);

    int max_fd = socket_fd > STDIN_FILENO
    ? socket_fd
    : STDIN_FILENO;

    int ready =
    select(max_fd + 1,
    &read_set,
    nullptr,
    nullptr,
    nullptr);

    if (ready == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("select");
    break;
    }

    if (stdin_open &&
    FD_ISSET(STDIN_FILENO, &read_set)) {
    char buffer[4096];
    ssize_t n = read(STDIN_FILENO,
    buffer,
    sizeof(buffer));

    if (n > 0) {
    if (!SendAll(socket_fd,
    buffer,
    static_cast<std::size_t>(n))) {
    close(socket_fd);
    return 1;
    }
    } else if (n == 0) {
    stdin_open = false;

    if (shutdown(socket_fd, SHUT_WR) == 1 &&
    errno != ENOTCONN) {
    std::perror("shutdown");
    }
    } else if (errno != EINTR) {
    std::perror("read(stdin)");
    break;
    }
    }

    if (FD_ISSET(socket_fd, &read_set)) {
    char buffer[4096];
    ssize_t n = recv(socket_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    std::cout.write(buffer, n);
    std::cout.flush();
    } else if (n == 0) {
    std::cout << "\\nserver closed connection\\n";
    close(socket_fd);
    return 0;
    } else if (errno != EINTR) {
    std::perror("recv");
    break;
    }
    }
    }

    close(socket_fd);
    return 1;
    }

    14.3 BuildClientEvents() 的意义

    事件掩码不能再写死:

    uint32_t BuildClientEvents(const ClientState& client) {
    uint32_t events = EPOLLET | EPOLLRDHUP;

    if (!client.peer_closed_write) {
    events |= EPOLLIN;
    }

    if (PendingBytes(client) > 0) {
    events |= EPOLLOUT;
    }

    return events;
    }

    它把连接状态映射成内核兴趣集合:

    连接仍可接收数据 -> EPOLLIN
    存在待发送数据 -> EPOLLOUT
    使用边缘触发 -> EPOLLET
    关注对端半关闭 -> EPOLLRDHUP

    14.4 为什么收到数据后先主动发送

    收到数据并加入发送队列后,不必等下一次 EPOLLOUT 才第一次尝试发送。

    套接字此时很可能已经可写:

    client.output.append(buffer, n);
    FlushOutput(fd, client);

    只有确实写不动时,才把 EPOLLOUT 加入事件掩码。

    14.5 EPOLLRDHUP 不能直接等同于“所有数据已经读完”

    这是一个非常容易写错的细节。

    TCP 数据和 FIN 的顺序可能是:

    业务数据 1
    业务数据 2
    FIN

    当 epoll 报告 EPOLLRDHUP 时,接收缓冲区中仍可能存在 FIN 之前到达但尚未读取的业务数据。

    因此正确做法是:

    EPOLLRDHUP / EPOLLHUP


    继续执行非阻塞 recv()

    ├─ n > 0:继续处理遗留数据
    └─ n == 0:才确认接收方向 EOF

    EPOLLRDHUP 是关闭方向变化的提示,recv() == 0 才是接收方向到达 EOF 的确认。

    14.6 本层修复了什么

    • 正确处理部分发送;
    • 正确处理 send() 返回 EAGAIN;
    • 不再因慢客户端暂时写不动而立即断开;
    • 仅在确实存在待发送数据时监听 EPOLLOUT;
    • 支持客户端执行 shutdown(SHUT_WR) 后继续接收回显;
    • 发送完最后的数据后再关闭连接;
    • 使用 MSG_NOSIGNAL 防止 SIGPIPE 结束进程。

    14.7 本层的新风险:内存无限增长

    如果客户端持续发送,但从不读取服务端回显:

    客户端不断发送


    服务端不断 append 到 output


    客户端不接收


    内核发送缓冲区满


    output 越积越多


    进程内存耗尽

    因此“有发送队列”还不够,必须增加背压。


    第九部分:背压是什么

    十五、背压的核心含义

    背压不是一个特定系统调用,而是一种流量控制策略:

    当下游处理速度跟不上上游生产速度时,限制或暂停上游继续产生数据。

    在 Echo 服务器中:

    客户端发送速度 = 上游生产速度
    客户端读取回显速度 = 下游消费速度

    慢客户端造成的问题是:

    recv 很快
    send 很慢
    output 持续增长

    十六、高低水位

    使用两个阈值而不是一个阈值,可以避免事件掩码频繁抖动:

    output < 128 KiB
    正常读取

    output >= 256 KiB
    暂停 EPOLLIN

    暂停后 output 逐渐发送

    output <= 128 KiB
    恢复 EPOLLIN

    ASCII 图:

    待发送字节

    1M │ 硬上限:超过则关闭

    256K├────────────── 高水位:暂停读取
    │ ╲
    │ ╲ 发送队列逐步下降
    128K├────────────────╲── 低水位:恢复读取

    0 └────────────────────────────────→ 时间

    十七、为什么还需要硬上限

    高水位只是暂停读取,但以下情况仍可能让发送队列突破预期:

    • 单次读取的数据本身很大;
    • 业务层一次生成大量响应;
    • 多个处理步骤在暂停生效前已经追加数据;
    • 代码缺陷导致继续入队。

    因此还应设置硬上限:

    if (pending + incoming > kOutputHardLimit) {
    close_connection();
    }

    这不是为了“惩罚客户端”,而是为了保护整个进程不被单个连接耗尽内存。


    第十部分:ET、公平预算和 EPOLLONESHOT

    十八、为什么不能无限排空一个非常活跃的连接

    ET 的经典规则是“一直处理到 EAGAIN”。

    但如果某个客户端持续高速发送,单线程可能长期停留在该连接的读取循环中:

    fd=7 持续有数据


    不断 recv()


    fd=8、fd=9 长时间得不到处理

    这叫 事件循环饥饿。

    一种工程化做法是给每个事件设置预算:

    每次最多读取 64 KiB
    每次最多发送 64 KiB

    预算耗尽后,把控制权交回事件循环,让其他就绪连接获得处理机会。

    十九、ET 下提前停止读取为什么危险

    如果使用普通 ET:

    读了 64 KiB
    但还没有读到 EAGAIN
    直接返回 epoll_wait

    接收缓冲区仍然可读,却不一定出现新的边缘通知,连接可能停滞。

    解决方法之一是:

    • 使用用户态就绪队列;
    • 或使用 EPOLLONESHOT,处理后通过 EPOLL_CTL_MOD 重新激活。

    本文最终层采用第二种方案。

    二十、EPOLLONESHOT 的语义

    注册:

    EPOLLIN | EPOLLET | EPOLLONESHOT

    事件被 epoll_wait() 交付后,这个监控项暂时禁用。

    应用处理完成后必须:

    epoll_ctl(epoll_fd,
    EPOLL_CTL_MOD,
    fd,
    &new_event);

    重新激活。

    这带来两个作用:

  • 可以在每轮处理后重新计算最新事件掩码;
  • 多线程模型中,可以避免同一监控项在未重新激活前被重复交付,但应用仍需正确管理连接对象生命周期和并发状态,不能把 EPOLLONESHOT 简化成“自动线程安全”。

  • 第十一部分:第 5 层——背压、预算与资源边界

    二十一、最终层的状态字段

    struct ClientState {
    std::string output;
    std::size_t sent_offset;

    bool peer_closed_write;
    bool read_paused;
    };

    其中:

    字段含义
    output 已生成但尚未完全发送的数据
    sent_offset output 中已经发送的前缀长度
    peer_closed_write recv() 已返回 0
    read_paused 因待发送数据过多而暂停监听 EPOLLIN

    21.1 完整服务端:05_epoll_et_backpressure/server.cpp

    #include <arpa/inet.h>
    #include <fcntl.h>
    #include <sys/epoll.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <algorithm>
    #include <cerrno>
    #include <cstdio>
    #include <iostream>
    #include <string>
    #include <unordered_map>
    #include <vector>

    namespace {
    constexpr uint16_t kPort = 8888;
    constexpr int kBacklog = 256;
    constexpr int kMaxEvents = 128;
    constexpr std::size_t kBufferSize = 16 * 1024;

    constexpr std::size_t kMaxClients = 10000;
    constexpr std::size_t kReadBudgetPerEvent = 64 * 1024;
    constexpr std::size_t kWriteBudgetPerEvent = 64 * 1024;

    constexpr std::size_t kOutputLowWatermark = 128 * 1024;
    constexpr std::size_t kOutputHighWatermark = 256 * 1024;
    constexpr std::size_t kOutputHardLimit = 1024 * 1024;

    struct ClientState {
    std::string peer;

    std::string output;
    std::size_t sent_offset = 0;

    bool peer_closed_write = false;
    bool read_paused = false;
    };

    bool SetNonBlocking(int fd) {
    int old_flags = fcntl(fd, F_GETFL, 0);
    if (old_flags == 1) {
    return false;
    }

    return fcntl(fd,
    F_SETFL,
    old_flags | O_NONBLOCK) != 1;
    }

    std::size_t PendingBytes(const ClientState& client) {
    return client.output.size() client.sent_offset;
    }

    void CompactOutput(ClientState& client) {
    if (client.sent_offset == 0) {
    return;
    }

    if (client.sent_offset == client.output.size()) {
    client.output.clear();
    client.sent_offset = 0;
    return;
    }

    if (client.sent_offset >= 4096 &&
    client.sent_offset * 2 >= client.output.size()) {
    client.output.erase(0, client.sent_offset);
    client.sent_offset = 0;
    }
    }

    uint32_t BuildClientEvents(const ClientState& client) {
    uint32_t events =
    EPOLLET | EPOLLONESHOT | EPOLLRDHUP;

    if (!client.peer_closed_write &&
    !client.read_paused) {
    events |= EPOLLIN;
    }

    if (PendingBytes(client) > 0) {
    events |= EPOLLOUT;
    }

    return events;
    }

    bool RearmClient(int epoll_fd,
    int client_fd,
    const ClientState& client) {
    epoll_event event{};
    event.events = BuildClientEvents(client);
    event.data.fd = client_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_MOD,
    client_fd,
    &event) == 1) {
    std::perror("epoll_ctl(MOD/rearm)");
    return false;
    }

    return true;
    }

    void CloseClient(
    int epoll_fd,
    int client_fd,
    std::unordered_map<int, ClientState>& clients) {
    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_DEL,
    client_fd,
    nullptr) == 1 &&
    errno != ENOENT &&
    errno != EBADF) {
    std::perror("epoll_ctl(DEL)");
    }

    close(client_fd);
    clients.erase(client_fd);
    std::cout << "closed fd=" << client_fd << '\\n';
    }

    bool QueueEcho(ClientState& client,
    const char* data,
    std::size_t size) {
    CompactOutput(client);

    if (PendingBytes(client) + size >
    kOutputHardLimit) {
    return false;
    }

    client.output.append(data, size);

    if (PendingBytes(client) >=
    kOutputHighWatermark) {
    client.read_paused = true;
    }

    return true;
    }

    bool FlushOutput(int client_fd,
    ClientState& client,
    std::size_t budget) {
    std::size_t written_this_round = 0;

    while (client.sent_offset < client.output.size() &&
    written_this_round < budget) {
    std::size_t remain =
    client.output.size() client.sent_offset;

    std::size_t allowed =
    std::min(remain,
    budget written_this_round);

    ssize_t n =
    send(client_fd,
    client.output.data() +
    client.sent_offset,
    allowed,
    MSG_NOSIGNAL);

    if (n > 0) {
    std::size_t written =
    static_cast<std::size_t>(n);

    client.sent_offset += written;
    written_this_round += written;
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    if (n == 1 &&
    (errno == EAGAIN || errno == EWOULDBLOCK)) {
    break;
    }

    std::perror("send");
    return false;
    }

    CompactOutput(client);

    if (client.read_paused &&
    PendingBytes(client) <=
    kOutputLowWatermark) {
    client.read_paused = false;
    }

    return true;
    }
    }

    int main() {
    int listen_fd = socket(AF_INET,
    SOCK_STREAM | SOCK_CLOEXEC,
    0);

    if (listen_fd == 1) {
    std::perror("socket");
    return 1;
    }

    if (!SetNonBlocking(listen_fd)) {
    std::perror("fcntl(listen_fd)");
    close(listen_fd);
    return 1;
    }

    int reuse = 1;
    if (setsockopt(listen_fd,
    SOL_SOCKET,
    SO_REUSEADDR,
    &reuse,
    sizeof(reuse)) == 1) {
    std::perror("setsockopt");
    close(listen_fd);
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
    server_addr.sin_port = htons(kPort);

    if (bind(listen_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("bind");
    close(listen_fd);
    return 1;
    }

    if (listen(listen_fd, kBacklog) == 1) {
    std::perror("listen");
    close(listen_fd);
    return 1;
    }

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
    if (epoll_fd == 1) {
    std::perror("epoll_create1");
    close(listen_fd);
    return 1;
    }

    epoll_event listen_event{};
    listen_event.events = EPOLLIN | EPOLLET;
    listen_event.data.fd = listen_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    listen_fd,
    &listen_event) == 1) {
    std::perror("epoll_ctl(ADD listen_fd)");
    close(epoll_fd);
    close(listen_fd);
    return 1;
    }

    std::vector<epoll_event> ready_events(kMaxEvents);
    std::unordered_map<int, ClientState> clients;

    std::cout
    << "epoll ET backpressure server listening on 0.0.0.0:"
    << kPort << '\\n';

    while (true) {
    int ready_count =
    epoll_wait(epoll_fd,
    ready_events.data(),
    static_cast<int>(ready_events.size()),
    1);

    if (ready_count == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("epoll_wait");
    break;
    }

    for (int i = 0; i < ready_count; ++i) {
    int current_fd = ready_events[i].data.fd;
    uint32_t events = ready_events[i].events;

    if (current_fd == listen_fd) {
    while (true) {
    sockaddr_in client_addr{};
    socklen_t client_len =
    sizeof(client_addr);

    int client_fd =
    accept(listen_fd,
    reinterpret_cast<sockaddr*>(
    &client_addr),
    &client_len);

    if (client_fd == 1) {
    if (errno == EINTR ||
    errno == ECONNABORTED) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("accept");
    break;
    }

    if (clients.size() >= kMaxClients) {
    std::cerr
    << "reject fd=" << client_fd
    << ": client limit reached\\n";
    close(client_fd);
    continue;
    }

    if (!SetNonBlocking(client_fd)) {
    std::perror("fcntl(client_fd)");
    close(client_fd);
    continue;
    }

    char ip[INET_ADDRSTRLEN]{};
    inet_ntop(AF_INET,
    &client_addr.sin_addr,
    ip,
    sizeof(ip));

    ClientState client;
    client.peer =
    std::string(ip) + ":" +
    std::to_string(
    ntohs(client_addr.sin_port));

    epoll_event client_event{};
    client_event.events =
    BuildClientEvents(client);
    client_event.data.fd = client_fd;

    if (epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &client_event) == 1) {
    std::perror(
    "epoll_ctl(ADD client)");
    close(client_fd);
    continue;
    }

    std::cout << "connected: "
    << client.peer
    << ", fd=" << client_fd << '\\n';

    clients.emplace(client_fd,
    std::move(client));
    }

    continue;
    }

    auto it = clients.find(current_fd);
    if (it == clients.end()) {
    continue;
    }

    ClientState& client = it->second;
    bool close_now = false;

    if ((events & EPOLLERR) != 0U) {
    int socket_error = 0;
    socklen_t error_len = sizeof(socket_error);
    getsockopt(current_fd,
    SOL_SOCKET,
    SO_ERROR,
    &socket_error,
    &error_len);

    std::cerr << "socket error on fd="
    << current_fd
    << ": errno=" << socket_error
    << '\\n';
    close_now = true;
    }

    if (!close_now &&
    !client.read_paused &&
    (events &
    (EPOLLIN | EPOLLRDHUP | EPOLLHUP)) != 0U) {
    std::size_t read_this_round = 0;

    while (read_this_round <
    kReadBudgetPerEvent) {
    char buffer[kBufferSize];

    std::size_t allowed =
    std::min(
    sizeof(buffer),
    kReadBudgetPerEvent
    read_this_round);

    ssize_t n = recv(current_fd,
    buffer,
    allowed,
    0);

    if (n > 0) {
    std::size_t received =
    static_cast<std::size_t>(n);

    if (!QueueEcho(client,
    buffer,
    received)) {
    std::cerr
    << "fd=" << current_fd
    << " output hard limit exceeded\\n";
    close_now = true;
    break;
    }

    read_this_round += received;

    if (client.read_paused) {
    break;
    }

    continue;
    }

    if (n == 0) {
    client.peer_closed_write = true;
    break;
    }

    if (errno == EINTR) {
    continue;
    }

    if (errno == EAGAIN ||
    errno == EWOULDBLOCK) {
    break;
    }

    std::perror("recv");
    close_now = true;
    break;
    }
    }

    // 不直接把 EPOLLRDHUP/EPOLLHUP 等同于“所有数据已读完”。
    // FIN 前面仍可能排着业务数据;只有 recv() 返回 0 才确认 EOF。

    if (!close_now &&
    PendingBytes(client) > 0) {
    if (!FlushOutput(current_fd,
    client,
    kWriteBudgetPerEvent)) {
    close_now = true;
    }
    }

    if (!close_now &&
    client.peer_closed_write &&
    PendingBytes(client) == 0) {
    close_now = true;
    }

    if (close_now) {
    CloseClient(epoll_fd,
    current_fd,
    clients);
    continue;
    }

    // EPOLLONESHOT 在事件交付后会禁用该 fd。
    // 这里通过 MOD 重新激活,同时写入最新的读写兴趣集合。
    if (!RearmClient(epoll_fd,
    current_fd,
    client)) {
    CloseClient(epoll_fd,
    current_fd,
    clients);
    }
    }
    }

    for (const auto& [client_fd, client] : clients) {
    (void)client;
    close(client_fd);
    }

    close(epoll_fd);
    close(listen_fd);
    return 0;
    }

    21.2 完整客户端:05_epoll_et_backpressure/client.cpp

    #include <arpa/inet.h>
    #include <sys/select.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <iostream>

    namespace {
    bool SendAll(int fd, const char* data, std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main() {
    int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (socket_fd == 1) {
    std::perror("socket");
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(8888);

    if (inet_pton(AF_INET,
    "127.0.0.1",
    &server_addr.sin_addr) != 1) {
    std::cerr << "invalid IPv4 address\\n";
    close(socket_fd);
    return 1;
    }

    if (connect(socket_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("connect");
    close(socket_fd);
    return 1;
    }

    std::cout << "connected to 127.0.0.1:8888\\n"
    << "type text; Ctrl+D stops sending but keeps receiving\\n";

    bool stdin_open = true;

    while (true) {
    fd_set read_set;
    FD_ZERO(&read_set);

    if (stdin_open) {
    FD_SET(STDIN_FILENO, &read_set);
    }
    FD_SET(socket_fd, &read_set);

    int max_fd = socket_fd > STDIN_FILENO
    ? socket_fd
    : STDIN_FILENO;

    int ready =
    select(max_fd + 1,
    &read_set,
    nullptr,
    nullptr,
    nullptr);

    if (ready == 1) {
    if (errno == EINTR) {
    continue;
    }

    std::perror("select");
    break;
    }

    if (stdin_open &&
    FD_ISSET(STDIN_FILENO, &read_set)) {
    char buffer[4096];
    ssize_t n = read(STDIN_FILENO,
    buffer,
    sizeof(buffer));

    if (n > 0) {
    if (!SendAll(socket_fd,
    buffer,
    static_cast<std::size_t>(n))) {
    close(socket_fd);
    return 1;
    }
    } else if (n == 0) {
    stdin_open = false;

    if (shutdown(socket_fd, SHUT_WR) == 1 &&
    errno != ENOTCONN) {
    std::perror("shutdown");
    }
    } else if (errno != EINTR) {
    std::perror("read(stdin)");
    break;
    }
    }

    if (FD_ISSET(socket_fd, &read_set)) {
    char buffer[4096];
    ssize_t n = recv(socket_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    std::cout.write(buffer, n);
    std::cout.flush();
    } else if (n == 0) {
    std::cout << "\\nserver closed connection\\n";
    close(socket_fd);
    return 0;
    } else if (errno != EINTR) {
    std::perror("recv");
    break;
    }
    }
    }

    close(socket_fd);
    return 1;
    }

    21.3 最终层的保护机制

    连接数上限

    constexpr std::size_t kMaxClients = 10000;

    真实可支持数量还受以下因素约束:

    • RLIMIT_NOFILE;
    • 系统级文件描述符限制;
    • epoll watch 限制;
    • 每连接用户态状态;
    • socket 缓冲区;
    • 进程内存;
    • CPU 和业务处理成本。

    因此“epoll 没有 1024 限制”不等于“可以无限建立连接”。

    单连接发送队列硬上限

    constexpr std::size_t kOutputHardLimit = 1024 * 1024;

    单个客户端不能无限占用内存。

    高低水位背压

    constexpr std::size_t kOutputLowWatermark = 128 * 1024;
    constexpr std::size_t kOutputHighWatermark = 256 * 1024;

    高水位暂停读取,低水位恢复读取。

    每事件读写预算

    constexpr std::size_t kReadBudgetPerEvent = 64 * 1024;
    constexpr std::size_t kWriteBudgetPerEvent = 64 * 1024;

    防止一个活跃连接长时间独占事件循环。

    EPOLLONESHOT 重装

    每轮处理结束后:

    RearmClient(epoll_fd, fd, client);

    把最新的读取、写入、暂停状态重新提交给 epoll。

    21.4 最终层完整流程图

    ┌──────────────────────────────────────────────────────────────┐
    │ epoll_wait │
    └──────────────────────────────┬───────────────────────────────┘


    取出本轮就绪 epoll_event

    ┌─────────────────┴──────────────────┐
    │ │
    ▼ ▼
    fd == listen_fd client fd
    │ │
    ▼ ▼
    循环 accept 到 EAGAIN 检查 EPOLLERR
    │ │
    ▼ ▼
    新 fd 设置非阻塞 按预算循环 recv
    │ │
    ▼ ▼
    注册 EPOLLIN/ET 数据加入 output


    达到高水位?
    │ │
    是 否
    │ │
    ▼ ▼
    暂停读 继续读


    按预算循环 send

    ┌───────────────────────┴─────────────┐
    │ │
    ▼ ▼
    全部发完 EAGAIN/预算耗尽
    │ │
    ▼ ▼
    移除 EPOLLOUT 保留 EPOLLOUT


    低于低水位时恢复 EPOLLIN


    EPOLL_CTL_MOD 重新激活 ONESHOT


    第十二部分:epoll 三大 API 逐项解析

    二十二、epoll_create1()

    int epoll_create1(int flags);

    推荐:

    int epoll_fd = epoll_create1(EPOLL_CLOEXEC);

    EPOLL_CLOEXEC 可以避免程序以后调用 exec() 时意外把 epoll fd 泄漏到新程序。

    返回值:

    >= 0 epoll 实例文件描述符
    -1 失败,检查 errno

    二十三、epoll_ctl()

    int epoll_ctl(int epoll_fd,
    int operation,
    int target_fd,
    epoll_event* event);

    三种常用操作:

    操作作用
    EPOLL_CTL_ADD 增加新的监控项
    EPOLL_CTL_MOD 修改已有监控项的事件掩码或用户数据
    EPOLL_CTL_DEL 删除监控项

    ADD

    epoll_event event{};
    event.events = EPOLLIN;
    event.data.fd = client_fd;

    epoll_ctl(epoll_fd,
    EPOLL_CTL_ADD,
    client_fd,
    &event);

    MOD

    event.events = EPOLLIN | EPOLLOUT | EPOLLET;

    epoll_ctl(epoll_fd,
    EPOLL_CTL_MOD,
    client_fd,
    &event);

    DEL

    epoll_ctl(epoll_fd,
    EPOLL_CTL_DEL,
    client_fd,
    nullptr);

    close() 会在相关 open file description 的引用最终消失时清理 epoll 关联,但显式 DEL 可以让单线程示例中的用户态生命周期更加清楚。复杂程序还必须考虑 dup()、fork() 和描述符复用产生的对象所有权问题。

    二十四、epoll_wait()

    int epoll_wait(int epoll_fd,
    epoll_event* events,
    int max_events,
    int timeout_ms);

    超时参数:

    值含义
    -1 一直等待
    0 不等待,立即返回
    > 0 最多等待对应毫秒数

    返回值:

    返回值含义
    > 0 返回的就绪事件数量
    0 超时
    -1 出错

    信号中断:

    if (count == 1 && errno == EINTR) {
    continue;
    }

    通常不是致命错误。


    第十三部分:常用事件标志

    二十五、事件位速查表

    标志含义常见处理
    EPOLLIN 可能可读 accept()、recv()、读取管道等
    EPOLLOUT 可能可写 继续发送待发送队列
    EPOLLRDHUP 流式连接对端关闭写方向 继续读完剩余数据,等待 recv()==0
    EPOLLERR 异步错误 用 getsockopt(SO_ERROR) 获取错误
    EPOLLHUP 挂断 仍应考虑读取残留数据
    EPOLLET 边缘触发 非阻塞并处理到 EAGAIN
    EPOLLONESHOT 一次交付后禁用 EPOLL_CTL_MOD 重新激活
    EPOLLEXCLUSIVE 某些多 epoll 等待场景的独占唤醒 主要用于缓解特定惊群场景

    EPOLLERR 和 EPOLLHUP 即使没有显式加入兴趣掩码,也可能被返回,因此事件循环必须处理它们。

    二十六、为什么使用 SO_ERROR

    int socket_error = 0;
    socklen_t error_len = sizeof(socket_error);

    getsockopt(fd,
    SOL_SOCKET,
    SO_ERROR,
    &socket_error,
    &error_len);

    EPOLLERR 只是告诉你套接字存在错误状态,SO_ERROR 才用于读取挂起的具体套接字错误码。


    第十四部分:accept、recv、send 的完整状态表

    二十七、非阻塞 accept

    返回含义
    fd >= 0 成功取得一个连接
    -1, EINTR 被信号中断,重试
    -1, ECONNABORTED 队列中的连接已经中止,可继续
    -1, EAGAIN/EWOULDBLOCK 当前连接队列已排空
    其他错误 记录并根据策略处理

    二十八、非阻塞 recv

    返回含义
    n > 0 收到 n 字节
    n == 0 接收方向 EOF
    -1, EINTR 重试
    -1, EAGAIN/EWOULDBLOCK 当前已无数据可读
    其他错误 连接异常

    二十九、非阻塞 send

    返回含义
    n > 0 发送了 n 字节,可能小于请求长度
    -1, EINTR 重试
    -1, EAGAIN/EWOULDBLOCK 当前发送缓冲区满,等待 EPOLLOUT
    -1, EPIPE 对端已经关闭相应方向
    其他错误 连接异常

    第十五部分:TCP 半关闭

    三十、shutdown(SHUT_WR) 的意义

    客户端按下 Ctrl+D 后,本文客户端不会立即 close(),而是:

    shutdown(socket_fd, SHUT_WR);

    这表示:

    客户端:
    不再发送新数据
    但仍然可以接收服务端数据

    对应的 TCP 状态是发送 FIN,但连接的另一个方向仍可继续传输。

    流程:

    客户端发送业务数据


    shutdown(SHUT_WR)


    服务端读到业务数据


    服务端 recv() 最终返回 0


    服务端发送完剩余回显


    服务端 close()


    客户端读到全部回显和 EOF

    这也是为什么服务端不能一看到 EPOLLRDHUP 就立即丢弃发送队列并关闭。


    第十六部分:LT 与 ET 应该如何选

    三十一、对比表

    维度LTET
    默认模式 否,需要 EPOLLET
    编程难度 较低 较高
    非阻塞要求 强烈建议 实际使用中应当采用
    读取策略 可以分批处理 必须确保不会遗留未处理就绪状态
    写入策略 仍要处理部分发送和 EAGAIN 同样需要,且事件状态更敏感
    重复通知 条件持续满足时可能重复 通知更依赖状态变化
    出错风险 相对低 忘记排空、错误暂停更容易停滞
    适用场景 大多数服务端的安全起点 明确理解状态机并需要精细控制时

    三十二、ET 一定比 LT 快吗

    不一定。

    性能取决于:

    • 活跃连接比例;
    • 每次事件的数据量;
    • 系统调用频率;
    • 业务处理成本;
    • 缓冲区大小;
    • 事件循环公平性;
    • 内存分配;
    • 锁竞争;
    • 网络延迟;
    • CPU 缓存行为。

    在许多服务中,业务逻辑、序列化、数据库和日志的成本远高于 LT/ET 的差异。

    选择原则应是:

    先实现正确、可测量的 LT


    通过性能分析确认事件调度成为瓶颈


    再评估 ET 是否能带来实际收益

    不要仅凭“ET 更高级”选择 ET。


    第十七部分:压力测试客户端

    三十三、完整大数据校验客户端

    该客户端:

  • 生成指定大小的 X 字节;
  • 完整发送;
  • shutdown(SHUT_WR);
  • 读取服务端全部回显;
  • 逐字节校验。
  • #include <arpa/inet.h>
    #include <sys/socket.h>
    #include <unistd.h>

    #include <cerrno>
    #include <cstdio>
    #include <cstdlib>
    #include <iostream>
    #include <string>

    namespace {
    bool SendAll(int fd,
    const char* data,
    std::size_t size) {
    std::size_t offset = 0;

    while (offset < size) {
    ssize_t n = send(fd,
    data + offset,
    size offset,
    MSG_NOSIGNAL);

    if (n > 0) {
    offset += static_cast<std::size_t>(n);
    continue;
    }

    if (n == 1 && errno == EINTR) {
    continue;
    }

    std::perror("send");
    return false;
    }

    return true;
    }
    }

    int main(int argc, char* argv[]) {
    std::size_t message_size = 1024 * 1024;

    if (argc >= 2) {
    message_size =
    static_cast<std::size_t>(
    std::strtoull(argv[1], nullptr, 10));
    }

    int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
    if (socket_fd == 1) {
    std::perror("socket");
    return 1;
    }

    sockaddr_in server_addr{};
    server_addr.sin_family = AF_INET;
    server_addr.sin_port = htons(8888);

    if (inet_pton(AF_INET,
    "127.0.0.1",
    &server_addr.sin_addr) != 1) {
    std::cerr << "invalid IPv4 address\\n";
    close(socket_fd);
    return 1;
    }

    if (connect(socket_fd,
    reinterpret_cast<sockaddr*>(&server_addr),
    sizeof(server_addr)) == 1) {
    std::perror("connect");
    close(socket_fd);
    return 1;
    }

    std::string payload(message_size, 'X');

    if (!SendAll(socket_fd,
    payload.data(),
    payload.size())) {
    close(socket_fd);
    return 1;
    }

    if (shutdown(socket_fd, SHUT_WR) == 1) {
    std::perror("shutdown");
    close(socket_fd);
    return 1;
    }

    std::string received;
    received.reserve(payload.size());

    while (true) {
    char buffer[16 * 1024];
    ssize_t n = recv(socket_fd,
    buffer,
    sizeof(buffer),
    0);

    if (n > 0) {
    received.append(
    buffer,
    static_cast<std::size_t>(n));
    continue;
    }

    if (n == 0) {
    break;
    }

    if (errno == EINTR) {
    continue;
    }

    std::perror("recv");
    close(socket_fd);
    return 1;
    }

    close(socket_fd);

    bool ok = received == payload;

    std::cout << "sent=" << payload.size()
    << ", received=" << received.size()
    << ", verify=" << (ok ? "OK" : "FAILED")
    << '\\n';

    return ok ? 0 : 1;
    }

    33.1 测试 1 MiB 回显

    ./bin/05_epoll_et_server

    另一个终端:

    ./bin/stress_client 1048576

    预期:

    sent=1048576, received=1048576, verify=OK

    33.2 为什么不能用一次 write 作为压力测试

    错误写法:

    write(fd, big_message.data(), big_message.size());

    如果返回值小于 big_message.size(),剩余数据没有被发送。

    压力测试本身也必须正确处理部分发送,否则无法区分:

    • 服务端错误;
    • 客户端发送不完整。

    第十八部分:编译与运行

    三十四、编译全部版本

    make

    使用的编译参数:

    g++ -std=c++17 -O2 -Wall -Wextra -Wpedantic

    三十五、运行指定层次

    第 0 层

    ./bin/00_select_server
    ./bin/00_select_client

    第 1 层

    ./bin/01_epoll_lt_server
    ./bin/01_epoll_lt_client

    第 2 层

    ./bin/02_epoll_lt_server
    ./bin/02_epoll_lt_client

    第 3 层

    ./bin/03_epoll_et_server
    ./bin/03_epoll_et_client

    第 4 层

    ./bin/04_epoll_et_server
    ./bin/04_epoll_et_client

    第 5 层

    ./bin/05_epoll_et_server
    ./bin/05_epoll_et_client

    客户端输入结束:

    Ctrl+D

    客户端会半关闭写方向,但继续读取服务端回显。


    第十九部分:调试方法

    三十六、使用 strace 观察事件循环

    strace -f \\
    -e trace=epoll_create1,epoll_ctl,epoll_wait,accept,recvfrom,sendto,close \\
    ./bin/05_epoll_et_server

    重点观察:

    epoll_wait(…)
    accept(…)
    recvfrom(…)
    sendto(…)
    epoll_ctl(… EPOLL_CTL_MOD …)

    三十七、使用 ss 查看连接

    ss -ltnp 'sport = :8888'

    查看监听状态:

    LISTEN

    查看已建立连接:

    ss -tnp 'sport = :8888 or dport = :8888'

    三十八、模拟慢客户端

    可以让客户端发送大量数据,却暂时不读取回显,从而观察:

    • 服务端发送队列增长;
    • 达到高水位后暂停 EPOLLIN;
    • 发送队列下降后恢复读取;
    • 超过硬上限后的断开策略。

    第二十部分:常见错误清单

    三十九、错误:ET 只加 EPOLLET,不改读取方式

    events = EPOLLIN | EPOLLET;
    read(fd, buf, sizeof(buf)); // 只读一次

    后果:未读数据可能滞留,处理停滞。

    四十、错误:ET 使用阻塞 fd 循环读

    while (read(fd, buf, sizeof(buf)) > 0) {
    }

    后果:缓冲区读空后阻塞整个事件循环。

    四十一、错误:把 EAGAIN 当成连接失败

    if (recv(...) == 1) {
    close(fd);
    }

    后果:正常的非阻塞控制流被误判为错误。

    四十二、错误:永久监听 EPOLLOUT

    events = EPOLLIN | EPOLLOUT;

    后果:没有待发送数据时仍可能持续被可写事件唤醒。

    四十三、错误:忽略部分发送

    send(fd, data, size, 0);

    后果:只发送前半部分,应用却以为全部完成。

    四十四、错误:发送队列没有上限

    后果:慢客户端可以使进程内存持续增长。

    四十五、错误:ET 设置处理预算但不重新激活

    后果:预算耗尽时缓冲区可能仍然可读,但没有新边缘触发。

    四十六、错误:EPOLLRDHUP 立即 close

    后果:可能丢弃 FIN 之前已经到达但尚未读取的数据,或丢弃尚未发送的响应。

    四十七、错误:在 epoll_event.data.ptr 中保存裸指针却不管理生命周期

    如果连接对象被释放,而 ready event 中仍保留旧指针,就可能出现 use-after-free。

    本文使用 data.fd 和 unordered_map<int, ClientState>,降低了示例中的生命周期复杂度。生产代码可以使用指针,但必须结合:

    • 稳定对象地址;
    • 引用计数;
    • generation 标识;
    • 延迟销毁;
    • 线程同步。

    第二十一部分:进一步工程化方向

    四十八、应用层协议拆包

    Echo 示例把收到的字节原样返回,不代表一次 recv() 就是一条业务消息。

    生产服务通常需要:

    • 固定长度协议;
    • 分隔符协议;
    • 长度字段 + Payload;
    • HTTP/WebSocket 等成熟协议解析器。

    长度头示意:

    ┌──────────────┬─────────────────────┐
    │ 4 字节长度 N │ N 字节 Payload │
    └──────────────┴─────────────────────┘

    每连接需要独立接收缓冲区:

    struct ClientState {
    std::string input;
    std::string output;
    };

    解析原则:

    input 不够一帧:
    保留,等待下次 recv

    input 包含完整一帧:
    取出并处理

    input 包含多帧:
    循环解析

    四十九、定时器

    可以使用:

    • timerfd;
    • 时间轮;
    • 最小堆;
    • 周期性 epoll_wait 超时。

    用于处理:

    • 空闲连接超时;
    • 握手超时;
    • 请求处理超时;
    • 心跳。

    五十、信号

    可以使用 signalfd 把信号也变成 epoll 事件,避免传统异步信号处理函数中的限制。

    五十一、多线程

    常见模型包括:

    单 acceptor + 多 event loop
    每线程一个 epoll
    主从 Reactor
    EPOLLONESHOT + 工作线程
    SO_REUSEPORT 多监听实例

    多线程下要额外解决:

    • 连接归属;
    • 同一 fd 是否会并发处理;
    • 对象销毁时机;
    • 跨线程唤醒;
    • 发送队列锁;
    • fd 复用;
    • 惊群。

    五十二、跨线程唤醒

    常使用:

    • eventfd;
    • pipe;
    • socketpair。

    工作线程向 event loop 提交任务后,通过 eventfd 唤醒 epoll_wait()。


    第二十二部分:FAQ

    五十三、close 前必须 EPOLL_CTL_DEL 吗

    简单单线程程序中,关闭文件描述符的最终引用后,内核会清理相关 epoll 监控关系。

    但显式 DEL 仍有价值:

    • 用户态生命周期清晰;
    • 错误更容易定位;
    • 复杂的 dup()/共享 open file description 场景中更容易控制;
    • 避免连接状态容器和 epoll 兴趣集合长时间不一致。

    不能简单总结成“close 一定会产生幽灵事件”或“DEL 永远必须”。

    五十四、accepted socket 会继承 O_NONBLOCK 吗

    在 Linux 上,accept() 返回的新套接字不会自动继承监听套接字的 O_NONBLOCK 文件状态标志。

    因此需要:

    fcntl(client_fd, F_SETFL, ... | O_NONBLOCK);

    或者:

    accept4(..., SOCK_NONBLOCK | SOCK_CLOEXEC);

    五十五、EAGAIN 和 EWOULDBLOCK 是否相同

    Linux 上通常相同,但可移植代码应同时检查:

    errno == EAGAIN ||
    errno == EWOULDBLOCK

    五十六、LT 可以使用阻塞套接字吗

    从接口语义上可以,但事件驱动服务器通常仍建议使用非阻塞套接字,以保证任何 I/O 操作都不会意外长时间阻塞整个事件循环。

    五十七、ET 是否每次都必须读到 EAGAIN

    最稳定、最容易验证的通用规则是:

    非阻塞读取到 EAGAIN

    对流式 fd,某些情况下短读也可以说明当前数据已经取尽,但把所有 ET 路径统一写成“循环到 EAGAIN”更不容易出错。

    如果为了公平预算提前停止,就必须使用用户态重新调度或 EPOLLONESHOT 重装等机制,不能直接回到普通 epoll_wait() 后假设一定会有新边缘。

    五十八、epoll 可以监控普通文件吗

    常规磁盘文件通常不适合通过 epoll 获取就绪通知,epoll_ctl() 可能返回 EPERM。epoll 主要面向支持轮询语义的对象,例如 socket、pipe、eventfd、timerfd、signalfd 等。

    五十九、epoll 能替代线程池吗

    不能。

    epoll 解决的是:

    大量 I/O 对象的就绪等待

    线程池解决的是:

    CPU 密集型或可能阻塞的业务任务并发执行

    两者经常组合使用。


    第二十三部分:六层演进总结

    六十、每一层到底学到了什么

    第 0 层:select

    理解:

    • fd 集合;
    • select() 破坏性修改;
    • 全量扫描;
    • FD_SETSIZE。

    第 1 层:epoll LT 最小迁移

    理解:

    • epoll 实例;
    • interest list;
    • ready events;
    • epoll_create1/ctl/wait。

    第 2 层:非阻塞 LT

    理解:

    • 非阻塞套接字;
    • EINTR;
    • EAGAIN;
    • 排空连接队列;
    • LT 也应避免阻塞事件循环。

    第 3 层:ET 正确读取

    理解:

    • EPOLLET;
    • 单次 read 的停滞问题;
    • 阻塞循环读的事件循环卡顿;
    • 非阻塞 + 处理到 EAGAIN。

    第 4 层:写事件状态机

    理解:

    • 部分发送;
    • 发送偏移;
    • EPOLLOUT 按需订阅;
    • 半关闭;
    • MSG_NOSIGNAL。

    第 5 层:工程边界

    理解:

    • 背压;
    • 高低水位;
    • 内存硬上限;
    • 每事件预算;
    • EPOLLONESHOT;
    • 公平性。

    六十一、最终核心原则

    可以把整篇文章压缩成六句话:

    1. epoll 只报告就绪,不替你执行 I/O。
    2. 非阻塞 I/O 的 EAGAIN 是正常边界,不是错误。
    3. ET 收到事件后,要把当前就绪状态处理完整。
    4. TCP send 可能部分完成,必须保存发送状态。
    5. EPOLLOUT 只在存在待发送数据时监听。
    6. 没有背压和资源上限的发送队列不是完整设计。


    附录 A:源码目录

    epoll_分层升级技术博客/
    ├── 00_select_original/
    │ ├── server.cpp
    │ └── client.cpp
    ├── 01_epoll_lt_minimal/
    │ ├── server.cpp
    │ └── client.cpp
    ├── 02_epoll_lt_nonblocking/
    │ ├── server.cpp
    │ └── client.cpp
    ├── 03_epoll_et_correct_read/
    │ ├── server.cpp
    │ └── client.cpp
    ├── 04_epoll_et_write_queue/
    │ ├── server.cpp
    │ └── client.cpp
    ├── 05_epoll_et_backpressure/
    │ ├── server.cpp
    │ └── client.cpp
    ├── pitfalls/
    │ ├── et_single_read_server.cpp
    │ └── et_blocking_loop_server.cpp
    ├── tools/
    │ └── stress_client.cpp
    ├── Makefile
    ├── README.md
    └── 从_select_到_epoll_C++17_完整技术博客.md

    附录 B:官方资料名称

    进一步核对接口语义时,应优先查阅 Linux man-pages 中的:

    • epoll(7)
    • epoll_create(2)
    • epoll_ctl(2)
    • epoll_wait(2)
    • accept(2)
    • recv(2)
    • send(2)
    • socket(7)
    • fcntl(2)

    本文的重点不是记住某份最终代码,而是建立可以迁移到真实网络项目的状态机思维:

    每个连接都不是一个“fd 数字”,而是一个持续变化的 I/O 状态对象;epoll 只是让这些状态对象在合适的时机获得处理机会。

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 11.升级简单的Epoll,从 select 到 epoll:C++17 TCP 服务器的六层渐进式重构
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!