CHAPTER 28 / Linux 与并发系统

TCP 字节流、分包、超时与背压

两次 send 为什么可能变成三次 recv,如何写出不会无界等待的接收器?

阅读与推演约 60 分钟练习时间另计

这一章要弄清楚

  • 为 TCP 字节流建立应用分帧协议
  • 正确推进部分 I/O 与 EOF 状态
  • 把时间与内存预算加入网络协议

先备知识:进程、系统调用与文件描述符 / 条件变量、有界队列与关闭协议

C++20;24-a 仅标准库,24-b 需要 macOS/Linux POSIX sockets,建立一次仅 127.0.0.1 的 TCP 连接,端口由系统分配。

网络交付的是字节顺序

应用想发送两条消息“cat”和“hi”,TCP 负责可靠、有序地传输字节流,但不会保留 send 调用的消息边界。发送端写了两次,接收端可能一次得到全部,也可能只得到半个消息,或者第一条完整消息加第二条的一半。这不是 TCP 损坏数据,而是应用把字节流接口误当成消息接口。分帧(framing)必须由应用协议定义。

本章使用一个刻意很小的格式:一个无符号字节表示负载长度,后面紧跟负载;负载上限为 8 字节。真实协议常用固定宽度大端长度头或分隔符,还要约定版本、最大长度与字符编码。长度包括头部还是仅负载必须说清楚,零长度也需要明确。本章允许零长度消息,因此长度零会生成一条空消息,而不是连接结束。

把接收写成状态机

解析器保留“正在等长度”或“正在等若干负载字节”的状态。每收到一段就消费能确定的部分;不足一帧时保留进度,下次继续;一段里有多帧就全部提取。先检查长度上限,再分配或累积,不要允许对端宣布巨大长度后耗尽内存。示例还为一次连接示范设置消息条数上限,使零长度洪泛也不能无界增长;真实服务通常处理完一帧就交给有界下游队列。

EOF 是连接读取方向结束,不自动代表最后一条消息完整。如果正等待负载却收到 EOF,应报告截断;如果正等待新头,才是干净结束。一个 TCP 连接可以半关闭:对端不再发送,但仍可接收。连接重置、超时和正常 EOF 应分别记录,便于调用方决定是否重试;重试有副作用的请求还要考虑重复执行。

部分 I/O 需要保存进度

按实际返回数量推进的基本合同见SYS-00 S4,标志的或、与读法见SYS-00 S1。下面继续应用到网络;fcntl、poll、连接状态、EAGAIN、超时和SIGPIPE各有额外合同,SYS的手填返回量模型不覆盖这些策略。

对非阻塞 fd,write 或 send 可能只接受一部分字节。下一次必须从已发送偏移继续,否则重复前缀;遇到 EAGAIN 等待可写,不是把它当永久失败。读也是一样:返回正数就处理实际长度,零表示 EOF,EINTR 按剩余时间重新评估,其他错误分类处理。面对已关闭连接,写操作还要按平台策略处理 SIGPIPE;本例始终保持读端有效,未把这一问题藏进测试。

第二个例子使用真正的本机 TCP loopback:只绑定 127.0.0.1,由端口 0 请求系统分配临时端口,再用 getsockname 查询。非阻塞 connect 进入进行中状态后,等待可写并检查 SO_ERROR,不能把可写本身当作连接成功。建立连接后发送六个字节并半关闭写方向,接收端每次最多读两个字节,最终得到 abcdef。这里的两字节是应用的读取窗口,不是 TCP packet 的边界。测试没有跨机器网络、可控拥塞或重传注入,因此不能用其时间推断远端系统性能。

时间预算与背压一起设计

超时最好用单调时钟的绝对截止时间,而不是每次重试都重新给五秒。否则恶意或缓慢对端每隔四秒发一个字节,就能无限延长请求。poll 返回可读后仍应准备面对 EOF、错误或其他消费者抢走数据;非阻塞模式让等待与读取不会把线程意外卡住。生产代码还需要连接、头部、整体请求或空闲超时的明确策略。

背压把下游拥堵传回读端:处理队列到达高水位时暂停继续读取,降到低水位再恢复。TCP 接收缓冲最终也会反馈流量,但若应用先无界读入内存,就绕过了这种保护。面试设计题应同时画出字节偏移、分帧状态、截止时间和队列容量;这四个数字比一句“用 epoll”更能说明系统是否可控。

动手改变 · 观察因果

同一字节流,换一种 read 分块

本次输入后能交付几帧?还缺多少字节?

分块表示 read 返回的数据段,不是 TCP packet。模型验证长度前缀解析;没有网络、重传或拥塞模拟。 对应 例题 24-a;图中的代码行是步骤提示,完整可编译源码见例题。

改输入后从第一步重新推演。Tab 选择控件,Enter/空格操作按钮;图内方向键平移,手机可横向滑动。

正在准备默认算例。下方例题包含完整源码与逐步解释。

第 1 步

先预测,再前进一步

本次输入后能交付几帧?还缺多少字节?

输入、边界、状态变化

完整文字推演与当前数据
    静态推演与完整文字(便于对照、打印)
    观察 · 推演

    cat 与 hi 的分段到达

    收到3,c负载c还需要201 / 03 · NETWORK收到3,c负载c还需要201 / 03 · NETWORK
    只到头和 c

    长度 3 已知,负载还差两个字节。

    1 / 3
    阅读完整推演文字
    1. 只到头和 c

      收到:3,c;负载:c;还需要:2

      长度 3 已知,负载还差两个字节。

    2. at 和新帧头

      收到:a,t,2;已提交:cat;还需要:2

      提交 cat,下一帧长度 2;当前负载为空。

    3. 收到 hi

      收到:h,i;已提交:cat / hi;状态:等待新头

      第二帧完整,回到等待长度状态。

    读完整TCP程序之前

    先完成NET-24:从socket到一次完整传输,逐项读懂地址、描述符拥有、就绪事件与返回量合同,再回到下面的完整示例。

    跟着例子,走完一遍

    例题 01C++20 · 本机可运行

    跨任意分块解析两条消息

    字节为 [3,c,a,t,2,h,i],分别以 1 字节和整段输入,并测试空帧、超长和截断。

    1. 长度先校验,最大为 8。
    2. 逐字节推进负载状态,满一帧立即提交。
    3. EOF 检查未完成状态。
    24-a.cpp
    下载
    #include <cstdint>
    #include <iostream>
    #include <stdexcept>
    #include <string>
    #include <vector>
    class Parser {
        int remaining=-1; std::string current;
    public:
        std::vector<std::string> frames;
        void feed(const std::vector<std::uint8_t>& bytes) {
            for(auto b:bytes) {
                if(remaining<0) {
                    if(b>8 || frames.size()>=16) throw std::runtime_error("limit");
                    remaining=b; current.clear();
                    if(remaining==0) {frames.emplace_back(); remaining=-1;}
                } else {
                    current.push_back(static_cast<char>(b));
                    if(--remaining==0) {frames.push_back(current); remaining=-1;}
                }
            }
        }
        void finish() const {if(remaining>=0) throw std::runtime_error("truncated");}
    };
    int main() {
        const std::vector<std::uint8_t> wire{3,'c','a','t',2,'h','i'};
        Parser one,chunk; for(auto b:wire) one.feed({b}); chunk.feed(wire);
        one.finish(); chunk.finish();
        if(one.frames!=chunk.frames || one.frames!=std::vector<std::string>{"cat","hi"})
            throw std::runtime_error("framing");
        int invalid=0;
        for(const auto& bytes:std::vector<std::vector<std::uint8_t>>{{9},{3,'x'}}) {
            try {Parser p;p.feed(bytes);p.finish();} catch(const std::runtime_error&) {++invalid;}
        }
        Parser empty;empty.feed({0});empty.finish();
        if(invalid!=2 || empty.frames!=std::vector<std::string>{""}) throw std::runtime_error("edges");
        std::cout<<"frames=cat,hi invalid="<<invalid<<" empty="<<empty.frames.size()<<'\n';
    }

    如何编译和运行下载的 .cpp 文件 →

    结果与解释

    frames=cat,hi invalid=2 empty=1

    相同字节流的任意分块得到相同消息序列;协议错误独立报告。

    在本机运行这个例子

    下载后,在文件所在目录执行。需要支持 C++20 的编译器;POSIX 示例还需要章节说明中的系统条件。

    clang++ -std=c++20 -Wall -Wextra -Wpedantic -Werror -pthread 24-a.cpp -o example && ./example

    预期标准输出:

    frames=cat,hi invalid=2 empty=1
    
    例题 02POSIX · 需对应系统

    TCP loopback:非阻塞连接与截止时间读取

    仅在 127.0.0.1 建立临时 TCP 连接,发送 abcdef 后 shutdown 写方向;接收端每次最多读取两个字节。

    1. bind 端口 0 后用 getsockname 取得系统分配的端口,listen 与 connect 建立本机连接。
    2. 非阻塞 connect 用 POLLOUT 和 SO_ERROR 确认结果;后续 I/O 都使用同一个单调时钟截止时间。
    3. send/recv 推进已完成字节数,处理 EINTR/EAGAIN、部分结果和 EOF;RAII 关闭三个描述符。
    24-b.cpp
    下载
    #include <sys/socket.h>
    #include <netinet/in.h>
    #include <unistd.h>
    #include <fcntl.h>
    #include <poll.h>
    #include <cerrno>
    #include <chrono>
    #include <iostream>
    #include <stdexcept>
    #include <string>
    
    class Fd {
    public:
        explicit Fd(int value):value_(value) { if(value<0) throw std::runtime_error("socket/accept"); }
        ~Fd(){::close(value_);}
        Fd(const Fd&)=delete;
        Fd& operator=(const Fd&)=delete;
        int get() const{return value_;}
    private:
        int value_;
    };
    using Clock=std::chrono::steady_clock;
    void nonblocking(int fd){
        const int flags=::fcntl(fd,F_GETFL,0);
        if(flags<0 || ::fcntl(fd,F_SETFL,flags|O_NONBLOCK)<0) throw std::runtime_error("fcntl");
    }
    void wait_ready(int fd,short events,Clock::time_point deadline){
        for(;;){
            const auto left=std::chrono::duration_cast<std::chrono::milliseconds>(deadline-Clock::now()).count();
            if(left<=0) throw std::runtime_error("deadline");
            pollfd event{fd,events,0};
            const int result=::poll(&event,1,static_cast<int>(left));
            if(result<0 && errno==EINTR) continue;
            if(result<=0 || (event.revents&(POLLERR|POLLNVAL))) throw std::runtime_error("poll");
            if(event.revents&(events|POLLHUP)) return;
        }
    }
    int main(){
        const auto deadline=Clock::now()+std::chrono::seconds(2);
        Fd listener(::socket(AF_INET,SOCK_STREAM,0));
        sockaddr_in address{};address.sin_family=AF_INET;
        address.sin_addr.s_addr=htonl(INADDR_LOOPBACK);address.sin_port=htons(0);
        if(::bind(listener.get(),reinterpret_cast<sockaddr*>(&address),sizeof address)!=0 ||
           ::listen(listener.get(),1)!=0) throw std::runtime_error("bind/listen");
        socklen_t length=sizeof address;
        if(::getsockname(listener.get(),reinterpret_cast<sockaddr*>(&address),&length)!=0)
            throw std::runtime_error("getsockname");
        nonblocking(listener.get());
        Fd client(::socket(AF_INET,SOCK_STREAM,0));nonblocking(client.get());
        if(::connect(client.get(),reinterpret_cast<sockaddr*>(&address),length)!=0 && errno!=EINPROGRESS)
            throw std::runtime_error("connect");
        wait_ready(client.get(),POLLOUT,deadline);
        int error=0;length=sizeof error;
        if(::getsockopt(client.get(),SOL_SOCKET,SO_ERROR,&error,&length)!=0 || error!=0)
            throw std::runtime_error("connect completion");
        wait_ready(listener.get(),POLLIN,deadline);
        Fd server(::accept(listener.get(),nullptr,nullptr));nonblocking(server.get());
        const std::string input="abcdef";std::size_t sent=0;
        while(sent<input.size()){
            wait_ready(client.get(),POLLOUT,deadline);
            const auto n=::send(client.get(),input.data()+sent,input.size()-sent,0);
            if(n<0 && (errno==EINTR || errno==EAGAIN || errno==EWOULDBLOCK))continue;
            if(n<=0)throw std::runtime_error("send");sent+=static_cast<std::size_t>(n);
        }
        if(::shutdown(client.get(),SHUT_WR)!=0)throw std::runtime_error("shutdown");
        std::string output;char buffer[2];
        for(;;){
            wait_ready(server.get(),POLLIN,deadline);
            const auto n=::recv(server.get(),buffer,sizeof buffer,0);
            if(n<0 && (errno==EINTR || errno==EAGAIN || errno==EWOULDBLOCK))continue;
            if(n<0)throw std::runtime_error("recv");if(n==0)break;
            output.append(buffer,static_cast<std::size_t>(n));
        }
        if(output!=input)throw std::runtime_error("stream");
        std::cout<<"data="<<output<<" eof=1\n";
    }
    结果与解释

    data=abcdef eof=1

    实际验证本机 TCP 字节流与 EOF。读取窗口刻意设为两个字节,不代表网络 packet 边界;未验证公网延迟、拥塞、重传或吞吐。

    在本机运行这个例子

    下载后,在文件所在目录执行。需要支持 C++20 的编译器;POSIX 示例还需要章节说明中的系统条件。

    clang++ -std=c++20 -Wall -Wextra -Wpedantic -Werror -pthread 24-b.cpp -o example && ./example

    预期标准输出:

    data=abcdef eof=1
    
    常见错误

    每个 recv 都直接反序列化一条对象

    头可能只到了一半,也可能一次收到两帧;解析器会读越界或丢掉后半消息。

    修正思路:保持跨读取的分帧状态,限制消息和缓冲区尺寸;EOF 时检查是否存在不完整帧。

    轮到你动手

    先写预测或代码,再按需打开提示。完整答案用于对照自己的推理。

    练习 1

    把 [3,c,a,t,2,h,i] 在每个可能位置切成两段,验证解析结果一致;再测试所有单字节分块。

    给我一点提示
    1. 切分位置包括 0 和总长度。
    2. 比较解析后的消息序列,不比较 read 次数。
    查看答案与推理

    对切分点 0..7,先 feed 前半再 feed 后半并 finish,全部得到 cat、hi。逐字节 feed 也相同;额外测试 0 长度帧以及只含头的截断帧。

    练习 2

    请求总预算 5 秒,每隔 4 秒收到一个字节,如何避免永不超时?

    给我一点提示
    1. 截止时间只在请求开始建立一次。
    2. EINTR、部分读、暂不可读都不能重置预算。
    查看答案与推理

    使用 steady_clock 保存 start+5s,每次等待之前计算剩余时间;到期停止整个请求。可另设 idle timeout,但它不能替代总截止时间。

    把理解说出来

    先用中文讲清因果,再用英文回答。问题依据技能主题编写,并非公司内部题库。

    解释

    TCP 会保留 send 边界吗?

    参考回答 / English answer

    不会,它交付有序字节流。应用必须用长度头、分隔符或其他协议恢复消息边界。

    TCP preserves byte order, not application write boundaries. The application needs its own framing protocol.
    预测

    recv 返回 0 与收到零长度应用帧相同吗?

    参考回答 / English answer

    不同,recv 为零表示读取方向结束;零长度帧由合法协议头表达。

    A zero return indicates end of stream. An empty application message is represented by the framing protocol.
    找错

    send 只返回 20,但请求长度 100,下一步怎么做?

    参考回答 / English answer

    保存偏移 20,从剩余 80 字节继续;暂不可写时等待,不能重新发送全部前缀。

    I advance the send offset by twenty bytes. Subsequent writes start at the unsent suffix.
    追问

    为什么长度头要先检查上限?

    参考回答 / English answer

    不可信长度可触发巨大分配或长期等待;同时限制消息数、总缓冲与截止时间。

    An untrusted length can exhaust memory or keep a request open indefinitely. I bound sizes, backlog, and time.
    找错

    每次 poll 都给五秒,是否等于请求超时五秒?

    参考回答 / English answer

    不等于。部分进度或信号会不断重置预算;应使用单调时钟的绝对截止时间。

    Repeated relative waits can extend the request indefinitely. I recompute each wait from one monotonic deadline.
    追问

    应用层如何落实背压?

    参考回答 / English answer

    限制处理队列,高水位暂停读取或拒绝工作,低水位恢复;避免先无界读入用户内存。

    I bound the downstream queue and pause intake at a high-water mark. Otherwise application buffering defeats transport-level flow control.

    继续查证

    公开资料用于查证;本章图解和例题是独立教学内容。CPU 逻辑模型不能证明设备性能。