Prhub

#32342 sglang rust server egress message

原始 PR 作者 rainj-me 合并时间 2026-07-30 08:16 文件变更 5 提交数 2 评论 6 代码增减 +2445 / -27

执行摘要

新增 Rust egress 消息层与完成原因模块

将 Rust 服务器消息层拆分为独立 PR,以增量方式替换 Python TokenizerManager 的 egress 处理路径。PR body 明确说明从 #29799 拆分,并依赖 #32240。目标是让 Rust 端接管响应序列化、完成原因分类和采样参数验证,减少 Python 调度器的负载。

该 PR 设计精巧,尤其关注了向前兼容(FinishReason::Unknown)和安全性(正则 admission 缓存、边界检查)。值得所有参与 Rust 服务器开发的成员精读,理解 egress 协议和数据类型约定。对于非 Rust 开发者,可以跳过实现细节,但应当通过 final_report 了解整体架构演进。

讨论亮点
  1. 数据重叠风险egress.rs:163):reviewer mrain 指出 take_f32/take_i32 等函数未检查 cvci 区间是否重叠。作者 rainj-me 回应:这是 hot path,调用方必须保证游标不重叠,否则被视为 malformed 帧。已接受该风险假设。

  2. SinkError 缺少 Displayegress.rs:25):mrain 注意到 SinkError 未实现 Display,无法直接用于日志。作者确认当前不记录该错误,未来需要时再添加。

  3. 直接解析为 Vecegress.rs 无行号):mrain 建议直接从字节流解析为 Vec<ChunkEvent> 以简化后续操作。作者未直接回应,该建议未在 PR 中采纳,可能留待后续优化。

实现拆解

  1. 新增 egress 模块 (rust/sglang-server/src/message/egress.rs):定义了 EgressSink(基于 mpsc 的 per-request 回传通道)、EgressItem(帧/完成/控制/错误四种事件)、帧标签常量以及列式批解码辅助函数(take_f32take_i32frame_egress_batch_cols 等)。这些类型构成了 Rust 侧响应流的基础。

  2. 新增完成原因模块 (rust/sglang-server/src/message/finish_reason.rs):实现了 FinishReason 枚举,其 Known(FinishKind) 分支覆盖 StopLengthAbort 三种完成类型,Unknown 分支保留原始 JSON Map 以确保向前兼容。配套 Matched 枚举和 AbortReason 结构体精确对应 Python 的 wire 格式。提供了 matched()abort_status() 方法供其他组件使用。

  3. 增强采样参数验证 (rust/sglang-server/src/message/sampling.rs):为 SamplingParams 添加了 normalize()verify() 方法,完整复制 Python 的 __post_init__normalize → verify 管线。引入了 defaulted! 宏简化字段默认值管理,并增加了对 stop_strsstop_regex 的长度与数量硬限制。

  4. 优化正则表达式 admission (rust/sglang-server/src/utils/regex.rs):引入基于 HashMap 的 pattern → bound 缓存(容量 512),避免重复解析已验证的正则表达式。cached_boundcache_bound 函数在 RegexPattern::build 中被调用,显著降低批量请求时的 CPU 开销。

  5. 注册新模块 (rust/sglang-server/src/message.rs):通过 mod egress; mod finish_reason; 将两个新模块纳入编译。

文件 模块 状态 重要度
rust/sglang-server/src/message/egress.rs 出方向 added 9.08
rust/sglang-server/src/message/finish_reason.rs 完成原因 added 8.94
rust/sglang-server/src/message/sampling.rs 采样参数 modified 8.65
rust/sglang-server/src/utils/regex.rs 正则缓存 modified 7.96
rust/sglang-server/src/message.rs 模块注册 modified 4.52

关键符号

try_send take_f32 take_i32 frame_egress_batch_cols take_flat take_poslens take_ragged take_hidden from matched abort_status fr abort_status_extracts_code_and_message non_error_finishes_stay_ok matched_reads_stops_only finish_reason_round_trips_python_shapes max_new_tokens_default normalize verify max_tokens_len deserialize cached_bound cache_bound admission_memo_agrees_with_admitting_afresh

关键源码片段

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

新增 egress 模块核心文件,定义 EgressSink、EgressItem、帧标签和列式批解码辅助函数,是响应管线的基石。

//! Egress 方向的核心数据类型 —— 响应回传通道与帧编码。use bytes::Bytes;
use tokio::sync::mpsc;
use crate::error::Error;/// 每个请求的回传通道,API 处理器从中 drain 响应。
/// 当前仅支持 Local(mpsc),未来可扩展为远程传输。
#[derive(Clone, Debug)]
pub enum EgressSink {
    Local(mpsc::Sender<EgressItem>),
}/// try_send 失败的原因。
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SinkError {
    Full, // 客户端反压
    Closed, // 客户端已断开
}impl EgressSink {
    /// 非阻塞发送。Full 和 Closed 均为流终止状态,调用方据此区分日志等级。
    pub fn try_send(&self, item: EgressItem) -> Result<(), SinkError> {
        match self {
            EgressSink::Local(tx) => tx.try_send(item).map_err(|e| match e {
                mpsc::error::TrySendError::Full(_) => SinkError::Full,
                mpsc::error::TrySendError::Closed(_) => SinkError::Closed,
            }),
        }
    }
}/// API 处理器在 egress 流上收到的四种事件。
#[derive(Debug)]
pub enum EgressItem {
    Frame(ChunkEvent), // 流式中间帧
    Done(ChunkEvent), // 最终完成帧
    Control(Bytes), // 控制请求的透传结果
    Error(Error), // 终端错误
}// 帧标签常量 —— 用于 egress ring 协议。
pub const EGRESS_TAG_RESULT: u8 = 1; // 单一控制结果
pub const EGRESS_TAG_BATCH: u8 = 2; // 解码 batch
pub const EGRESS_TAG_ERROR: u8 = 3; // 单请求错误
rust/sglang-server/src/message/finish_reason.rs core-logic

新增完成原因模块,精确对应 Python FinishReasonDict,提供向前兼容的 Unknown 分支和 abort_status 辅助方法。

//! FinishReason 枚举 —— 精确映射 Python `FinishReasonDict`。
//! 未知类型被保留为 raw map,向前兼容。use serde::{Serialize, Deserialize};/// stop 匹配的值:单个 token id、字符串或多 token 序列。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum Matched {
    Token(i64),
    Str(String),
    Tokens(Vec<i64>),
}/// 已知的完成类型,通过 JSON 字段 `type` 区分。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "lowercase")]
pub enum FinishKind {
    Stop {
        #[serde(default)]
        matched: Option<Matched>,
    },
    Length {
        #[serde(default)]
        length: Option<u64>,
    },
    // Abort 被 Box 包装以减少 ChunkEvent 的大小(它很少出现)。
    Abort(Box<AbortReason>),
}/// Abort 附带的详细信息。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AbortReason {
    pub message: Option<String>,
    pub status_code: Option<u16>,
    pub err_type: Option<String>,
}/// 完成原因对外表示。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum FinishReason {
    Known(FinishKind),
    // 未知类型保留原始 JSON Map,保证 batch 解析不会因新增完成类型而整体失败。
    Unknown(serde_json::Map<String, serde_json::Value>),
}impl FinishReason {
    /// 如果原因是 stop,返回匹配到的 stop 值。
    pub fn matched(&self) -> Option<&Matched> {
        match self {
            FinishReason::Known(FinishKind::Stop { matched }) => matched.as_ref(),
            _ => None,
        }
    }    /// 如果原因是有状态码的 abort,返回 (code, message)。
    pub fn abort_status(&self) -> Option<(u16, &str)> {
        match self {
            FinishReason::Known(FinishKind::Abort(a)) => {
                Some((a.status_code?, a.message.as_deref().unwrap_or("request aborted")))
            }
            _ => None,
        }
    }
}

评论区精华

take_f32/take_i32 区间重叠风险 正确性

reviewer mrain 指出 `[cv..cv + 4*l]` 与 `[ci..ci + 4*l]` 可能重叠,导致未定义读取。

结论:作者承认问题,但强调这是 hot path,调用方必须保证游标不重叠,否则视为 malformed 帧拒绝。 · 已解决

SinkError 缺少 Display 实现 style

reviewer mrain 询问 SinkError 为何未实现 Display,是否影响日志记录。

结论:作者回应当前不记录该错误,未来需要时再添加。 · 已解决

直接解析为 Vec<ChunkEvent> 的建议 设计

reviewer mrain 建议直接解析 egress 消息为 `Vec<ChunkEvent>`,避免后续转换开销。

结论:作者未直接回应,建议未被采纳;该设计决策留待后续优化。 · pending

风险与影响

  1. 缺少测试覆盖:本次变更未包含任何测试文件,egress 帧解析和完成原因序列化的正确性依赖 Python 侧的集成测试。关键路径如批解码边界条件、填充字节对齐等未被 Rust 单元测试覆盖。
  2. cursoring 重叠假设take_f32/take_i32 假设调用方提供的游标区间不重叠,若调用方出错将导致静默错误解析(非 panic,返回 None),可能导致整个 batch 被丢弃。
  3. 缓存一致性问题:正则 admission 缓存 (ADMISSION_CACHE) 使用全局 Mutex<HashMap>,在多线程 ingress 场景下可能引入锁竞争;全量清空策略虽简单但可能降低高频重复 pattern 的命中率。
  4. 采样参数拒绝 n>1sampling.rs 明确拒绝 n > 1(并行采样),这与 Python 端的行为一致,但若 Rust 服务器被配置为允许并行采样请求,将返回 400 而非降级。

影响范围:核心消息路径的 Rust 化。依赖此 PR 的后续拆分(如 #32240、#32343)才能构建完整的 Rust 处理管线。对用户透明,所有接口行为与 Python 版本保持兼容。对团队而言,此 PR 定义了 egress 数据类型和协议约定,后续开发者必须遵循。

性能:正则 admission 缓存可减少重复请求的 CPU 开销;采样参数验证前置到 Rust 侧可减轻调度器负载。但当前未合并完整管线,性能收益需与 ingress 等模块配合后才能体现。

缺少测试覆盖 核心路径变更 cursoring 重叠假设 未实现 Display 影响调试

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论