Prhub

#50033 [Rust Frontend][gRPC] Add KV event source discovery

原始 PR 作者 connorcarpenter15 合并时间 2026-07-30 04:46 文件变更 13 提交数 2 评论 0 代码增减 +179 / -5

执行摘要

添加 gRPC 发现 KV 事件源的能力

Rust 前端需要动态发现已配置的 ZMQ KV 事件发布者,避免重复解析端点和冗余配置。PR 描述明确要求 'Add Control.GetKvEventSources for discovering configured ZMQ KV-event publishers',并强调 'Expose the publisher's resolved runtime configuration through the EngineCore ready handshake so the Rust frontend does not duplicate rank-based endpoint resolution'。

值得精读,尤其是跨语言配置传递的设计模式和 gRPC 接口的简洁实现。关键关注 kv_event_source 过滤逻辑、EngineCoreReadyResponse 的握手扩展方式,以及 Python 侧如何暴露内部配置给 Rust 前端。

讨论亮点

本 PR 无实质性的 Review 讨论。claude[bot] 自动评论说明从 Fork 提交时自动审查被禁用,njhillmgoin 均直接批准,未留下进一步意见。

实现拆解

  1. Python 侧暴露发布者配置:在 vllm/distributed/kv_events.pyEventPublisher 基类新增 get_publisher_config 方法(默认返回 None),并在 ZmqEventPublisher 中实现以返回完整的 KVEventsConfig 对象,包含端点、主题、队列设置等信息。
  2. 调度器传递配置:在 vllm/v1/core/sched/interface.pyscheduler.py 添加 get_kv_event_publisher_config 方法,通过调用 publisher 实例的方法获取配置,使上层能获取该配置。
  3. 握手协议扩展:在 rust/src/engine-core-client/src/protocol/handshake.rs 新增 KvEventsConfig 结构体,并作为可选字段加入 EngineCoreReadyResponse,使 Python 端配置能序列化传递给 Rust 前端。
  4. Rust gRPC 实现:在 rust/src/server/src/grpc/control.rs 实现 get_kv_event_sources RPC 处理函数,遍历所有已连接的 EngineCore 响应,通过 kv_event_source 函数过滤出启用了 ZMQ 的事件源并转换为 protobuf 消息。
  5. 测试与兼容性:Rust 侧新增单元测试验证过滤逻辑(kv_event_source_filters_and_exposes_zmq_publisher);更新 Python/Rust MessagePack 兼容性测试以保证 KVEventsConfig 序列化一致;调整分布式事件测试以覆盖新配置路径。
文件 模块 状态 重要度
rust/src/server/src/grpc/control.rs gRPC 服务 modified 6.97
rust/src/server/src/grpc/tests.rs gRPC 测试 modified 7.07
vllm/distributed/kv_events.py KV 事件层 modified 6.49
vllm/v1/core/sched/scheduler.py 调度器 modified 6.21
rust/src/engine-core-client/src/protocol/handshake.rs 握手协议 modified 6.15
vllm/v1/core/sched/interface.py 调度接口 modified 6.11

关键符号

get_kv_event_sources kv_event_source get_publisher_config get_kv_event_publisher_config KvEventsConfig

关键源码片段

rust/src/server/src/grpc/control.rs core-logic

核心 gRPC 实现:添加 `get_kv_event_sources` RPC 和 `kv_event_source` 过滤函数,是 PR 的主要功能载体。

// rust/src/server/src/grpc/control.rs - 新增 RPC 实现与过滤逻辑async fn get_kv_event_sources(
    &self,
    _request: Request<pb::GetKvEventSourcesRequest>,
) -> Result<Response<pb::GetKvEventSourcesResponse>, Status> {
    let client = self.state.engine_core_client();
    // 遍历所有 EngineCore 的 ready 响应,过滤出 ZMQ 事件源
    let sources = client
        .ready_responses()
        .into_iter()
        .filter_map(kv_event_source)
        .collect();
    Ok(Response::new(pb::GetKvEventSourcesResponse { sources }))
}/// 从 EngineCoreReadyResponse 中提取 ZMQ 事件源配置
pub(super) fn kv_event_source(
    response: &EngineCoreReadyResponse,
) -> Option<pb::KvEventSource> {
    let config = response.kv_events_config.as_ref()?;
    // 仅当启用 KV 缓存事件且发布者类型为 "zmq" 时才暴露
    if !config.enable_kv_cache_events || config.publisher != "zmq" {
        return None;
    }
    Some(pb::KvEventSource {
        transport: "zmq".to_string(),
        endpoint: config.endpoint.clone(),
        topic: config.topic.clone(),
        replay_endpoint: config.replay_endpoint.clone().unwrap_or_default(),
        data_parallel_rank: Some(response.data_parallel_rank),
        encoding: "msgpack".to_string(),
        schema_version: 1,
        buffer_steps: config.buffer_steps,
        hwm: config.hwm,
        max_queue_size: config.max_queue_size,
    })
}
rust/src/server/src/grpc/tests.rs core-logic

Rust 单元测试验证事件源过滤逻辑的正确性,确保只有启用 ZMQ 且配置有效的发布者才被暴露。

// rust/src/server/src/grpc/tests.rs - 事件源过滤单元测试#[test]
fn kv_event_source_filters_and_exposes_zmq_publisher() {
    let mut ready = default_ready_response();
    ready.data_parallel_rank = 2;
    ready.kv_events_config = Some(KvEventsConfig {
        enable_kv_cache_events: false,
        publisher: "null".to_string(),
        endpoint: "tcp://*:5559".to_string(),
        replay_endpoint: Some("tcp://*:5560".to_string()),
        buffer_steps: 10_000,
        hwm: 100_000,
        max_queue_size: 100_000,
        topic: "kv".to_string(),
    });    // 未启用时返回 None
    assert!(kv_event_source(&ready).is_none());    // 启用并改为 ZMQ 后应返回正确配置
    let config = ready.kv_events_config.as_mut().unwrap();
    config.enable_kv_cache_events = true;
    config.publisher = "zmq".to_string();
    let source = kv_event_source(&ready).expect("configured ZMQ event source");
    assert_eq!(source.transport, "zmq");
    assert_eq!(source.endpoint, "tcp://*:5559");
    assert_eq!(source.topic, "kv");
    assert_eq!(source.replay_endpoint, "tcp://*:5560");
    assert_eq!(source.data_parallel_rank, Some(2));
    assert_eq!(source.encoding, "msgpack");
    assert_eq!(source.schema_version, 1);
    assert_eq!(source.buffer_steps, 10_000);
    assert_eq!(source.hwm, 100_000);
    assert_eq!(source.max_queue_size, 100_000);
}
vllm/distributed/kv_events.py core-logic

Python 侧关键变更:为 `EventPublisher` 和 `ZmqEventPublisher` 新增 `get_publisher_config` 方法,将运行配置暴露给上层。

# vllm/distributed/kv_events.py - 发布者配置暴露class EventPublisher(ABC):
    # ...
    def get_publisher_config(self) -> KVEventsConfig | None:
        """Return the publisher's resolved runtime configuration."""
        return Noneclass ZmqEventPublisher(EventPublisher):
    def __init__(self, data_parallel_rank: int, endpoint: str = ..., ...):
        # ... 初始化后构建配置对象
        self._publisher_config = KVEventsConfig(
            enable_kv_cache_events=True,
            publisher="zmq",
            endpoint=self._endpoint,
            replay_endpoint=self._replay_endpoint,
            buffer_steps=buffer_steps,
            hwm=hwm,
            max_queue_size=max_queue_size,
            topic=topic,
        )
        # ...
​
    def get_publisher_config(self) -> KVEventsConfig:
        return self._publisher_config

评论区精华

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

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

风险与影响

整体风险较低。主要风险点:

1) 跨语言序列化兼容性KVEventsConfig 通过握手消息从 Python 传递到 Rust,若字段类型或默认值不一致可能导致解析失败。代码中 Rust 侧使用了 #[serde(default)] 并保持了与 Python 侧命名的对齐,风险可控。
2) 配置过滤逻辑依赖硬编码字符串kv_event_source 函数中检查 publisher != "zmq",若未来添加新发布者类型需同步更新。
3) 无生产环境完整测试:PR 作者提到未在完整 Python 环境本地运行分布式测试,但 CI 会覆盖,回归风险较小。

对用户无直接影响(仅新增内部 gRPC API)。对分离式服务架构,前端可动态获取 KV 事件源信息,有助于未来实现事件驱动的调度或监控。对团队,提供了一种清晰的跨语言配置共享模式,便于后续扩展其他发布者类型。

跨语言序列化 配置过滤硬编码

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论