Prhub

#46833 [Rust Frontend] Start current wave for a stale DP FirstRequest

原始 PR 作者 blasrodri 合并时间 2026-06-30 13:13 文件变更 2 提交数 1 评论 3 代码增减 +84 / -9

执行摘要

修复 Rust DP 协调器中过时 FirstRequest 回退 wave 的问题

在进程内 DP 协调器中,notify_first_request 在引擎暂停时入队 FirstRequest{wave: W},但 runner 稍后才广播 START_DP_WAVE。如果引擎在该间隙报告 WaveComplete(W)current_wave 会前进到 W+1。此时过时的 FirstRequest{wave: W} 被无条件应用,将 current_wave 回退到 W 并广播一个已被取代的 wave。此行为与 Python 协调器的前端路径不一致:Python 版从不回退 wave,而是广播当前 wave 并在请求 wave 过时时唤醒所有引擎(exclude = None)。

值得精读,特别是关注 Rust 前端与 Python 协调器行为一致性设计的同学。start_wave_for_first_request 方法的设计清晰体现了竞态处理逻辑,单元测试覆盖了关键场景。

讨论亮点

本 PR 的 review 评论较少,主要由 BugenZhao 批准,表示 LGTM。无重大争议或未解决的疑虑。Codex 自动审查未发现重大问题。

实现拆解

  1. rust/src/engine-core-client/src/coordinator/inproc.rs:修改 broadcast_start_wave 方法,将 exclude_engine_index 参数类型从 u32 改为 Option<u32>,以支持 None(唤醒所有引擎)。修改 handle_commandFirstRequest 分支:不再直接设置 state.lock().current_wave = wave,而是调用 state.start_wave_for_first_request(wave, target_engine_index) 获取当前 wave 和 exclude 值,然后用当前 wave 广播。更新另一处 broadcast_start_wave 调用(引擎通知过时 wave 时)为 Some(engine_index)

  2. rust/src/engine-core-client/src/coordinator/handle.rs:在 CoordinatorStateSnapshot 上新增 start_wave_for_first_request 方法。该方法将 engines_running 置为 true,然后根据 request_wave >= self.current_wave 决定 exclude:如果请求 wave 不小于当前 wave(非过时),排除目标引擎;否则(过时)exclude 为 None,唤醒所有引擎。始终返回 self.current_wave,确保不回退 wave。

  3. rust/src/engine-core-client/src/coordinator/inproc.rs(测试):添加两个单元测试 first_request_for_current_wave_excludes_targetstale_first_request_starts_current_wave_for_all_engines,验证正常路径和过时路径的行为。

文件 模块 状态 重要度
rust/src/engine-core-client/src/coordinator/inproc.rs 协调器 modified 8.21
rust/src/engine-core-client/src/coordinator/handle.rs 协调器 modified 7.04

关键符号

broadcast_start_wave start_wave_for_first_request handle_command

关键源码片段

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

核心变更文件:修改 broadcast_start_wave 参数类型,重构 handle_command 中 FirstRequest 分支逻辑,并新增单元测试。

// rust/src/engine-core-client/src/coordinator/inproc.rs ( 关键变更部分 )/// 广播 START_DP_WAVE 消息给所有引擎。
/// exclude_engine_index 为 None 时唤醒所有引擎(用于过时请求 wave)。
async fn broadcast_start_wave(
    &mut self,
    wave: u32,
    exclude_engine_index: Option<u32>, // 从 u32 改为 Option<u32>
) -> Result<()> {
    let payload = encode_msgpack(&StartDpWaveMessage {
        wave,
        exclude_engine_index,
    })?;
    self.coordinator_input
        .send(
            ZmqMessage::try_from(vec![
                EngineCoreRequestType::StartDpWave.to_frame(),
                payload.into(),
            ])
            .expect("coordinator START_DP_WAVE message must contain two frames"),
        )
        .await?;
    Ok(())
}// handle_command 中 FirstRequest 分支,不再直接设置 current_wave,
// 而是委托给 start_wave_for_first_request 处理过时判断。
CoordinatorCommand::FirstRequest {
    target_engine_id,
    wave,
} => {
    let target_engine_index = target_engine_id.engine_index().ok_or_else(|| {
        Error::UnsupportedCoordinatorEngineId {
            engine_id: target_engine_id.to_vec(),
        }
    })?;
    // 获取当前 wave 和 exclude 值,注意这里不改变 current_wave
    let (current_wave, exclude) = {
        let mut state = self.state.lock();
        state.start_wave_for_first_request(wave, target_engine_index)
    };
    debug!(
        current_wave,
        request_wave = wave,
        ?exclude,
        "starting DP wave after first request while engines were paused"
    );
    // 广播当前 wave,exclude 为 None 时唤醒所有引擎
    self.broadcast_start_wave(current_wave, exclude).await?;
}
// 单元测试:过时请求 broadcast 全部引擎
#[test]
fn stale_first_request_starts_current_wave_for_all_engines() {
    let mut state = CoordinatorStateSnapshot {
        current_wave: 4,
        engines_running: false,
    };
    // 请求 wave 3 小于 current_wave 4,说明 wave 已过时
    let (wave, exclude) = state.start_wave_for_first_request(3, 2);
    assert_eq!(wave, 4); // 广播当前 wave,不回退
    assert_eq!(exclude, None); // 唤醒所有引擎
    assert!(state.engines_running);
    assert_eq!(state.current_wave, 4); // current_wave 保持不变
}
rust/src/engine-core-client/src/coordinator/handle.rs core-logic

新增 start_wave_for_first_request 方法,封装过时 wave 判断逻辑;修改 exclude 字段类型为 Option<u32>。

// rust/src/engine-core-client/src/coordinator/handle.rs ( 新增方法 )impl CoordinatorStateSnapshot {
    /// 为 FirstRequest 启动 wave,返回要广播的 wave 和要排除的引擎索引。
    /// 如果 request_wave 小于 current_wave(过时),则 exclude 为 None(唤醒所有引擎),
    /// 且绝不回退 current_wave。
    pub(crate) fn start_wave_for_first_request(
        &mut self,
        request_wave: u32,
        target_engine_index: u32,
    ) -> (u32, Option<u32>) {
        self.engines_running = true;
        // 当 request_wave >= current_wave 时,说明请求 wave 是当前的,
        // 排除目标引擎(它已经接收过请求);否则(过时)None 唤醒所有引擎。
        let exclude = (request_wave >= self.current_wave).then_some(target_engine_index);
        (self.current_wave, exclude) // 始终返回 current_wave,不回退
    }
}
// StartDpWaveMessage 结构体中的 exclude 字段类型改变
struct StartDpWaveMessage {
    wave: u32,
    /// Engine index that already received the triggering request and so does not
    /// need an extra wakeup. `None` wakes every engine (used when the triggering
    /// request was for a stale wave).
    exclude_engine_index: Option<u32>, // 从 u32 改为 Option<u32>
}

评论区精华

没有提炼出高价值讨论线程

当前评论区没有形成足够清晰的争议点或结论,后续有更多讨论时会体现在这里。

风险与影响

风险较低。变更集中在 Rust 引擎核心客户端协调器中,修改了 wave 广播逻辑,确保与 Python 实现一致。涉及竞态条件修复,测试覆盖了正常和过时路径。主要风险是:如果其他部分依赖旧行为(wave 回退),可能出现意外行为,但旧行为本身就是 bug,影响面小。

直接影响:修复了 Rust DP 协调器中 wave 管理的一个竞态条件 bug,确保在 FirstRequest 过时时正确广播当前 wave 并唤醒所有引擎。影响范围仅限于 Rust 实现的进程内 DP 协调器模块(rust/src/engine-core-client),不影响 Python 协调器或其他系统。改善了 DP 数据并行场景下的稳定性。

竞态条件修复 测试覆盖关键路径

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论