跳到内容

2.2 消息边界、Deadline 与 Backpressure

第一条回显连接已经通了。你把两封法术信函接连塞进传送口,对面的信标塔却只收到一条连续的字节带:信封不见了,封口也不见了。TCP 忠实地交付字节,却从未答应替应用保存两次 send 的边界。

驿道管理员必须自己规定怎样拆信。可以在末尾画分隔符,也可以先写长度;无论选哪一种,都要限制一封信有多大、半封信等多久、接收大厅堆满以后谁先停手。映射回正式术语,这四件事分别是 framing、size limit、deadline 和 backpressure。

这一课实现四字节长度前缀协议。例子会故意一次只写一个 byte,逼 parser 面对任意分段;随后把 timeout、deadline、pipelining 和 backpressure 接到同一条连接状态机里。

1. 四种常见 framing

固定长度

每条消息恰好 N 字节。解析最简单,但短消息浪费空间,变长字段需要额外规则。适合硬件记录或已知尺寸 block。

分隔符

以换行或特殊序列结束:

text
PING\r\n
SET key value\r\n

实现必须处理分隔符跨越两次 recv、payload 中转义、最大行长和未结束行。不能无限等待并增长 buffer。

长度前缀

text
[4-byte length][payload]

可以传二进制和空 payload。读取 header 后必须先验证上限,再分配和读取 body;否则一个 0xffffffff 就可能触发内存耗尽。

自描述格式

JSON object、HTTP message 或 Protobuf 等有自己的语法/长度规则。它们仍需要增量 parser,并在不完整输入时保存状态。不能假设一次 recv 得到一个完整 JSON。

2. 正确读取固定字节数

recv_exact 必须区分:

  • 收满 N 字节;
  • 在任何字节前遇到 EOF,表示消息流正常结束;
  • 读了一部分后遇到 EOF,表示截断 frame。
python
import socket
import struct
import threading

HEADER = struct.Struct("!I")
MAX_FRAME = 1024 * 1024


def recv_exact(connection, length, allow_clean_eof=False):
    data = bytearray()
    while len(data) < length:
        chunk = connection.recv(length - len(data))
        if chunk == b"":
            if allow_clean_eof and len(data) == 0:
                return None
            raise EOFError(
                f"stream ended after {len(data)} of {length} bytes"
            )
        data.extend(chunk)
    return bytes(data)


def encode_frame(payload):
    if len(payload) > MAX_FRAME:
        raise ValueError("frame too large")
    return HEADER.pack(len(payload)) + payload


def recv_frame(connection):
    header = recv_exact(
        connection,
        HEADER.size,
        allow_clean_eof=True,
    )
    if header is None:
        return None

    (length,) = HEADER.unpack(header)
    if length > MAX_FRAME:
        raise ValueError(f"declared frame too large: {length}")
    return recv_exact(connection, length)


def fragmented_writer(connection, frames):
    with connection:
        encoded = b"".join(encode_frame(frame) for frame in frames)
        for byte in encoded:
            connection.sendall(bytes([byte]))
        connection.shutdown(socket.SHUT_WR)


left, right = socket.socketpair()
expected = [b"alpha", b"", b"omega"]
writer = threading.Thread(
    target=fragmented_writer,
    args=(left, expected),
)
writer.start()

with right:
    actual = []
    while True:
        frame = recv_frame(right)
        if frame is None:
            break
        actual.append(frame)

writer.join()
assert actual == expected
print(actual)

socketpair 不经过真实网络,但逐字节发送稳定地证明 parser 不依赖 packet 或 sendall 边界。!I 表示 network byte order 的 unsigned 32-bit length。

3. Header 也可能被拆开

不少错误实现会写:

python
length = struct.unpack("!I", connection.recv(4))[0]

recv(4) 可以只返回 1–3 字节。即使当前 header 都到了,signal、调度和 buffer 状态也会让读取结果不同。header 与 body 都要走精确读取循环。

另一种错误是把 EOF 当成“暂时没数据”:

python
if connection.recv(4096) == b"":
    continue

这会在已关闭连接上忙循环。blocking socket 暂时没数据时会阻塞或 timeout;返回空字节表示有序 EOF。

4. 长度字段属于不可信输入

读取到长度后,在分配前检查:

  • 是否超过协议最大 frame;
  • 是否允许零长度;
  • header 算法是否可能整数溢出;
  • 解压后的大小是否另有上限;
  • 一个连接允许多少在途 frame;
  • 总 buffer 是否受全局预算控制。

“最大 1 MiB”只是示例。真实上限应由业务语义、内存预算与代理链限制共同决定。

压缩协议还要防高压缩比 payload。wire length 很小不等于解压内存小;应同时限制压缩输入、解压输出和 CPU 工作量。

5. Per-operation timeout 与总 deadline

给每次 recv 设置 5 秒 timeout,不代表整个请求最多 5 秒。对端可以每 4 秒发一个字节,让一条 1 MiB 消息持续数周。

端到端 deadline 用 monotonic clock 计算剩余预算:

python
import time


def recv_exact_before(connection, length, deadline):
    data = bytearray()
    while len(data) < length:
        remaining = deadline - time.monotonic()
        if remaining <= 0:
            raise TimeoutError("message deadline exceeded")
        connection.settimeout(remaining)

        chunk = connection.recv(length - len(data))
        if chunk == b"":
            raise EOFError("truncated message")
        data.extend(chunk)
    return bytes(data)

这段函数表达总预算,但还不是完整生产策略:

  • timeout 后连接中的半条 frame 如何处理?
  • 是否关闭连接,还是协议支持安全 resynchronize?
  • request deadline 是否包含排队和下游调用?
  • response 已部分发送时能否重试?
  • 取消怎样传递给 worker 和 I/O?

deadline 是状态机的一部分,不是给 socket 加一个数字。

6. Backpressure:接收大厅满了以后

长度前缀解决了怎样拆信,却没有扩大接收大厅。若对端拆信很慢,信函会先堆在 socket buffer,再堆到应用自己的发送队列。继续无上限接单,只是把驿道拥堵搬进进程内存。

回到正式路径,当应用调用 send

  1. 字节先进入本机 socket send buffer;
  2. TCP 按接收窗口与拥塞窗口发送;
  3. 对端内核放入 receive buffer;
  4. 对端应用最终 recv

若对端读得慢,buffer 会逐级填满。blocking sendall 最终阻塞,nonblocking send 返回 would-block,async writer 的待发送队列会增长。

错误做法是每个 producer 无限制地把响应追加进用户态列表。这样把网络 backpressure 变成进程 OOM。

需要明确:

  • 每连接最大 queued bytes;
  • 全局最大 queued bytes;
  • 超限时暂停读取上游、拒绝新请求还是关闭慢连接;
  • 哪些消息可以丢弃或合并;
  • 写入 deadline;
  • 公平调度,避免一个大响应饿死小响应。

在 asyncio 等框架里,write 往往只把数据交给 transport buffer;相应的 drain/await 或 high-water mark 才参与 backpressure。具体 API 应以 runtime 文档为准。

7. Request/response 顺序与 request ID

最简单协议每连接一次只允许一个请求:

text
send request A
receive response A
send request B
receive response B

pipelining 允许 A、B 同时在途。如果 response 必须按请求顺序返回,慢 A 会阻塞已完成的 B;如果允许乱序返回,每条 frame 就要带 request ID:

text
[length][request-id][type][payload]

还要定义:

  • ID 是否可复用,何时可复用;
  • response、error、cancel 怎样关联;
  • server push 使用哪个 ID 空间;
  • 同一 ID 重复出现如何处理;
  • 连接重建后旧 ID 是否仍有意义。

HTTP/2、gRPC 与 QUIC stream 已经解决了大量多路复用状态,普通应用不应轻易重造。

8. Blocking、nonblocking 与 readiness

nonblocking socket 在暂时不能完成时返回 would-block。epoll/kqueue 等 readiness API 告诉 event loop 某 fd “现在可能可读/可写”,不保证下一次操作一定完成全部请求。

事件处理器仍要:

  • 循环 read 到 would-block 或达到公平预算;
  • 保存半个 header/body 的 parser 状态;
  • 只在有待发送数据时关注 writable;
  • 处理 hangup/error 与剩余可读数据;
  • 避免一个永远活跃的连接垄断 loop;
  • 在关闭前取消 timer 和业务任务。

edge-triggered 模式通常要求 drain 到 would-block,否则可能等不到下一次边沿。level-triggered 会在条件仍成立时继续通知,实现更直观但也可能反复唤醒。

async/await 把状态机藏进编译器/runtime,不会消除这些协议规则。

9. TLS 会再增加一层 framing

TLS 把应用字节包装成 record。一次应用 write 可能对应多个 TLS record 和多个 TCP segment;一次 TCP recv 也可能只包含半个 TLS record。

应用应通过 TLS library 的 read/write API,让库维护 record、认证和重组状态,不要对加密 socket 的底层 TCP 字节自行切片。TLS clean shutdown 与 TCP EOF 也有区别,安全敏感协议需要判断是否收到规范的 close notification。

10. UDP framing 的不同边界

UDP 保留 datagram boundary,所以一个长度前缀通常不用于跨多个 recvfrom 重组同一个 datagram。但应用仍可能在一个 datagram 内放多条 record,或把大消息分成多个 datagram。

自定义 UDP 分片需要 message ID、片号、总数、timeout、去重、内存上限和拥塞控制。丢一个片可能让整个消息不可用。若需求是可靠、安全、多路流传输,优先评估 QUIC 而不是逐步复刻它。

11. 协议测试矩阵

不要只测“正常发送一条消息”。至少覆盖:

输入预期
header 每次只到 1 字节正确重组
body 分成任意 chunk正确重组
两个 frame 合并到一次 recv解析两条
零长度 frame按协议接受或拒绝
header 中途 EOF截断错误
body 中途 EOF截断错误
length 超上限分配前拒绝
慢速逐字节输入总 deadline 生效
对端不读 responsequeued bytes 有界
cancel 与 close 同时发生只释放一次,状态一致

property-based test 可以随机切分同一编码字节串,验证 parser 对所有 chunk boundary 得到相同消息序列。

12. 小结

TCP 协议设计必须主动补上字节流没有提供的东西:

  • framing 定义消息边界;
  • recv_exact 处理 header/body 的任意切分;
  • 长度和解压结果必须在分配前受限;
  • operation timeout 不能代替 end-to-end deadline;
  • backpressure 让慢下游限制上游,而不是无限堆内存;
  • pipelining 需要 request ID、取消和响应顺序规则;
  • readiness/async 仍要维护半包、错误与关闭状态;
  • 测试要随机化 chunk boundary,并覆盖截断与慢客户端。

到这里,信标塔管理员已经给连续字节划出消息边界,也知道大厅塞满时必须让上游停手。下一章打开驿道观测镜,先看 TCP 连接怎样建立和关闭,再追踪序号、ACK、流量控制与重传。

Built with VitePress | Software Systems Atlas