Prhub

#50960 [Bugfix] Fix ZMQ port TOCTOU race in shm_broadcast MessageQueue

原始 PR 作者 aoshen02 合并时间 2026-08-08 09:19 文件变更 2 提交数 5 评论 3 代码增减 +68 / -4

执行摘要

修复广播队列端口 TOCTOU 竞态,内核原子分配端口

Issue #28498「[Bug][RL]: Port Conflict」记录了 RL 场景下 DP worker 各自独立找空闲端口的竞态,报告明确指出 _get_open_port 对重复选择同一端口没有任何防御,日志中可见 zmq.error.ZMQError: Address already in use。PR body 用两节点拓扑拆解了根因:每个逻辑队列的 XPUB writer 都先 get_open_port() 探测端口(bind 探测 socket、读端口、关闭),随后才绑定真 socket;探测与绑定之间端口归「无主」,同主机任何进程都能抢占。单实例罕见,但多实例冷启动(RL rollout fleet、multi-instance serving)时探测窗口完全重叠,且 VLLM_PORT 固定时确定性扫描会让并发探测选中同一端口。关键触发条件是 n_reader > n_local_reader(TP/PP/PCP 跨节点即可进入),而非 DP 本身;作者用 8 进程 × 3 轮压测复现:main 上 19/24 构造失败,而同一路径在单实例下几乎不可见。

值得精读,尤其适合负责分布式通信与多实例部署的工程师。源码只有 7 行净改动,但示范了用内核原子语义消除经典 TOCTOU 竞态的标准做法,并处理了「从归一化端点恢复信息」的隐蔽边界——只取端口、用原始地址重建,避免 IPv6 归一化漂移破坏 handle 契约。测试设计(pin 探测端口让旧代码确定性失败)对同类竞态的回归测试有很强的迁移价值。后续可参照此模式清理其余 get_open_port 使用点(如 #50965 正在处理的 DP 保留端口 livelock),并为 VLLM_PORT 行为变化补充文档。

讨论亮点

该 PR 没有人工 review 评论(review_comments_count = 0),claude[bot] 因 fork PR 未自动审查,维护者 ywang96 直接 APPROVED。实质性的设计讨论以 PR body 与 commit message 形式沉淀,要点如下:

  • 为什么只从 LAST_ENDPOINT 取端口、不用整个端点:commit aa9719 说明 LAST_ENDPOINT 返回内核归一化端点,其地址形式可能与配置的 connect_ip 漂移(如 IPv6 规范化),整体采用会改变 reader 连接的 handle 字符串;因此仅解析端口号,地址用原生 connect_ip 重建。这是本 PR 相对 #44495 原始实现的关键增量。

  • 回归测试如何做到确定性失败:commit 429fb1 把 get_open_port 固定为同一个端口,让旧路径每个 writer 都绑定同一端口必然冲突,而晚绑定实现完全不查询探测端口;作者验证 fails on base, passes with the fix,避免竞态类测试 flaky。

  • 避免重复劳动:PR body 明确本 PR 取代停滞的 #44495(合并冲突搁置),并检查过所有引用 shm_broadcast 与 #28498 的 open PR,确认无重复提交;原作者的首个 commit 原封保留、署名与 sign-off 完整,是接手他人 PR 的正面示范。

实现拆解

1. 变更入口与触发路径:核心改动在 vllm/distributed/device_communicators/shm_broadcast.pyMessageQueue.__init__n_remote_reader > 0 分支(远端 XPUB socket 创建),这是广播队列为远端 SUB reader 准备 TCP 监听地址的唯一路径;同时从 vllm.utils.network_utils 导入中移除不再使用的 get_open_port

2. 消除 TOCTOU 窗口:删除「探测-释放-再绑定」流程(先 remote_subscribe_port = get_open_port(),再 bind(f"tcp://{connect_ip}:{port}")),改为直接 self.remote_socket.bind(f"tcp://{connect_ip}:0"),由内核原子分配空闲端口,探测与绑定之间不再存在无主窗口。同时把 is_valid_ipv6_address(connect_ip) 判断与地址括号化前置到 bind 之前,保证 bind 与公告使用同一份规范化后的 connect_ip

3. 端口回读与地址重建:绑定后通过 getsockopt(zmq.LAST_ENDPOINT) 回读实际端点,但只用 last_endpoint.decode().rsplit(":", 1)[1] 取端口号;公告地址仍由原始 connect_ip 拼装为 tcp://{connect_ip}:{port}。原因:LAST_ENDPOINT 返回内核归一化端点,IPv6 等场景下地址形式可能漂移,整体采用会改变 handle 中 remote_subscribe_addr 字符串,导致远端 reader 连向错误地址。

4. 确定性回归测试tests/distributed/test_shm_broadcast.py 新增 test_remote_subscribe_addr_unique_concurrent_writers:用 monkeypatchshm_broadcast.get_open_port 固定返回同一端口,模拟 VLLM_PORT 固定时的确定性竞争——旧路径 32 个并发 writer 必然绑定同一端口失败,新路径不查询探测端口故全部成功;断言 32 个 remote_subscribe_addr 唯一且以 tcp:// 开头,并额外启动 reader 线程完成一次 ping 往返,验证重建地址真实可连。

5. 演进与配套:5 个 commit 中首个为被取代 PR #44495 原作者 RTCartist 的原始实现(原样保留,署名与 sign-off 完整),后续由 aoshen02 补充「仅取端口 + connect_ip 重建地址」与回归测试,期间两次 merge main 同步。无配置、schema 或部署配套改动;压测与 IPv4/IPv6 双栈往返验证结果记录在 PR body。

文件 模块 状态 重要度
vllm/distributed/device_communicators/shm_broadcast.py 广播队列 modified 5.63
tests/distributed/test_shm_broadcast.py 回归测试 modified 6.15

关键符号

MessageQueue.__init__ test_remote_subscribe_addr_unique_concurrent_writers

关键源码片段

vllm/distributed/device_communicators/shm_broadcast.py core-logic

修复核心所在:MessageQueue.__init__ 中远端 XPUB socket 的端口分配从「探测 - 释放 - 再绑定」改为「bind 端口 0 + 从 zmq.LAST_ENDPOINT 回读」,原子消除 TOCTOU 窗口;同时处理 IPv6 判断顺序与公告地址重建。

# vllm/distributed/device_communicators/shm_broadcast.py
# MessageQueue.__init__ 中远端 XPUB 创建分支(n_remote_reader > 0 时进入)
if n_remote_reader > 0:
    # 远端 reader 通过 XPUB-SUB 传输大块数据;XPUB 只 bind,SUB 只 connect
    if not connect_ip:
        connect_ip = get_ip()
    self.remote_socket = context.socket(XPUB)
    self.remote_socket.setsockopt(XPUB_VERBOSE, True)
​
    # IPv6 判断必须前置到 bind 之前:地址需要加方括号后才能构造 tcp 地址
    if is_valid_ipv6_address(connect_ip):
        self.remote_socket.setsockopt(IPV6, 1)
        remote_addr_ipv6 = True
        connect_ip = f"[{connect_ip}]"
​
    # 修复核心:不再 get_open_port() 探测后再绑定(探测 - 释放窗口内端口无主,
    # 同主机并发进程可抢占,多实例冷启动时确定性 EADDRINUSE);
    # 直接 bind 端口 0,由内核原子分配空闲端口,从根上消除 TOCTOU
    self.remote_socket.bind(f"tcp://{connect_ip}:0")
​
    # 只从 LAST_ENDPOINT 取端口号:该值返回内核归一化端点,其地址形式
    # 可能与配置的 connect_ip 不同(如 IPv6 漂移),整体采用会改变
    # reader 连接的 handle 字符串,因此公告地址必须用原始 connect_ip 重建
    last_endpoint = self.remote_socket.getsockopt(zmq.LAST_ENDPOINT)
    remote_subscribe_port = last_endpoint.decode().rsplit(":", 1)[1]
    remote_subscribe_addr = f"tcp://{connect_ip}:{remote_subscribe_port}"
else:
    remote_subscribe_addr = None
    self.remote_socket = None
tests/distributed/test_shm_broadcast.py test-coverage

新增确定性回归测试 test_remote_subscribe_addr_unique_concurrent_writers:monkeypatch 将 get_open_port 固定为同一端口使旧路径必然冲突,32 个并发 writer 断言地址全部唯一,并做 ping 往返验证可连接性。

# tests/distributed/test_shm_broadcast.py(新增回归测试)
def test_remote_subscribe_addr_unique_concurrent_writers(
    monkeypatch: pytest.MonkeyPatch,
):
    """并发 writer 将远端 socket 绑定端口 0(内核分配),
    因此不会竞争同一探测端口,公告地址保证可连接。    修复前 writer 先用 get_open_port() 探测端口再绑定;
    把探测固定到同一个端口,旧路径每次都会绑定同一端口而确定性失败,
    而晚绑定实现完全不查询探测端口——该测试在 base 上必失败、带修复通过。
    """
    from vllm.distributed.device_communicators import shm_broadcast
​
    colliding_port = get_open_port()
    # 固定探测端口:模拟 VLLM_PORT 固定时多实例冷启动的确定性竞争
    monkeypatch.setattr(
        shm_broadcast, "get_open_port", lambda: colliding_port, raising=False
    )
​
    n_writers = 32
    queues: list[MessageQueue] = []
    lock = threading.Lock()
​
    def make_writer():
        # n_local_reader=0 且 n_reader=1 触发远端 XPUB 绑定分支
        q = MessageQueue(
            n_reader=1,
            n_local_reader=0,
            max_chunk_bytes=4096,
            max_chunks=2,
            connect_ip="127.0.0.1",
        )
        with lock:
            queues.append(q)
​
    threads = [threading.Thread(target=make_writer) for _ in range(n_writers)]
    for t in threads:
        t.start()
    for t in threads:
        t.join()
​
    assert len(queues) == n_writers
    addrs = [q.export_handle().remote_subscribe_addr for q in queues]
    # 核心断言:32 个并发 writer 必须各自拿到唯一端口
    assert all(addr and addr.startswith("tcp://") for addr in addrs)
    assert len(set(addrs)) == n_writers
​
    # 往返验证:用导出的 handle 创建远端 reader,确保重建地址真实可连
    writer = queues[0]
    received = []
​
    def reader_main():
        reader = MessageQueue.create_from_handle(writer.export_handle(), rank=0)
        reader.wait_until_ready()
        received.append(reader.dequeue())
        reader.remote_socket.close(linger=0)
​
    reader_thread = threading.Thread(target=reader_main)
    reader_thread.start()
    writer.wait_until_ready()
    writer.enqueue("ping")
    reader_thread.join(timeout=30)
    assert not reader_thread.is_alive()
    assert received == ["ping"]
​
    for q in queues:
        q.remote_socket.close(linger=0)

评论区精华

LAST_ENDPOINT 端点归一化漂移与地址重建 设计

commit aa9719 说明 LAST_ENDPOINT 返回内核归一化后的端点,其地址形式可能与配置的 connect_ip 不同(如 IPv6 规范化),若整体作为公告地址会改变 reader 实际连接的 handle 字符串。

结论:只从 LAST_ENDPOINT 解析端口号(rsplit(":", 1)[1]),公告地址用原始 connect_ip 重建,保证 handle 契约不变。 · 已解决

回归测试的确定性失败设计 测试

commit 429fb1 说明:把 get_open_port 固定为同一个空闲端口,旧路径每个 writer 都会绑定同一端口而确定性失败;新实现完全不查询探测端口。作者验证 fails on base, passes with the fix。

结论:通过 monkeypatch pin 探测端口,使测试在 base 上必失败、带修复必通过,避免偶发竞态导致的 flaky 测试。 · 已解决

取代 #44495 与避免重复工作 other

PR body 说明本 PR 取代因合并冲突停滞的 #44495:其原实现作为首个 commit 原样保留(作者署名与 sign-off 完整),之上叠加地址重建与回归测试;并检查了所有引用 shm_broadcast 与 #28498 的 open PR 确认无重复。

结论:继承 #44495 原 commit,以两段式提交(原始修复 + 增强与测试)推进,避免重复劳动。 · 已解决

风险与影响

  • 行为变化:远端广播端口从「get_open_port 探测(受 VLLM_PORT 确定性扫描影响)」变为「内核随机分配」,VLLM_PORT 对远端广播端口不再生效;依赖固定端口白名单或防火墙策略的部署需要重新评估(对绝大多数场景这是期望行为)。
  • 依赖 ZMQ 语义:端口回读依赖 zmq.LAST_ENDPOINT 的字符串格式与 libzmq 行为,rsplit(":", 1)[1] 取最后一段端口;若未来端点格式演进(如 IPv6 scope id)存在解析风险,但 TCP 端点当前格式稳定。
  • 回归范围:改动位于所有 n_reader > n_local_reader 路径(跨节点 TP/PP/PCP、多实例)的远端 XPUB 绑定,属分布式通信核心路径;好在新增并发回归测试直接覆盖该分支,用真实 XPUB bind 验证。
  • 测试环境残留:PR body 报告本机 6 个既有用例因 harness 硬编码 MASTER_PORT=12345 被无关进程占用而失败,属既有环境冲突(恰是 #28498 静态端口家族的又一实例),非本 PR 引入。
  • 用户侧:多节点拓扑(executor 与 worker 分宿主)与多实例冷启动(RL rollout 集群、multi-instance serving)不再随机遭遇 EADDRINUSE 启动崩溃,issue #28498 的端口冲突问题在当前端口分配路径上被系统性消除。
  • 系统侧:广播队列创建路径不再存在无主端口窗口;压测从 main 的 19/24 失败提升到修复后 24/24 成功且公告地址全部唯一,200-writer 于 ::1 并发构造也全部唯一。
  • 团队侧:沉淀了「bind 0 + LAST_ENDPOINT」端口晚绑定模式(与 DP coordinator #37452 一致)和「pin 探测端口使旧路径确定性失败」的竞态回归测试方法论,对同家族其他 get_open_port 使用点有示范意义。
分布式通信核心路径变更 VLLM_PORT 对远端端口不再生效 依赖 ZMQ LAST_ENDPOINT 语义 既有测试环境端口冲突(非本 PR 引入)

关联 Issue

#28498 [Bug][RL]: Port Conflict

完整报告

参与讨论