Prhub

#52671 [Rust Frontend] Wait for all utility calls to finish

原始 PR 作者 connorcarpenter15 合并时间 2026-08-19 02:23 文件变更 4 提交数 3 评论 8 代码增减 +187 / -85

执行摘要

多引擎 utility 调用等待全部收敛后再报错,支持安全补偿

PR body 明确说明:"Make multi-engine utility fanout wait for every engine outcome before returning an error",并指出要让调用方"compensate partially applied mutations only after every engine call has finished"。旧实现使用 try_join_all 做 fail-fast,某个 engine 一旦失败就立刻返回错误,而其余 engine 可能仍在执行同一个 mutation(如 add_lora、pause/resume),调用方无法判断哪些引擎已应用变更,补偿操作不安全。

值得精读。三个看点:(1) join_alltry_join_all 在"补偿安全"场景下的语义差异,是并发扇出错误模型的重要一课;(2) 用局部 RAII 守卫 UtilityCallGuard 统一异步取消、准备失败、发送失败的清理路径,是 Rust 异步代码清理资源的通用范式;(3) BugenZhao 的简化 commit 示范了如何把多分支错误处理收敛成单一 Drop 路径,并顺带修复普通 call 的取消泄漏。若后续要扩展 utility 协议或新增多引擎操作,建议直接复用这一模式。

讨论亮点

核心 review 交锋来自维护者 BugenZhao。他在 approval 中肯定这是一个 good catch,并追加两个 commit:一是简化 call_utility 的实现,把准备失败、发送失败、取消三条清理路径收敛成单一 RAII Drop 路径;二是把相同的取消安全语义对齐到普通生成请求 call(先构造输出流再发送,保证 cancel 时经 stream 的 Drop 清理注册,避免悬挂 waiter)。原 PR body 的核心矛盾点是:fail-fast 会让调用方在其余引擎尚未完成 mutation 时就收到错误,无法安全补偿;BugenZhao 的简化没有改变这一语义,只是让实现更紧凑。其余为流程噪音:mergify 提示 pre-commit 失败与合并冲突,作者 rebase 后重新触发 Buildkite CI #84410 并通过。

实现拆解

本 PR 的改造集中在 Rust engine-core client 的 utility 扇出路径,分四步演进:

  1. 定位问题入口:EngineCoreClient::call_utilityrust/src/engine-core-client/src/client.rs)。旧实现分三阶段:先为每个 engine 分配 call_id 并注册 waiter;再用 try_join_all 并发发送请求(phase 2 fail-fast);最后用 try_join_all 等待响应(phase 3 fail-fast)。任一阶段出错立即返回,导致其余 engine 的 mutation 状态未知,且已注册的 waiter 需要手动回滚,多地散落 unregister_utility_calls 调用,容易漏清理。

  2. 核心语义改造:把"发送 + 等待响应 + 类型化解码"合并到每个 engine 的单个 future 里,统一用 join_all 驱动。join_all 不会在第一个错误出现时取消其余 future,而是等所有结果都收敛,最后通过 outcomes.into_iter().collect()Vec<Result<T>> 折叠为 Result<Vec<T>>,保持 per-engine 顺序的结果列表。这样任何 engine 失败时,其余引擎的 mutation 也都已落定,调用方才能安全补偿。

  3. 清理路径收敛为 RAII:新增局部结构体 UtilityCallGuard,其 Drop 实现统一 draincall_ids 并调用 ClientInner::unregister_utility_calls。无论是请求准备失败、发送失败,还是整个 call_utility future 被 cancel(drop),都走同一条清理路径,不再需要手写回滚分支。

  4. 取消安全对齐与测试配套:BugenZhao 在 review 中追加 commit,把 call 改为在发送前先构造 EngineCoreOutputStream,使 future 被取消时仍会经 stream 的 Drop 清理请求注册;同时新增 #[cfg(test)]ClientInner::pending_utility_call_countUtilityRegistry::len 测试钩子。测试方面:tests/client.rs 将原 call_utility_failure_message_surfaces_as_error 重写为 call_utility_waits_for_all_engines_before_returning_error(双 mock engine:engine0 先回 failure、engine1 挂起,断言错误不会提前返回且 pending 计数收敛为 1,释放后才透出错误),并新增 cancelled_utility_call_unregisters_waiter 验证取消后注册表清零;原有 dispatcher_failure_propagates_to_waiting_utility_calls 保留。

配套改动仅 Rust crate 内部,wire format、Python 侧的 utility 协议与成功路径行为均未变化。

文件 模块 状态 重要度
rust/src/engine-core-client/src/client.rs 引擎客户端 modified 7.63
rust/src/engine-core-client/src/tests/client.rs 单元测试 modified 7.03
rust/src/engine-core-client/src/client/imp.rs 引擎客户端 modified 5.17
rust/src/engine-core-client/src/client/state.rs 状态注册表 modified 5.17

关键符号

call_utility call UtilityCallGuard::drop pending_utility_call_count UtilityRegistry::len call_utility_waits_for_all_engines_before_returning_error cancelled_utility_call_unregisters_waiter

关键源码片段

rust/src/engine-core-client/src/client.rs core-logic

核心改造文件:call_utility 由 fail-fast 改为 join_all 全量收敛,引入 UtilityCallGuard RAII 清理;call 改为先构造输出流再做发送以确保取消安全。

/// Call a typed utility method on all connected engines, returning one
/// decoded result per connected engine if all calls succeed.
///
/// 关键语义:与旧的 fail-fast(`try_join_all`)不同,这里会等待所有
/// engine 都返回结果(成功或失败)后才对外暴露错误,这样调用方只有在
/// 每个引擎的变更都落定之后,才能安全地补偿部分已应用的 mutation。
pub async fn call_utility<T, A>(&self, method: &str, args: A) -> Result<Vec<T>>
where
    T: serde::de::DeserializeOwned,
    A: serde::Serialize + std::fmt::Debug,
{
    // RAII 守卫:无论函数提前返回、发送中途失败,还是整个 future 被
    // cancel(drop),`Drop` 都会统一把已分配的 call_id 从注册表中删除,
    // 避免 waiter 泄漏到 shutdown。旧代码需要手写多个回滚分支,容易漏。
    struct UtilityCallGuard<'a> {
        inner: &'a ClientInner,
        call_ids: Vec<u64>,
    }    impl Drop for UtilityCallGuard<'_> {
        fn drop(&mut self) {
            self.inner.unregister_utility_calls(self.call_ids.drain(..));
        }
    }    let mut call_guard = UtilityCallGuard {
        inner: self.inner.as_ref(),
        call_ids: Vec::with_capacity(self.engines.len()),
    };    // 为每个 connected engine 分配唯一的 call_id 并注册响应 waiter,
    // 同时按 Python 侧 `(client_index, call_id, method_name, args)`
    // 的元组契约构建请求 payload。
    let mut prepared_calls = Vec::with_capacity(self.engines.len());
    for engine in &self.engines {
        let (call_id, rx) = self.inner.allocate_and_register_utility_call()?;
        call_guard.call_ids.push(call_id);
        let request = EngineCoreUtilityRequest::new(
            self.config.client_index,
            call_id,
            method,
            &args,
        )?;
        prepared_calls.push((&engine.engine_id, call_id, rx, request));
    }    // 并发发送 + 等待响应。`join_all` 不会像 `try_join_all` 那样在第一个
    // 错误出现时取消其余 future,而是等所有结果收敛后再统一处理。
    let outcomes = join_all(prepared_calls.into_iter().map(
        |(engine_id, call_id, rx, request)| async move {
            self.inner
                .send_to_engine(engine_id, EngineCoreRequestType::Utility, &request)
                .await?;
            rx.await
                .map_err(|_| Error::UtilityCallClosed {
                    method: method.to_string(),
                    call_id,
                })??
                .into_typed_result(method)
        },
    ))
    .await;    // `Vec<Result<T>>` 折叠为 `Result<Vec<T>>`,结果列表保持 engine 顺序;
    // 此时所有 engine 均已返回,调用方可以安全地执行补偿逻辑。
    outcomes.into_iter().collect()
}
// `call` 侧与 `call_utility` 对齐取消安全:先构造输出流,再真正向 engine
// 发送。这样一旦在发送过程中 future 被 cancel,stream 被 drop 时仍会走
// `EngineCoreOutputStream` 的 `Drop` 逻辑把请求注册回滚掉,不会留下
// 悬挂的注册项。
pub async fn call(&self, mut req: EngineCoreRequest) -> Result<EngineCoreOutputStream> {
    req.client_index = self.config.client_index;
    req.validate()?;    let (engine_id, rx) =
        self.inner.register_request(req.request_id.clone(), lora_name, data_parallel_rank)?;    let stream = EngineCoreOutputStream::new(
        req.request_id.clone(),
        engine_id.engine_index().unwrap_or(0),
        self.abort_tx.clone(),
        rx,
    );    let result: Result<()> = async {
        // 协调器存在时,用协调器快照覆盖当前 wave,并在引擎未运行时
        // 通知第一个请求,触发引擎启动流程。
        if let Some(coordinator) = self.coordinator.as_ref() {
            let snapshot = coordinator.snapshot();
            req.current_wave = snapshot.current_wave;
            if !snapshot.engines_running {
                coordinator.notify_first_request(engine_id.clone())?;
            }
        }
        self.inner.send_request_to_engine(&engine_id, req).await?;
        Ok(())
    }
    .await;    // 发送失败则回滚注册,并丢弃已构造的 stream(其 `Drop` 兜底清理)。
    if let Err(error) = result {
        self.inner.rollback_request(&req.request_id);
        return Err(error);
    }    Ok(stream)
}
rust/src/engine-core-client/src/tests/client.rs test-coverage

回归测试核心文件:把原单引擎失败测试重写为双引擎等待测试,新增取消清理测试,验证本 PR 的两个关键行为。

#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cancelled_utility_call_unregisters_waiter() {
    init_tracing();
    let ipc = IpcNamespace::new().unwrap();
    let handshake_address = ipc.handshake_endpoint();
    let (received_tx, received_rx) = oneshot::channel();    // mock engine:收到 utility 请求后只上报 " 已收到 ",然后永远挂起,
    // 模拟一个迟迟不给响应的远端。
    let (_shutdown, engine_task) = spawn_mock_engine_task(
        handshake_address.clone(),
        vec![0x00, 0x00],
        |dealer, _push| {
            Box::pin(async move {
                let _utility = recv_engine_message(dealer).await;
                let _ = received_tx.send(());
                std::future::pending::<()>().await;
            })
        },
    );    let client = std::sync::Arc::new(connect_client_with_ipc(/* ... */).await);    // 发起 utility 调用后,注册表中应当立即出现 1 个 pending waiter。
    let call_client = client.clone();
    let call =
        tokio::spawn(async move { call_client.call_utility::<bool, _>("add_lora", ()).await });
    received_rx.await.unwrap();
    assert_eq!(client.pending_utility_call_count(), 1);    // 取消调用:注册表中不再有残留 waiter;同时 Arc 还能被 `try_unwrap`,
    // 说明没有任何任务在结束后仍持有 client。
    call.abort();
    assert!(call.await.unwrap_err().is_cancelled());
    assert_eq!(client.pending_utility_call_count(), 0);    engine_task.abort();
    let client = std::sync::Arc::try_unwrap(client)
        .unwrap_or_else(|_| panic!("utility task retained client after cancellation"));
    client.shutdown().await.unwrap();
}

评论区精华

实现简化与取消安全对齐 设计

BugenZhao 在 approval 中肯定这是 good catch,并追加两个 commit:一是简化 call_utility 实现,把准备失败、发送失败、取消三条清理路径收敛为单一 RAII Drop 路径;二是把相同的取消安全语义对齐到普通生成请求 call——先构造 EngineCoreOutputStream 再发送,保证 cancel 时经 stream 的 Drop 清理请求注册。

结论:合并时已包含这两个 commit;call_utility 与 call 的注册清理均改为依赖 Drop 兜底,原 PR 的手动回滚分支被移除。 · 已解决

CI 与合并流程 other

mergify bot 提示 pre-commit 检查失败,随后 PR 出现合并冲突;作者 rebase 后重新触发 Buildkite CI。

结论:冲突与 pre-commit 问题在最终 head 中解决,Buildkite CI #84410 通过后由 njhill 合并。 · 已解决

风险与影响

  1. 错误返回延迟显著拉长:新实现即使某个 engine 已失败,也要等全部引擎返回(成功或失败)才报错;若某个 engine 响应异常挂起且链路没有超时兜底,call_utility 会被无限阻塞,而旧 fail-fast 行为会立即暴露问题。这是"补偿安全 vs 时效性"的刻意取舍,建议上层确认 utility 调用是否有超时保护。
  2. 级联行为变化:call_utility_consensuscollective_rpc 都委托 call_utility,它们的失败语义也随之从 fail-fast 变为全量收敛,影响所有 utility 调用方。
  3. 取消清理依赖 Drop 兜底:call 的取消安全完全依赖 EngineCoreOutputStream::Drop 的回滚实现(见 client.rs 中新增注释),未来若改动 stream 的 Drop 语义,注册清理会被静默破坏。
  4. 测试钩子暴露内部状态:pending_utility_call_count()UtilityRegistry::len() 虽为 #[cfg(test)],仍为注册表内部结构新增了 API 面,测试专用路径需防止误入生产逻辑。

影响范围:仅 Rust 前端 vllm-engine-core-client crate,Python 侧 utility 协议与 wire format 未变,推理输出与精度不受影响(PR body 明确说明 model evaluation 不适用)。受益方:多 engine(data-parallel)v1 部署下,LoRA 增删、pause/resume、is_sleeping 等 utility 调用获得确定的"全引擎收敛"错误语义,调用方可以放心补偿部分已应用的变更;同时显著减少 waiter 泄漏与取消后遗留的悬挂注册。代价:所有 utility 调用方的错误返回延迟变长。对团队而言,本 PR 确立了一个可复用的异步取消安全模式(RAII guard + join_all 收敛),后续新增 utility 逻辑应沿用。

无超时保护的全量等待语义 并发路径语义变更(fail-fast → 全量收敛) 取消安全依赖 Drop 兜底 影响所有 utility 调用方

关联 Issue

未识别关联 Issue

当前没有检测到明确关联的 Issue 链接,后续同步到相关引用后会出现在这里。

完整报告

参与讨论