2.2 消息边界、Deadline 与 Backpressure
第一条回显连接已经通了。你把两封法术信函接连塞进传送口,对面的信标塔却只收到一条连续的字节带:信封不见了,封口也不见了。TCP 忠实地交付字节,却从未答应替应用保存两次 send 的边界。
驿道管理员必须自己规定怎样拆信。可以在末尾画分隔符,也可以先写长度;无论选哪一种,都要限制一封信有多大、半封信等多久、接收大厅堆满以后谁先停手。映射回正式术语,这四件事分别是 framing、size limit、deadline 和 backpressure。
这一课实现四字节长度前缀协议。例子会故意一次只写一个 byte,逼 parser 面对任意分段;随后把 timeout、deadline、pipelining 和 backpressure 接到同一条连接状态机里。
1. 四种常见 framing
固定长度
每条消息恰好 N 字节。解析最简单,但短消息浪费空间,变长字段需要额外规则。适合硬件记录或已知尺寸 block。
分隔符
以换行或特殊序列结束:
PING\r\n
SET key value\r\n实现必须处理分隔符跨越两次 recv、payload 中转义、最大行长和未结束行。不能无限等待并增长 buffer。
长度前缀
[4-byte length][payload]可以传二进制和空 payload。读取 header 后必须先验证上限,再分配和读取 body;否则一个 0xffffffff 就可能触发内存耗尽。
自描述格式
JSON object、HTTP message 或 Protobuf 等有自己的语法/长度规则。它们仍需要增量 parser,并在不完整输入时保存状态。不能假设一次 recv 得到一个完整 JSON。
2. 正确读取固定字节数
recv_exact 必须区分:
- 收满 N 字节;
- 在任何字节前遇到 EOF,表示消息流正常结束;
- 读了一部分后遇到 EOF,表示截断 frame。
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 也可能被拆开
不少错误实现会写:
length = struct.unpack("!I", connection.recv(4))[0]recv(4) 可以只返回 1–3 字节。即使当前 header 都到了,signal、调度和 buffer 状态也会让读取结果不同。header 与 body 都要走精确读取循环。
另一种错误是把 EOF 当成“暂时没数据”:
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 计算剩余预算:
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:
- 字节先进入本机 socket send buffer;
- TCP 按接收窗口与拥塞窗口发送;
- 对端内核放入 receive buffer;
- 对端应用最终
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
最简单协议每连接一次只允许一个请求:
send request A
receive response A
send request B
receive response Bpipelining 允许 A、B 同时在途。如果 response 必须按请求顺序返回,慢 A 会阻塞已完成的 B;如果允许乱序返回,每条 frame 就要带 request ID:
[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 生效 |
| 对端不读 response | queued 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、流量控制与重传。