Prhub

#32242 sglang rust server request message

原始 PR 作者 rainj-me 合并时间 2026-07-30 01:54 文件变更 7 提交数 4 评论 82 代码增减 +1568 / -38

执行摘要

新增 Rust 服务器消息层:请求体、调度器 wire 协议与改进的请求 ID

从 PR #29799 中拆分,作为 Rust 服务器替换 Python TokenizerManager 线路的一部分。需要为 /generate 请求路径在 Rust 端提供强类型、高性能的消息定义,保持与现有 Python 调度器的 wire 兼容性,并解决动态批处理中重复 rid 碰撞等一致性风险。

值得详细阅读,特别是 wire_struct! 宏设计、OneOrMany 的 sealed 模式、Rid 唯一化策略。讨论中的输入验证和错误处理共识可作为团队代码规范参考。建议在后续 egress PR 合并后,补充端到端兼容性测试。

讨论亮点

Review 中重点讨论了类型安全、输入验证和错误处理的一致性:

  • 类型安全:mrain 指出 OneOrManyuntagged 在与自描述类型(如 serde_json::Value)搭配时会导致 Many 分支永远无法匹配。作者通过添加密封 trait OneOrManyItem 限制允许的类型,从根本上消除了该风险。
  • 重复 rid 处理:merrymercy 指出注释声称 Python 同等行为但实际上 Python 会拒绝重复 rid,且 Rust 端缺失此检查。作者最终采用 Rid::from_client 通过唯一化(追加随机 base + 计数器)消除碰撞,无需运行时检查。
  • 错误类型统一:merrymercy 要求 split 方法的 Result<_, String> 应与其余代码的 Result<_, Error> 统一。作者改为 Err(Error::Validation(String)),保持全路径一致。
  • wire 格式维护:merrymercy 对 wire_struct! 的顺序依赖表达担忧,作者最终保留了宏但没有完全消除顺序约束,表示将在后续迭代中考虑更简洁的方案。

实现拆解

  1. 定义共享 wire 类型rust/sglang-server/src/message/types.rs):引入 TokenIds 类型别名、OneOrMany 枚举(支持标量或列表),并通过 sealed trait 限制允许的类型避免序列化二义性。同时定义 wire_struct! 宏以自动生成 msgspec array_like=True 的编解码。

  2. 定义请求结构rust/sglang-server/src/message/request.rs):实现 GenerateBody 结构体(对应 Python GenerateReqInput),包含 text/input_ids/sampling_params 等字段,并将 unknown fields allowed 以保证后向兼容;实现 into_requests 方法完成批处理展开与验证(空 input_ids 拒绝、重复 rid 拒绝等)。

  3. 定义调度器 wire structrust/sglang-server/src/message/io_struct.rs):通过 wire_struct!control_messages! 宏声明 TokenizedGenerateReqInputAbortReqGetInternalStateReq 等 msgspec 结构,保证字段顺序与 Python 端一致;实现 From<&GenerateRequest> 转换并补充单元测试验证 msgpack 形状。

  4. 重写请求标识符rust/sglang-server/src/ids.rs):将旧的 RidHash 替换为功能更完整的 Rid 类型,包含 from_client(唯一化客户端传入 ID 以防止碰撞)和 client_facing(截取唯一前缀还原客户端原始 ID)方法,以及 newnew_health_check 工厂方法;通过自定义 PartialEq/Hash 确保仅基于 ID 字符串比较。

  5. 采样参数类型rust/sglang-server/src/message/sampling.rs):初步定义 SamplingParams 结构体(后续补充 normalize/verify),以及用于 HTTP body 装箱的 SamplingParamsInputOneOrMany 形式);标志为 deny_unknown_fields 以尽早捕获参数错误。

  6. 测试与配置:在 io_struct.rs 内联了 abort_req_msgpack_shapeto_header_msgpack_is_positionally_aligned 等单元测试;未涉及性能测试脚本。

文件 模块 状态 重要度
rust/sglang-server/src/message/request.rs Rust 服务器 added 9.08
rust/sglang-server/src/message/types.rs Rust 服务器 added 8.89
rust/sglang-server/src/message/io_struct.rs Rust 服务器 added 8.77
rust/sglang-server/src/ids.rs Rust 服务器 modified 8.69
rust/sglang-server/src/message/sampling.rs Rust 服务器 added 8.04
rust/sglang-server/src/message.rs Rust 服务器 added 5.21
rust/sglang-server/src/lib.rs Rust 服务器 modified 3.78

关键符号

GenerateBody::into_requests Rid::from_client Rid::client_facing wire_struct! OneOrMany

关键源码片段

rust/sglang-server/src/message/request.rs core-logic

核心文件:定义 HTTP 请求体(`GenerateBody`)和到调度器的编码逻辑,替换 Python 端的分批处理。

/// The `/generate` wire body before batch splitting.
/// Unknown keys are IGNORED, matching Python's `GenerateReqInput`.
#[derive(Debug, Clone, Default, Deserialize)]
pub struct GenerateBody {
    /// Optional client-supplied request id(s): a single string (which fans out
    /// as `{rid}_{i}`) or one per item.
    #[serde(default)]
    pub rid: Option<OneOrMany<String>>,
    #[serde(default)]
    pub text: Option<OneOrMany<String>>,
    #[serde(default)]
    pub input_ids: Option<OneOrMany<TokenIds>>,
    #[serde(default)]
    pub stream: bool,
    /// One params object (broadcast) or a list of them (per item).
    #[serde(default)]
    pub sampling_params: Option<SamplingParamsInput>,
    // ... other fields omitted for brevity
}
rust/sglang-server/src/message/types.rs core-logic

提供共享 wire 类型(`OneOrMany`、`TokenIds`、`wire_struct!` 宏),确保序列化安全与 Python 端对齐。

/// A field taking a bare `T` **or** `[T, …]` (e.g. `text: "hi"` or `text: ["a","b"]`).
/// `untagged` takes the first variant that matches, so a `T` that itself accepts
/// a sequence would make `Many` unreachable — hence the [`OneOrManyItem`] gate.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum OneOrMany<T: OneOrManyItem> {
    One(T),
    Many(Vec<T>),
}/// Types vetted for [`OneOrMany`]. Sealed, so adding one is a deliberate act.
pub trait OneOrManyItem: sealed::SealedItem {}
impl<T: sealed::SealedItem> OneOrManyItem for T {}mod sealed {
    pub trait SealedItem {}
    impl SealedItem for bool {}
    impl SealedItem for i64 {}
    impl SealedItem for String {}
    impl SealedItem for super::TokenIds {}
}
rust/sglang-server/src/ids.rs core-logic

重写请求标识符:引入 `Rid` 结构体替代旧的 `RidHash`,支持客户端 ID 唯一化以防止重复 rid 碰撞。

/// Separates a client-supplied rid from the uniquifier appended to it.
/// `Rid::new()` is uuid hex and `Rid::new_health_check()` adds only
/// `HEALTH_CHECK_`, so its presence at the fixed offset below is what lets
/// `Rid::client_facing()` recognize a suffix without carrying a flag.
/// This matters because `Rid` rides on every `ChunkEvent`.
const UNIQ_SEP: u8 = b'#';
const UNIQ_DIGITS: usize = 16;
const UNIQ_SUFFIX_LEN: usize = 1 + UNIQ_DIGITS;#[derive(Clone, Debug)]
pub struct Rid {
    id: String,
    /// Partition key, derived from `id`. Never part of identity.
    hash: u64,
}// Identity is the ID, not the digest.
impl PartialEq for Rid {
    fn eq(&self, other: &Self) -> bool { self.id == other.id }
}
impl Eq for Rid {}
impl Hash for Rid {
    fn hash<H: std::hash::Hasher>(&self, state: &mut H) { self.id.hash(state); }
}impl Rid {
    pub fn new() -> Self {
        let id = Uuid::new_v4().simple().to_string();
        Rid::from(id)
    }    pub fn new_health_check() -> Self {
        let id = format!("{}_" , crate::HEALTH_CHECK_RID_PREFIX)
            + &Uuid::new_v4().simple().to_string();
        Rid::from(id)
    }    /// A CLIENT-SUPPLIED rid, made unique for internal use by appending a
    /// uniquifier. Prevents hash collisions (duplicate client rid) from
    /// causing two concurrent requests to share the same downstream sink.
    pub fn from_client(id: &str) -> Self {
        use std::sync::atomic::{AtomicU32, Ordering};
        static BASE: OnceLock<u32> = OnceLock::new();
        static NEXT: AtomicU32 = AtomicU32::new(0);
        let base = *BASE.get_or_init(|| Uuid::new_v4().as_u128() as u32);
        let n = NEXT.fetch_add(1, Ordering::Relaxed);
        Rid::from(format!("{id}{sep}{base:08x}{n:08x}" , sep = UNIQ_SEP as char))
    }    /// The rid as written by the client — what `meta_info.id` must echo.
    pub fn client_facing(&self) -> &str {
        let b = self.id.as_bytes();
        let Some(cut) = b.len().checked_sub(UNIQ_SUFFIX_LEN) else {
            return &self.id;
        };
        if b[cut] == UNIQ_SEP && b[cut + 1..].iter().all(u8::is_ascii_hexdigit) {
            &self.id[..cut]
        } else {
            &self.id
        }
    }    #[inline]
    pub fn shard(&self, n: usize) -> usize {
        debug_assert!(n > 0);
        (self.hash as usize) % n
    }
}impl From<String> for Rid {
    fn from(id: String) -> Self {
        let hash = {
            let state = std::collections::hash_map::RandomState::new();
            state.hash_one(&id)
        };
        Rid { id, hash }
    }
}

评论区精华

OneOrMany 的 untagged 安全性与 sealed trait 设计

mrain 指出若 `T` 为 `serde_json::Value` 等自描述类型,`untagged` 优先匹配 `One` 导致 `Many` 不可达。作者通过添加 `OneOrManyItem` 密封 trait 限制允许的类型,从根本上消除风险。

结论:添加密封 trait 限制类型,确保 `OneOrMany` 的 untagged 行为正确。 · 已解决

重复客户端 rid 处理 正确性

merrymercy 指出 Rust 端注释称 Python 同等忽略重复 rid,但实际 Python 会拒绝。建议添加检查。作者最终采用 `Rid::from_client` 唯一化(追加随机 base + 计数器),无需运行时检查即可消除碰撞。

结论:通过唯一化策略解决,无需逐请求检查重复。 · 已解决

错误类型统一:String vs Error style

merrymercy 要求 `split` 方法返回值类型从 `Result<_, String>` 改为 `Result<_, Error>` 以保持全路径一致。作者改为 `Err(Error::Validation(String))`。

结论:统一使用 `Error` 类型,增加可维护性。 · 已解决

风险与影响

  1. wire 兼容性:Rust 端 TokenizedGenerateReqInput 等结构必须与 Python 的 io_struct.py 严格同步,任何字段顺序或默认值变化都会导致解码失败。该 PR 通过单元测试(abort_req_msgpack_shape 等)覆盖基本形状,但 Python 端新增字段时需要同时更新 Rust 定义。
  2. 采样参数未完成SamplingParams::normalize()verify()todo!(),当前仅存储参数,过渡期需确保 Python 端仍主导规范化,但 is_normalized=true 标记可能会跳过 Python 端的检查,导致不一致。
  3. 缺失字段:未支持 prioritysession_idcustom_logit_processor 等字段,虽不影响现有客户端(未知字段忽略),但依赖这些功能的用户无法走 Rust 路径。
  4. 新 Rid 稳定Ridhash 字段使用 RandomState::new() 每次进程可能不同?不是,是每个进程创建一次 RandomState,所以同一进程内 hash 一致。但跨进程 hash 不同,影响 detok 分片?shard 方法用于 detok 分片,这只能在同一域内使用,应该没问题。

该 PR 对 Rust 服务器线路产生根本性影响,是后续 Rust 预处理、调度器交互的基础。对用户而言,无直接功能变化(Python 服务器仍为主入口);对系统而言,引入新的消息格式和 ID 策略,需确保与解码器兼容;对团队而言,后续必须跟进 egress 协议、detok 集成等。影响范围仅限于 Rust 服务器线路,不影响现有 Python 端。

wire 兼容性维护成本 部分采样参数未完成 缺少 egress 集成测试

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论