PP+MTP r36 de-GLOO: fix + verdict (PP+MTP loses i8k 1.81x to B', graph mode correctness-broken, A16 restored)

This commit is contained in:
yy-fighting 2026-09-08 22:21:20 +08:00
parent 123023b6ae
commit ffda226f5a
10 changed files with 3478 additions and 1 deletions

View File

@ -14,7 +14,7 @@
| 60.5 | `glm53-nvfp4`Up 29h | **NVFP4 团队生产**deploy_glm53_605.shmd5 fcd9109b。生产机铁律不实验、不重启、不覆盖脚本 | `profiles/pro6000/glm53_nvfp4_pro6000_sglang_tp8eagle.env`(方案 A 口径) | | 60.5 | `glm53-nvfp4`Up 29h | **NVFP4 团队生产**deploy_glm53_605.shmd5 fcd9109b。生产机铁律不实验、不重启、不覆盖脚本 | `profiles/pro6000/glm53_nvfp4_pro6000_sglang_tp8eagle.env`(方案 A 口径) |
| 60.6 | 无容器,但 8 卡被外部裸金属实验占用(`/data/hzy/sparse-opd-*`09-08 晚实测) | 外部任务,勿动(此前台账漏记) | — | | 60.6 | 无容器,但 8 卡被外部裸金属实验占用(`/data/hzy/sparse-opd-*`09-08 晚实测) | 外部任务,勿动(此前台账漏记) | — |
| 60.7 | 基本空4 卡仍有 `/home/user/dirA_exp` 外部小任务09-08 晚实测) | 09-08 已拆除清空(方案 F 前身单机实验 + 场景一深优资产留盘),不再恢复 | — | | 60.7 | 基本空4 卡仍有 `/home/user/dirA_exp` 外部小任务09-08 晚实测) | 09-08 已拆除清空(方案 F 前身单机实验 + 场景一深优资产留盘),不再恢复 | — |
| 60.8 | `glm53-nvfp4`Up:30000restart=unless-stopped | **A16 口径在役**09-08 i8k/o1k/c16 压测收官保留tp8eagle.env 场景二变体 = A + MRR32 + graph bs 4/8/12/16池 276,864质量门 7/7。同场景实测最优为 BTP4PP2out 265 vs A16 219 tok/s如转纯吞吐用途可切 B。压测前经授权清理了 GPU3/6 的 dirA_ext 外部 eval 进程 | `profiles/pro6000/glm53_nvfp4_pro6000_sglang_tp8eagle.env`+注释中场景二高并发变体);压测记录 `experiments/pro6000/glm53_nvfp4_i8k_o1k_c16_bench/` | | 60.8 | `glm53-nvfp4`Up:30000restart=unless-stopped | **A16 口径在役**09-08 i8k/o1k/c16 压测收官保留tp8eagle.env 场景二变体 = A + MRR32 + graph bs 4/8/12/16池 276,864质量门 7/7。同场景实测最优为 BTP4PP2out 265 vs A16 219 tok/s如转纯吞吐用途可切 B。压测前经授权清理了 GPU3/6 的 dirA_ext 外部 eval 进程。**09-08 晚 PP+MTP r36 实验后已恢复本口径**(实验判决:去 GLOO 修复成立且稳定,但 i8k 场景 PP+MTP 111-116s 差 B' 61.5s 1.81×,图模式 r36g 为正确性灾难——详见 `experiments/pro6000/glm53_ppmtp_r36_degloo/`;语料消耗至 21,235,008/21,296,780 无复测余量) | `profiles/pro6000/glm53_nvfp4_pro6000_sglang_tp8eagle.env`+注释中场景二高并发变体);压测记录 `experiments/pro6000/glm53_nvfp4_i8k_o1k_c16_bench/`、`experiments/pro6000/glm53_ppmtp_r36_degloo/` |
## 方案 A-F 一览GLM-5.3-NVFP4 @ pro60002026-09-08 双场景报告口径) ## 方案 A-F 一览GLM-5.3-NVFP4 @ pro60002026-09-08 双场景报告口径)

View File

@ -0,0 +1,95 @@
# GLM-5.3-NVFP4 PP+MTP r36去边界 GLOO修复与判决 — 2026-09-08 @174.1.60.8
接续 `glm53_nvfp4_i8k_o1k_c16_bench`A16/B 基线)与 r33 PP+MTP 深挖ppmtp_deepdive
目标:修复自研 TP4PP2+MTP 补丁并让 MTP 与 PP 双优势同时兑现;判负则 60.8 恢复 A16 在役。
## 结论速览
| 判决项 | 结果 |
|---|---|
| 去 GLOOr36 | **成立**稳定conc_test×3 + cc16 杀手×2 零崩溃、可复现 ±1%prefill 密集档 8×16384 **-22%**46.2→36.0s |
| i8k/o1024/cc16/nreq32 判决PG19 语料) | **PP+MTP 判负**e2e P50 **111.2/116.3s** vs B' 同栈对照 **61.5/61.6s**≡B 基线 61.8svs A16 66.2s**差 1.81×**,双门未过 |
| 缺口归属 | **100% 属 MTP 机制**B' 在同一 r36 栈上完美复现 B 基线栈本身零回退verify-step ~235ms@cc16 / ~148ms@cc8break-even 需 <148ms=52ms×accept2.85 |
| Phase4 图模式r36g | **正确性灾难,判死**`cuda graph: True` 生效但 accept **2.85→1.05**(接受率 0.02、GSM8K 输出数字噪声draft+verify 双坏)。**FORCE_EAGER 双开关是承重墙**,非性能旋钮 |
| sanitizer 根治路线 | **不可行**320 错误全为 NCCL 探测噪声error 209真死因=memcheck 开销 ~10GB/卡 吃掉 0.88 的 KV 余量(最低可行 memfrac 0.959,复现条件被破坏)。沿用 mask=1274.9ms/轮)+ 实证稳定性,风险显式记录 |
| 在役状态 | **A16 已恢复**TP8 EAGLE 4/1/5@0.90 + MRR32 + graph bs 4/8/12/16restart=unless-stopped |
## r35→r36一次死锁与一次修复本次核心工程产出
**r35**scheduler_pp_mixin_r35.pyr33 fb8c96f7 基础上去 GLOO
- dict 通道元数据:`send_object`2 条 GLOO/字典 × 3 字典/轮)→ **8KB 固定尺寸 NCCL uint8 缓冲**pickle 零填充,`pickle.loads` 天然忽略 STOP 后字节——已验证CPU 张量随车 GPU 化 + `__pp_cpu_keys__` 回填。
- pyobj 通道(请求中继/控制面GLOO 阻塞 `dist.send`**两段式 NCCL**int64 size + uint8 payload空载优化 size=0`world_group.device_group`(与 dict 通道的 pp_group 通信器隔离FIFO 不串扰)。
- `_pp_nccl_stage_guard`fill-event 门 + 双流 record_streamr33 实证isend 内核跑在调度器 ambient 流而非 send_streamguard 必须在 isend 入队前)。
- 修掉两处自捕 bugguard 原在 isend 之后(无效门控)→ 重排;`P2PWork` 须钉住实际发送对象contiguous() 副本)。
**r35 首部署死锁**py-spy 双端定格):
- PP0_TP0 卡在 `_pp_send_pyobj_nccl` isend 内(懒加载 2-rank NCCL communicator 创建自旋102% CPU
- PP1_TP0 停在 `request_receiver.py:142 _pull_raw_reqs` 的 GLOO `point_to_point_pyobj` irecv
- **根因:通道断裂**——请求中继的发送端在 mixin 里(已切 NCCL接收端却在 `scheduler_components/request_receiver.py`(补丁树之外,仍走 GLOO两端永不相遇。全库仅 2 个 `point_to_point_pyobj` 使用者,无第三者顺序风险。
**r36 修复**request_receiver_degloo.py第 8 个挂载文件):
- `_pull_raw_reqs` p2p 分支pp_rank≠0 且 attn_tp/cp_rank==0同步切两段式 NCCL 接收,与发送端同通道同 FIFO 语义(含空载 size=0 锁步)。
- `SGLANG_PP_DEGLOO` 门控,=0 整体回退 r33 GLOO 行为。
- 部署即刻通过:健康 380s与 r33 完全一致)、池 384,960、冒烟 ALL-OK、accept 2.09-3.54。
## 判决数据i8k/o1024/cc16/nreq32PG19 真实语料温度0warm+2
| 配置 | run | e2e P50 | out tok/s | TPOT P50 | TTFT P50 | accept | retract |
|---|---|---|---|---|---|---|---|
| PP+MTP r36 | 9402 | 111.2s | 136.0 | 105.2ms | 4.09s | 2.873 | 0 |
| PP+MTP r36 | 9403 | 116.3s | 133.6 | 106.5ms | 4.07s | 2.842 | 0 |
| B'(同栈 nomtp | 9405 | **61.5s** | **266.7** | 51.4ms | 9.24s | — | 0 |
| B'(同栈 nomtp | 9406 | **61.6s** | **265.7** | 51.8ms | 9.29s | — | 0 |
| 参照B 基线 | 09-08 上午 | 61.8s | 264.9 | 50.3ms | — | — | 0 |
| 参照A16 基线 | 09-08 上午 | 66.2s | 218.9 | 55ms | 4.0s | 2.50 | 7/7 |
语料窗口PP+MTP [19,662,144, 20,448,576)、B' [20,448,576, 21,235,008),各 3×262,144=warm+2i8k 全唯一 prompt 窗口=32×8192
消耗至 **21,235,008 / 21,296,780**(余 61,772无复测余量复测须换语料。run-id 9401-94069401/9404 为 warm不入判决
缓存纪律验证:两配置 hit_rate=0.0窗口互不重叠、radix 干净)。
**MTP 增量拆解**:稳态 16 并发 decode 190 tok/s → per-req 11.9 tok/s → step≈239ms8 并发段 step≈148ms≈A16 的 137ms
步长随批超线性 + spec 每轮 4 次 forward3 draft + 1 verify全 eager+ 边界中继 → 2.85× 放大率不敌 4.5× 步长。MTP 在本栈 cc16 为**净负收益 1.59×**。
## cc16 杀手 profile16×16384/512 random-ids历史竞态触发器语料免费
| 配置 | 种子 | wall | out tok/s | E2E 中位 | TPOT 中位 | accept |
|---|---|---|---|---|---|---|
| r36 eager | 6402 | 105.3s | 77.8 | 72.2s | 72.5ms | 3.76 |
| r36 eager | 6403 | 106.0s | 77.3 | 72.6s | 73.1ms | 3.69 |
零崩溃、±1% 可复现SYNC_MASK=127 掩码下;历史 SYNC_MASK=0 会触发 async-IMA 的同一 profile
注意random-ids 的 accept 3.7+ 相对真实语料 2.85 虚高,与既有结论一致。
## Phase4 图模式判决r36g = r36 + FORCE_EAGER_DRAFT/VERIFY=0
- `cuda graph: True` 确认生效conc_test ALL-OK 但 **accept len 1.05-1.11 / 接受率 0.02-0.04**draft 提议几乎全灭)
- 质量门GSM8K-1 输出即数字噪声 → **verify 图同样损坏**(非仅 draft
- r33 注释"draft-decode cuda-graph capture issues on PP stages"实证为 draft+verify 双重捕获缺陷
- **修复方向**(未来工程):按 cudagraph×spec 冲突清单(桶形状/地址预分配/host 控制流等 8 类)做预分配补丁 + 图外 canary非配置级可解
## 文件清单md5
| 文件 | md5 |
|---|---|
| patches/scheduler_pp_mixin_r35.py | efce72af34b2e59f20d2eb0874f9f2b2 |
| patches/request_receiver_degloo.py | ad6c978a3e1efea395ec3d537a4fedd8 |
| scripts/deploy_ppmtp_r35.sh | 86fb4b576fb04eef7b71628324796ab7 |
| scripts/deploy_ppmtp_r36.sh | 7b80f4af0debccef8e87df4a2532d8dc |
| scripts/deploy_ppmtp_r36g.sh | dca9dd387a5e980dc81e25590b895646 |
| scripts/deploy_ppmtp_sanitize.sh | 5bb5a7dfaf00569d5408e487267129a6 |
| scripts/killer_cc16.sh | ac3dda7e05b416aecf1da978325e2584 |
| patches/request_receiver_orig_0901.py | vanilla 参考(镜像原版) |
results/bench_logs/9401-9406 六轮语料 SUMMARY、killer s6402/6403、质量门 r367/7/r36g灾难证据
四个部署日志、san_crash_prehealth_s6401.log290Ksanitizer 不可行的完整证据链320×error209 噪声 + KV OOM 死因)。
60.8 服务器侧资产:/root/scheduler_pp_mixin_r35.py、request_receiver_degloo.py、deploy_ppmtp_r35/36/36g.sh、
sglang_patch2/r33 七件套不动、corpus_ids.json消耗至 21,235,008
## 风险与遗留显式记录
1. **async-IMA 竞态未根治**sanitizer 路线不可行见上r36 全部验证在 SYNC_MASK=127 掩码下通过。
动掩码/改流拓扑/图化任一变动都可能重新暴露,动前必须重跑全套稳定性门。
2. **PP+MTP 的兑现条件**后续若再战verify-step 需压到 <148ms 235@cc16/148@cc8)——图捕获修复Phase4 深水区
是最大单点杠杆;否则 MTP 在 PP 内为净负收益,不如 BTP4PP2 nomtp
3. A16 三参数口径MRR32 + cuda-graph-max-bs-decode 16 + bs 4/8/12/16已从档案找回并用于恢复部署。

View File

@ -0,0 +1,336 @@
from __future__ import annotations
from dataclasses import dataclass
from http import HTTPStatus
from typing import (
TYPE_CHECKING,
Any,
Callable,
List,
Optional,
Union,
)
import os
import pickle
import torch
import torch.distributed
import zmq
from torch.distributed import barrier
from sglang.srt.disaggregation.utils import prepare_abort
from sglang.srt.environ import envs
from sglang.srt.managers.io_struct import (
BatchTokenizedEmbeddingReqInput,
BatchTokenizedGenerateReqInput,
TokenizedEmbeddingReqInput,
TokenizedGenerateReqInput,
sock_recv,
)
from sglang.srt.managers.mm_utils import (
has_shm_features,
unwrap_shm_features,
)
from sglang.srt.runtime_context import get_disagg, get_parallel, is_ep_scale_joiner
from sglang.srt.utils import (
broadcast_pyobj,
point_to_point_pyobj,
)
from sglang.srt.utils.nvtx_utils import scheduler_nvtx_method
if TYPE_CHECKING:
from sglang.srt.configs.model_config import ModelConfig
from sglang.srt.distributed.parallel_state_wrapper import ParallelState
from sglang.srt.rust_server.server import RustServer
from sglang.srt.server_args import ServerArgs
from sglang.test.scripted_runtime.scheduler_hook import ScriptedSchedulerHook
from sglang.test.scripted_runtime.tokenizer_recv_proxy import (
ScriptedTokenizerRecvProxy,
)
# [r36 de-GLOO] Request-relay receiver side. The sender moved to a two-phase
# NCCL exchange on world_group.device_group (scheduler_pp_mixin
# _pp_send_pyobj_nccl); this end MUST use the same channel, or the pair never
# meets — a GLOO irecv here while the peer isends on NCCL wedges the first
# request (observed as the 09-08 r35 startup deadlock: PP0 stuck inside
# isend's lazy 2-rank comm init, PP1 parked in the GLOO irecv).
_PP_DEGLOO = os.getenv("SGLANG_PP_DEGLOO", "1") == "1"
def _pp_recv_pyobj_nccl(world_group, global_src: int):
"""[r36 de-GLOO] Two-phase NCCL object recv: int64 size, then uint8
payload. Mirrors scheduler_pp_mixin._pp_send_pyobj_nccl, including the
empty optimization (size=0 means no payload follows), so the relay stays
in FIFO lockstep on the same pair communicator the sender uses."""
device = torch.device(torch.cuda.current_device())
group = world_group.device_group
size_tensor = torch.empty(1, dtype=torch.int64, device=device)
work = torch.distributed.irecv(size_tensor, global_src, group=group)
work.wait()
size = int(size_tensor.item())
if size == 0:
return []
payload_tensor = torch.empty(size, dtype=torch.uint8, device=device)
work = torch.distributed.irecv(payload_tensor, global_src, group=group)
work.wait()
return pickle.loads(bytes(payload_tensor.cpu().numpy().tobytes()))
@dataclass(kw_only=True, slots=True, frozen=True)
class SchedulerRequestReceiver:
recv_from_tokenizer: Union[zmq.Socket, ScriptedTokenizerRecvProxy, RustServer]
recv_from_rpc: Optional[zmq.Socket]
recv_skipper: Any
input_blocker: Any
mm_receiver: Any
ps: ParallelState
tp_group: Any
tp_cpu_group: Any
attn_tp_group: Any
attn_tp_cpu_group: Any
attn_cp_group: Any
attn_cp_cpu_group: Any
world_group: Any
server_args: ServerArgs
model_config: ModelConfig
max_recv_per_poll: int
stream_output: Callable[..., None]
get_last_batch: Callable[[], Any]
scripted_scheduler_hook: Optional[ScriptedSchedulerHook] = None
def recv_limit_reached(self, num_recv_reqs: int) -> bool:
if self.max_recv_per_poll < 0:
return False
return num_recv_reqs >= self.max_recv_per_poll
@scheduler_nvtx_method("scheduler.recv_requests")
def recv_requests(
self,
) -> List[Union[TokenizedGenerateReqInput, TokenizedEmbeddingReqInput, Any]]:
"""Receive results at tp_rank = 0 and broadcast it to all other TP ranks."""
if self.scripted_scheduler_hook is not None:
self.scripted_scheduler_hook.step()
if self.recv_skipper is not None:
if not self.recv_skipper.handle(self.get_last_batch()):
return []
recv_reqs = self._pull_raw_reqs()
if self.input_blocker is not None:
recv_reqs = self.input_blocker.handle(recv_reqs)
recv_reqs = self._broadcast_reqs_across_ranks(recv_reqs)
if self.ps.pp_rank == 0:
self.unwrap_pickle_wrapper(recv_reqs)
recv_reqs = self._apply_mm_receiver(recv_reqs)
self._finalize_shm_features(recv_reqs)
return recv_reqs
def _pull_raw_reqs(self) -> Optional[List]:
if self.ps.pp_rank == 0:
if self.ps.attn_tp_rank == 0 and self.ps.attn_cp_rank == 0:
recv_reqs = []
# Rust ringbuffer backend: drain the in-process ring fed by the
# embedded Rust TokenizerManager instead of a zmq socket. Same
# non-blocking, msgpack-decoded contract as the zmq path below.
if envs.SGLANG_RUST_SERVER.get():
recv_reqs.extend(
self.recv_from_tokenizer.drain(self.max_recv_per_poll)
)
return recv_reqs
while True:
try:
if self.recv_limit_reached(len(recv_reqs)):
break
recv_req = sock_recv(self.recv_from_tokenizer, zmq.NOBLOCK)
except zmq.ZMQError:
break
recv_reqs.append(recv_req)
while True:
try:
if self.recv_limit_reached(len(recv_reqs)):
break
recv_rpc = sock_recv(self.recv_from_rpc, zmq.NOBLOCK)
except zmq.ZMQError:
break
recv_reqs.append(recv_rpc)
else:
recv_reqs = None
else:
if self.ps.attn_tp_rank == 0 and self.ps.attn_cp_rank == 0:
dp_offset = (
self.ps.attn_dp_rank * self.ps.attn_cp_size * self.ps.attn_tp_size
)
if _PP_DEGLOO:
recv_reqs = _pp_recv_pyobj_nccl(
self.world_group,
(self.ps.pp_rank - 1) * self.ps.tp_size + dp_offset,
)
else:
recv_reqs = point_to_point_pyobj(
[],
self.ps.pp_rank * self.ps.tp_size + dp_offset,
self.world_group.cpu_group,
(self.ps.pp_rank - 1) * self.ps.tp_size + dp_offset,
self.ps.pp_rank * self.ps.tp_size + dp_offset,
)
else:
recv_reqs = None
return recv_reqs
def _broadcast_reqs_across_ranks(self, recv_reqs: Optional[List]) -> List:
if get_parallel().enable_dp_attention:
if self.ps.attn_tp_rank == 0 and self.ps.attn_cp_rank == 0:
work_reqs, control_reqs = self._split_work_and_control_reqs(recv_reqs)
else:
work_reqs = None
control_reqs = None
if self.ps.attn_tp_size != 1:
work_reqs = broadcast_pyobj(
work_reqs,
self.attn_tp_group.rank,
self.attn_tp_cpu_group,
src=self.attn_tp_group.ranks[0],
)
if self.ps.attn_cp_size != 1:
work_reqs = broadcast_pyobj(
work_reqs,
self.attn_cp_group.rank,
self.attn_cp_cpu_group,
src=self.attn_cp_group.ranks[0],
)
# When dp_attention_local_control_broadcast is enabled, each DP
# group leader already receives control messages from the DP
# controller, so we broadcast within attn_tp_group + attn_cp_group
# instead of the full tp_group. This avoids an expensive
# all-ranks gloo sync.
_local_ctrl = (
get_parallel().enable_dp_attention_local_control_broadcast
or is_ep_scale_joiner()
)
if _local_ctrl:
if self.ps.attn_tp_size != 1:
control_reqs = broadcast_pyobj(
control_reqs,
self.attn_tp_group.rank,
self.attn_tp_cpu_group,
src=self.attn_tp_group.ranks[0],
)
if self.ps.attn_cp_size != 1:
control_reqs = broadcast_pyobj(
control_reqs,
self.attn_cp_group.rank,
self.attn_cp_cpu_group,
src=self.attn_cp_group.ranks[0],
)
elif self.ps.tp_size != 1:
control_reqs = broadcast_pyobj(
control_reqs,
self.tp_group.rank,
self.tp_cpu_group,
src=self.tp_group.ranks[0],
)
recv_reqs = work_reqs + control_reqs
elif self.ps.tp_size != 1:
recv_reqs = broadcast_pyobj(
recv_reqs,
self.tp_group.rank,
self.tp_cpu_group,
src=self.tp_group.ranks[0],
)
return recv_reqs
def unwrap_pickle_wrapper(self, recv_reqs: Optional[List]) -> None:
if not recv_reqs:
return
for req in recv_reqs:
if isinstance(req, (TokenizedGenerateReqInput, TokenizedEmbeddingReqInput)):
req.unwrap_pickle_fields()
elif isinstance(
req, (BatchTokenizedGenerateReqInput, BatchTokenizedEmbeddingReqInput)
):
for sub_req in req:
sub_req.unwrap_pickle_fields()
def _apply_mm_receiver(self, recv_reqs: List) -> List:
# Process MM requests under EPD-disaggregation mode
if (
self.ps.pp_rank == 0
and get_disagg().language_only
and get_disagg().encoder_transfer_backend
in ["zmq_to_scheduler", "mooncake"]
):
recv_reqs, abort_reqs = self.mm_receiver.process_waiting_requests(recv_reqs)
for req, error_msg, error_code in abort_reqs:
if error_code is None:
status_code = HTTPStatus.INTERNAL_SERVER_ERROR
elif isinstance(error_code, HTTPStatus):
status_code = error_code
else:
status_code = HTTPStatus(int(error_code))
prepare_abort(req, error_msg, status_code=status_code)
self.stream_output([req], req.return_logprob)
return recv_reqs
def _finalize_shm_features(self, recv_reqs: Optional[List]) -> None:
# Unwrap shared memory features AFTER all broadcasts complete,
# so that ShmPointerMMData metadata (not full tensor data) is what
# gets serialized during broadcast_pyobj.
if recv_reqs:
if self.model_config.is_multimodal and has_shm_features(recv_reqs):
# The broadcast source returns with its original objects while
# peer ranks may still be unpickling ShmPointerMMData
# (-> shm_open). Synchronize the same CPU groups that carried
# SHM-backed work requests before materialize() unlinks them.
if get_parallel().enable_dp_attention:
if self.ps.attn_tp_size > 1:
barrier(group=self.attn_tp_cpu_group)
if self.ps.attn_cp_size > 1:
barrier(group=self.attn_cp_cpu_group)
elif self.ps.tp_size > 1:
barrier(group=self.tp_cpu_group)
for req in recv_reqs:
unwrap_shm_features(req)
def _split_work_and_control_reqs(self, recv_reqs: List):
work_reqs = [
req
for req in recv_reqs
if isinstance(
req,
(
TokenizedGenerateReqInput,
TokenizedEmbeddingReqInput,
BatchTokenizedGenerateReqInput,
BatchTokenizedEmbeddingReqInput,
),
)
]
control_reqs = [
req
for req in recv_reqs
if not isinstance(
req,
(
TokenizedGenerateReqInput,
TokenizedEmbeddingReqInput,
BatchTokenizedGenerateReqInput,
BatchTokenizedEmbeddingReqInput,
),
)
]
return work_reqs, control_reqs

View File

@ -0,0 +1,298 @@
from __future__ import annotations
from dataclasses import dataclass
from http import HTTPStatus
from typing import (
TYPE_CHECKING,
Any,
Callable,
List,
Optional,
Union,
)
import zmq
from torch.distributed import barrier
from sglang.srt.disaggregation.utils import prepare_abort
from sglang.srt.environ import envs
from sglang.srt.managers.io_struct import (
BatchTokenizedEmbeddingReqInput,
BatchTokenizedGenerateReqInput,
TokenizedEmbeddingReqInput,
TokenizedGenerateReqInput,
sock_recv,
)
from sglang.srt.managers.mm_utils import (
has_shm_features,
unwrap_shm_features,
)
from sglang.srt.runtime_context import get_disagg, get_parallel, is_ep_scale_joiner
from sglang.srt.utils import (
broadcast_pyobj,
point_to_point_pyobj,
)
from sglang.srt.utils.nvtx_utils import scheduler_nvtx_method
if TYPE_CHECKING:
from sglang.srt.configs.model_config import ModelConfig
from sglang.srt.distributed.parallel_state_wrapper import ParallelState
from sglang.srt.rust_server.server import RustServer
from sglang.srt.server_args import ServerArgs
from sglang.test.scripted_runtime.scheduler_hook import ScriptedSchedulerHook
from sglang.test.scripted_runtime.tokenizer_recv_proxy import (
ScriptedTokenizerRecvProxy,
)
@dataclass(kw_only=True, slots=True, frozen=True)
class SchedulerRequestReceiver:
recv_from_tokenizer: Union[zmq.Socket, ScriptedTokenizerRecvProxy, RustServer]
recv_from_rpc: Optional[zmq.Socket]
recv_skipper: Any
input_blocker: Any
mm_receiver: Any
ps: ParallelState
tp_group: Any
tp_cpu_group: Any
attn_tp_group: Any
attn_tp_cpu_group: Any
attn_cp_group: Any
attn_cp_cpu_group: Any
world_group: Any
server_args: ServerArgs
model_config: ModelConfig
max_recv_per_poll: int
stream_output: Callable[..., None]
get_last_batch: Callable[[], Any]
scripted_scheduler_hook: Optional[ScriptedSchedulerHook] = None
def recv_limit_reached(self, num_recv_reqs: int) -> bool:
if self.max_recv_per_poll < 0:
return False
return num_recv_reqs >= self.max_recv_per_poll
@scheduler_nvtx_method("scheduler.recv_requests")
def recv_requests(
self,
) -> List[Union[TokenizedGenerateReqInput, TokenizedEmbeddingReqInput, Any]]:
"""Receive results at tp_rank = 0 and broadcast it to all other TP ranks."""
if self.scripted_scheduler_hook is not None:
self.scripted_scheduler_hook.step()
if self.recv_skipper is not None:
if not self.recv_skipper.handle(self.get_last_batch()):
return []
recv_reqs = self._pull_raw_reqs()
if self.input_blocker is not None:
recv_reqs = self.input_blocker.handle(recv_reqs)
recv_reqs = self._broadcast_reqs_across_ranks(recv_reqs)
if self.ps.pp_rank == 0:
self.unwrap_pickle_wrapper(recv_reqs)
recv_reqs = self._apply_mm_receiver(recv_reqs)
self._finalize_shm_features(recv_reqs)
return recv_reqs
def _pull_raw_reqs(self) -> Optional[List]:
if self.ps.pp_rank == 0:
if self.ps.attn_tp_rank == 0 and self.ps.attn_cp_rank == 0:
recv_reqs = []
# Rust ringbuffer backend: drain the in-process ring fed by the
# embedded Rust TokenizerManager instead of a zmq socket. Same
# non-blocking, msgpack-decoded contract as the zmq path below.
if envs.SGLANG_RUST_SERVER.get():
recv_reqs.extend(
self.recv_from_tokenizer.drain(self.max_recv_per_poll)
)
return recv_reqs
while True:
try:
if self.recv_limit_reached(len(recv_reqs)):
break
recv_req = sock_recv(self.recv_from_tokenizer, zmq.NOBLOCK)
except zmq.ZMQError:
break
recv_reqs.append(recv_req)
while True:
try:
if self.recv_limit_reached(len(recv_reqs)):
break
recv_rpc = sock_recv(self.recv_from_rpc, zmq.NOBLOCK)
except zmq.ZMQError:
break
recv_reqs.append(recv_rpc)
else:
recv_reqs = None
else:
if self.ps.attn_tp_rank == 0 and self.ps.attn_cp_rank == 0:
dp_offset = (
self.ps.attn_dp_rank * self.ps.attn_cp_size * self.ps.attn_tp_size
)
recv_reqs = point_to_point_pyobj(
[],
self.ps.pp_rank * self.ps.tp_size + dp_offset,
self.world_group.cpu_group,
(self.ps.pp_rank - 1) * self.ps.tp_size + dp_offset,
self.ps.pp_rank * self.ps.tp_size + dp_offset,
)
else:
recv_reqs = None
return recv_reqs
def _broadcast_reqs_across_ranks(self, recv_reqs: Optional[List]) -> List:
if get_parallel().enable_dp_attention:
if self.ps.attn_tp_rank == 0 and self.ps.attn_cp_rank == 0:
work_reqs, control_reqs = self._split_work_and_control_reqs(recv_reqs)
else:
work_reqs = None
control_reqs = None
if self.ps.attn_tp_size != 1:
work_reqs = broadcast_pyobj(
work_reqs,
self.attn_tp_group.rank,
self.attn_tp_cpu_group,
src=self.attn_tp_group.ranks[0],
)
if self.ps.attn_cp_size != 1:
work_reqs = broadcast_pyobj(
work_reqs,
self.attn_cp_group.rank,
self.attn_cp_cpu_group,
src=self.attn_cp_group.ranks[0],
)
# When dp_attention_local_control_broadcast is enabled, each DP
# group leader already receives control messages from the DP
# controller, so we broadcast within attn_tp_group + attn_cp_group
# instead of the full tp_group. This avoids an expensive
# all-ranks gloo sync.
_local_ctrl = (
get_parallel().enable_dp_attention_local_control_broadcast
or is_ep_scale_joiner()
)
if _local_ctrl:
if self.ps.attn_tp_size != 1:
control_reqs = broadcast_pyobj(
control_reqs,
self.attn_tp_group.rank,
self.attn_tp_cpu_group,
src=self.attn_tp_group.ranks[0],
)
if self.ps.attn_cp_size != 1:
control_reqs = broadcast_pyobj(
control_reqs,
self.attn_cp_group.rank,
self.attn_cp_cpu_group,
src=self.attn_cp_group.ranks[0],
)
elif self.ps.tp_size != 1:
control_reqs = broadcast_pyobj(
control_reqs,
self.tp_group.rank,
self.tp_cpu_group,
src=self.tp_group.ranks[0],
)
recv_reqs = work_reqs + control_reqs
elif self.ps.tp_size != 1:
recv_reqs = broadcast_pyobj(
recv_reqs,
self.tp_group.rank,
self.tp_cpu_group,
src=self.tp_group.ranks[0],
)
return recv_reqs
def unwrap_pickle_wrapper(self, recv_reqs: Optional[List]) -> None:
if not recv_reqs:
return
for req in recv_reqs:
if isinstance(req, (TokenizedGenerateReqInput, TokenizedEmbeddingReqInput)):
req.unwrap_pickle_fields()
elif isinstance(
req, (BatchTokenizedGenerateReqInput, BatchTokenizedEmbeddingReqInput)
):
for sub_req in req:
sub_req.unwrap_pickle_fields()
def _apply_mm_receiver(self, recv_reqs: List) -> List:
# Process MM requests under EPD-disaggregation mode
if (
self.ps.pp_rank == 0
and get_disagg().language_only
and get_disagg().encoder_transfer_backend
in ["zmq_to_scheduler", "mooncake"]
):
recv_reqs, abort_reqs = self.mm_receiver.process_waiting_requests(recv_reqs)
for req, error_msg, error_code in abort_reqs:
if error_code is None:
status_code = HTTPStatus.INTERNAL_SERVER_ERROR
elif isinstance(error_code, HTTPStatus):
status_code = error_code
else:
status_code = HTTPStatus(int(error_code))
prepare_abort(req, error_msg, status_code=status_code)
self.stream_output([req], req.return_logprob)
return recv_reqs
def _finalize_shm_features(self, recv_reqs: Optional[List]) -> None:
# Unwrap shared memory features AFTER all broadcasts complete,
# so that ShmPointerMMData metadata (not full tensor data) is what
# gets serialized during broadcast_pyobj.
if recv_reqs:
if self.model_config.is_multimodal and has_shm_features(recv_reqs):
# The broadcast source returns with its original objects while
# peer ranks may still be unpickling ShmPointerMMData
# (-> shm_open). Synchronize the same CPU groups that carried
# SHM-backed work requests before materialize() unlinks them.
if get_parallel().enable_dp_attention:
if self.ps.attn_tp_size > 1:
barrier(group=self.attn_tp_cpu_group)
if self.ps.attn_cp_size > 1:
barrier(group=self.attn_cp_cpu_group)
elif self.ps.tp_size > 1:
barrier(group=self.tp_cpu_group)
for req in recv_reqs:
unwrap_shm_features(req)
def _split_work_and_control_reqs(self, recv_reqs: List):
work_reqs = [
req
for req in recv_reqs
if isinstance(
req,
(
TokenizedGenerateReqInput,
TokenizedEmbeddingReqInput,
BatchTokenizedGenerateReqInput,
BatchTokenizedEmbeddingReqInput,
),
)
]
control_reqs = [
req
for req in recv_reqs
if not isinstance(
req,
(
TokenizedGenerateReqInput,
TokenizedEmbeddingReqInput,
BatchTokenizedGenerateReqInput,
BatchTokenizedEmbeddingReqInput,
),
)
]
return work_reqs, control_reqs

File diff suppressed because it is too large Load Diff

View File

@ -0,0 +1,72 @@
#!/bin/bash
# PP+MTP r35 部署60.8, nightly-dev-cu13-20260901-07c8f729
# 与 deploy_ppmtp_mask.sh 差异scheduler 挂载切到 r35r33 + 去边界 GLOO
# SGLANG_PP_DEGLOO 默认开,=0 回退 r33 行为。eagle worker 仍挂 mask 版保留热开关。
# 用法: bash deploy_ppmtp_r35.sh "<并行参数>" [mtp|nomtp] [chunk] [memfrac] [degloo]
# 例: bash deploy_ppmtp_r35.sh "--tp 4 --pp-size 2 --disable-overlap-schedule --max-prefill-tokens 16384" mtp 8192 0.88 1
PAR=${1:?usage: deploy_ppmtp_r35.sh "<flags>" [mtp|nomtp] [chunk] [memfrac] [degloo]}
MTPMODE=${2:-nomtp}
CHUNK=${3:-8192}
MEMFRAC=${4:-0.88}
DEGLOO=${5:-1}
MTPARGS=""
if [ "$MTPMODE" = "mtp" ]; then
MTPARGS="--speculative-algorithm EAGLE --speculative-num-steps 3 --speculative-eagle-topk 1 --speculative-num-draft-tokens 4"
fi
docker update --restart=no glm53-nvfp4 >/dev/null 2>&1
for i in 1 2 3 4 5; do
docker rm -f glm53-nvfp4 >/dev/null 2>&1
sleep 2
docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$' || break
done
if docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$'; then
echo "ERROR: old container cannot be removed"; exit 1
fi
for i in $(seq 1 15); do ss -ltn 2>/dev/null | grep -q ":30000 " || break; sleep 2; done
docker run -d --name glm53-nvfp4 --gpus all --shm-size 64g --ipc=host --cap-add SYS_PTRACE \
-v /root/sglang_patch2/layer_setup.py:/sgl-workspace/sglang/python/sglang/srt/model_executor/model_runner_components/layer_setup.py:ro \
-v /root/sglang_patch2/validation_hook.py:/sgl-workspace/sglang/python/sglang/srt/arg_groups/validation_hook.py:ro \
-v /root/eagle_worker_v2_mask.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_v2.py:ro \
-v /root/sglang_patch2/eagle_worker_common.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_common.py:ro \
-v /root/sglang_patch2/deepseek_nextn.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_nextn.py:ro \
-v /root/scheduler_pp_mixin_r35.py:/sgl-workspace/sglang/python/sglang/srt/managers/scheduler_pp_mixin.py:ro \
-v /root/sglang_patch2/deepseek_v2.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_v2.py:ro \
-e SGLANG_PP_SPEC_DEBUG=0 \
-e SGLANG_PP_DEGLOO=${DEGLOO} \
-e SGLANG_PP_SPEC_FORCE_EAGER_DRAFT=1 \
-e SGLANG_PP_FORCE_EAGER_VERIFY=1 \
--restart no -p 30000:30000 \
-v /data/hf_models:/data/hf_models \
lmsysorg/sglang:nightly-dev-cu13-20260901-07c8f729 \
python3 -m sglang.launch_server \
--model-path /data/hf_models/GLM-5.3-NVFP4 \
--tp 8 \
--mem-fraction-static ${MEMFRAC} \
--max-running-requests 16 \
--chunked-prefill-size ${CHUNK} \
--disable-shared-experts-fusion \
--moe-runner-backend flashinfer_cutlass \
--disable-flashinfer-autotune \
--reasoning-parser glm45 --tool-call-parser glm47 \
--enable-hierarchical-cache --hicache-ratio 3 \
${MTPARGS} \
${PAR} \
--host 0.0.0.0 --port 30000
echo "deployed r35 (DEGLOO=${DEGLOO}): par=[${PAR}] mtp=${MTPMODE}; waiting for health..."
for i in $(seq 10 10 1800); do
code=$(curl -s -o /dev/null -m3 -w '%{http_code}' http://127.0.0.1:30000/health 2>/dev/null)
if [ "$code" = "200" ]; then
echo "healthy after ${i}s"
docker logs glm53-nvfp4 2>&1 | grep -oE "max_total_num_tokens = [0-9]+" | head -1
exit 0
fi
if ! docker ps --format '{{.Names}}' | grep -q '^glm53-nvfp4$'; then
echo "CONTAINER DIED after ${i}s"; docker logs --tail 60 glm53-nvfp4 2>&1 | grep -iE "error|assert|not support|incompatible" | tail -8; exit 1
fi
sleep 10
done
echo "TIMEOUT waiting for health"; exit 1

View File

@ -0,0 +1,74 @@
#!/bin/bash
# PP+MTP r36 部署60.8, nightly-dev-cu13-20260901-07c8f729
# r35 首部署死锁修复版r35 只把请求中继发送端切到 NCCL接收端在
# request_receiver.py:_pull_raw_reqs 仍走 GLOO通道断裂 → 首请求 wedge。
# r36 = r35 scheduler + request_receiver_degloo.py接收端同步切两段式 NCCL
# 用法: bash deploy_ppmtp_r36.sh "<并行参数>" [mtp|nomtp] [chunk] [memfrac] [degloo]
# 例: bash deploy_ppmtp_r36.sh "--tp 4 --pp-size 2 --disable-overlap-schedule --max-prefill-tokens 16384" mtp 8192 0.88 1
PAR=${1:?usage: deploy_ppmtp_r36.sh "<flags>" [mtp|nomtp] [chunk] [memfrac] [degloo]}
MTPMODE=${2:-nomtp}
CHUNK=${3:-8192}
MEMFRAC=${4:-0.88}
DEGLOO=${5:-1}
MTPARGS=""
if [ "$MTPMODE" = "mtp" ]; then
MTPARGS="--speculative-algorithm EAGLE --speculative-num-steps 3 --speculative-eagle-topk 1 --speculative-num-draft-tokens 4"
fi
docker update --restart=no glm53-nvfp4 >/dev/null 2>&1
for i in 1 2 3 4 5; do
docker rm -f glm53-nvfp4 >/dev/null 2>&1
sleep 2
docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$' || break
done
if docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$'; then
echo "ERROR: old container cannot be removed"; exit 1
fi
for i in $(seq 1 15); do ss -ltn 2>/dev/null | grep -q ":30000 " || break; sleep 2; done
docker run -d --name glm53-nvfp4 --gpus all --shm-size 64g --ipc=host --cap-add SYS_PTRACE \
-v /root/sglang_patch2/layer_setup.py:/sgl-workspace/sglang/python/sglang/srt/model_executor/model_runner_components/layer_setup.py:ro \
-v /root/sglang_patch2/validation_hook.py:/sgl-workspace/sglang/python/sglang/srt/arg_groups/validation_hook.py:ro \
-v /root/eagle_worker_v2_mask.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_v2.py:ro \
-v /root/sglang_patch2/eagle_worker_common.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_common.py:ro \
-v /root/sglang_patch2/deepseek_nextn.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_nextn.py:ro \
-v /root/scheduler_pp_mixin_r35.py:/sgl-workspace/sglang/python/sglang/srt/managers/scheduler_pp_mixin.py:ro \
-v /root/request_receiver_degloo.py:/sgl-workspace/sglang/python/sglang/srt/managers/scheduler_components/request_receiver.py:ro \
-v /root/sglang_patch2/deepseek_v2.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_v2.py:ro \
-e SGLANG_PP_SPEC_DEBUG=0 \
-e SGLANG_PP_DEGLOO=${DEGLOO} \
-e SGLANG_PP_SPEC_FORCE_EAGER_DRAFT=1 \
-e SGLANG_PP_FORCE_EAGER_VERIFY=1 \
--restart no -p 30000:30000 \
-v /data/hf_models:/data/hf_models \
lmsysorg/sglang:nightly-dev-cu13-20260901-07c8f729 \
python3 -m sglang.launch_server \
--model-path /data/hf_models/GLM-5.3-NVFP4 \
--tp 8 \
--mem-fraction-static ${MEMFRAC} \
--max-running-requests 16 \
--chunked-prefill-size ${CHUNK} \
--disable-shared-experts-fusion \
--moe-runner-backend flashinfer_cutlass \
--disable-flashinfer-autotune \
--reasoning-parser glm45 --tool-call-parser glm47 \
--enable-hierarchical-cache --hicache-ratio 3 \
${MTPARGS} \
${PAR} \
--host 0.0.0.0 --port 30000
echo "deployed r36 (DEGLOO=${DEGLOO}): par=[${PAR}] mtp=${MTPMODE}; waiting for health..."
for i in $(seq 10 10 1800); do
code=$(curl -s -o /dev/null -m3 -w '%{http_code}' http://127.0.0.1:30000/health 2>/dev/null)
if [ "$code" = "200" ]; then
echo "healthy after ${i}s"
docker logs glm53-nvfp4 2>&1 | grep -oE "max_total_num_tokens = [0-9]+" | head -1
exit 0
fi
if ! docker ps --format '{{.Names}}' | grep -q '^glm53-nvfp4$'; then
echo "CONTAINER DIED after ${i}s"; docker logs --tail 60 glm53-nvfp4 2>&1 | grep -iE "error|assert|not support|incompatible" | tail -8; exit 1
fi
sleep 10
done
echo "TIMEOUT waiting for health"; exit 1

View File

@ -0,0 +1,74 @@
#!/bin/bash
# PP+MTP r36 部署60.8, nightly-dev-cu13-20260901-07c8f729
# r35 首部署死锁修复版r35 只把请求中继发送端切到 NCCL接收端在
# request_receiver.py:_pull_raw_reqs 仍走 GLOO通道断裂 → 首请求 wedge。
# r36 = r35 scheduler + request_receiver_degloo.py接收端同步切两段式 NCCL
# 用法: bash deploy_ppmtp_r36.sh "<并行参数>" [mtp|nomtp] [chunk] [memfrac] [degloo]
# 例: bash deploy_ppmtp_r36.sh "--tp 4 --pp-size 2 --disable-overlap-schedule --max-prefill-tokens 16384" mtp 8192 0.88 1
PAR=${1:?usage: deploy_ppmtp_r36.sh "<flags>" [mtp|nomtp] [chunk] [memfrac] [degloo]}
MTPMODE=${2:-nomtp}
CHUNK=${3:-8192}
MEMFRAC=${4:-0.88}
DEGLOO=${5:-1}
MTPARGS=""
if [ "$MTPMODE" = "mtp" ]; then
MTPARGS="--speculative-algorithm EAGLE --speculative-num-steps 3 --speculative-eagle-topk 1 --speculative-num-draft-tokens 4"
fi
docker update --restart=no glm53-nvfp4 >/dev/null 2>&1
for i in 1 2 3 4 5; do
docker rm -f glm53-nvfp4 >/dev/null 2>&1
sleep 2
docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$' || break
done
if docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$'; then
echo "ERROR: old container cannot be removed"; exit 1
fi
for i in $(seq 1 15); do ss -ltn 2>/dev/null | grep -q ":30000 " || break; sleep 2; done
docker run -d --name glm53-nvfp4 --gpus all --shm-size 64g --ipc=host --cap-add SYS_PTRACE \
-v /root/sglang_patch2/layer_setup.py:/sgl-workspace/sglang/python/sglang/srt/model_executor/model_runner_components/layer_setup.py:ro \
-v /root/sglang_patch2/validation_hook.py:/sgl-workspace/sglang/python/sglang/srt/arg_groups/validation_hook.py:ro \
-v /root/eagle_worker_v2_mask.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_v2.py:ro \
-v /root/sglang_patch2/eagle_worker_common.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_common.py:ro \
-v /root/sglang_patch2/deepseek_nextn.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_nextn.py:ro \
-v /root/scheduler_pp_mixin_r35.py:/sgl-workspace/sglang/python/sglang/srt/managers/scheduler_pp_mixin.py:ro \
-v /root/request_receiver_degloo.py:/sgl-workspace/sglang/python/sglang/srt/managers/scheduler_components/request_receiver.py:ro \
-v /root/sglang_patch2/deepseek_v2.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_v2.py:ro \
-e SGLANG_PP_SPEC_DEBUG=0 \
-e SGLANG_PP_DEGLOO=${DEGLOO} \
-e SGLANG_PP_SPEC_FORCE_EAGER_DRAFT=0 \
-e SGLANG_PP_FORCE_EAGER_VERIFY=0 \
--restart no -p 30000:30000 \
-v /data/hf_models:/data/hf_models \
lmsysorg/sglang:nightly-dev-cu13-20260901-07c8f729 \
python3 -m sglang.launch_server \
--model-path /data/hf_models/GLM-5.3-NVFP4 \
--tp 8 \
--mem-fraction-static ${MEMFRAC} \
--max-running-requests 16 \
--chunked-prefill-size ${CHUNK} \
--disable-shared-experts-fusion \
--moe-runner-backend flashinfer_cutlass \
--disable-flashinfer-autotune \
--reasoning-parser glm45 --tool-call-parser glm47 \
--enable-hierarchical-cache --hicache-ratio 3 \
${MTPARGS} \
${PAR} \
--host 0.0.0.0 --port 30000
echo "deployed r36g GRAPH (DEGLOO=${DEGLOO}): par=[${PAR}] mtp=${MTPMODE}; waiting for health..."
for i in $(seq 10 10 1800); do
code=$(curl -s -o /dev/null -m3 -w '%{http_code}' http://127.0.0.1:30000/health 2>/dev/null)
if [ "$code" = "200" ]; then
echo "healthy after ${i}s"
docker logs glm53-nvfp4 2>&1 | grep -oE "max_total_num_tokens = [0-9]+" | head -1
exit 0
fi
if ! docker ps --format '{{.Names}}' | grep -q '^glm53-nvfp4$'; then
echo "CONTAINER DIED after ${i}s"; docker logs --tail 60 glm53-nvfp4 2>&1 | grep -iE "error|assert|not support|incompatible" | tail -8; exit 1
fi
sleep 10
done
echo "TIMEOUT waiting for health"; exit 1

View File

@ -0,0 +1,75 @@
#!/bin/bash
# PP+MTP 竞态定位部署60.8, nightly-dev-cu13-20260901-07c8f729
# 与 deploy_ppmtp_mask.sh 差异:
# 1) 入口用 compute-sanitizer --tool memcheck 包裹LD_PRELOAD 随 mp.spawn 传满 8 rank
# 2) SGLANG_PP_SPEC_SYNC_MASK=0 —— 拆掉 7 同步遮罩暴露竞态
# 3) 健康等待窗口放宽到 3600ssanitizer 下启动 ~20min+
# 用法: bash deploy_ppmtp_sanitize.sh "<并行参数>" [mtp|nomtp] [chunk] [memfrac]
# 例: bash deploy_ppmtp_sanitize.sh "--tp 4 --pp-size 2 --disable-overlap-schedule --max-prefill-tokens 16384" mtp
PAR=${1:?usage: deploy_ppmtp_sanitize.sh "<flags>" [mtp|nomtp] [chunk] [memfrac]}
MTPMODE=${2:-nomtp}
CHUNK=${3:-8192}
MEMFRAC=${4:-0.88}
MTPARGS=""
if [ "$MTPMODE" = "mtp" ]; then
MTPARGS="--speculative-algorithm EAGLE --speculative-num-steps 3 --speculative-eagle-topk 1 --speculative-num-draft-tokens 4"
fi
# 顽固容器清理(同 bench 版)
docker update --restart=no glm53-nvfp4 >/dev/null 2>&1
for i in 1 2 3 4 5; do
docker rm -f glm53-nvfp4 >/dev/null 2>&1
sleep 2
docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$' || break
done
if docker ps -a --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$'; then
echo "ERROR: old container cannot be removed"; exit 1
fi
for i in $(seq 1 15); do ss -ltn 2>/dev/null | grep -q ":30000 " || break; sleep 2; done
docker run -d --name glm53-nvfp4 --gpus all --shm-size 64g --ipc=host --cap-add SYS_PTRACE \
-v /root/sglang_patch2/layer_setup.py:/sgl-workspace/sglang/python/sglang/srt/model_executor/model_runner_components/layer_setup.py:ro \
-v /root/sglang_patch2/validation_hook.py:/sgl-workspace/sglang/python/sglang/srt/arg_groups/validation_hook.py:ro \
-v /root/eagle_worker_v2_mask.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_v2.py:ro \
-v /root/sglang_patch2/eagle_worker_common.py:/sgl-workspace/sglang/python/sglang/srt/speculative/eagle_worker_common.py:ro \
-v /root/sglang_patch2/deepseek_nextn.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_nextn.py:ro \
-v /root/sglang_patch2/scheduler_pp_mixin.py:/sgl-workspace/sglang/python/sglang/srt/managers/scheduler_pp_mixin.py:ro \
-v /root/sglang_patch2/deepseek_v2.py:/sgl-workspace/sglang/python/sglang/srt/models/deepseek_v2.py:ro \
-e SGLANG_PP_SPEC_DEBUG=0 \
-e SGLANG_PP_SPEC_SYNC_MASK=0 \
-e SGLANG_PP_SPEC_FORCE_EAGER_DRAFT=1 \
-e SGLANG_PP_FORCE_EAGER_VERIFY=1 \
--restart no -p 30000:30000 \
-v /data/hf_models:/data/hf_models \
lmsysorg/sglang:nightly-dev-cu13-20260901-07c8f729 \
compute-sanitizer --tool memcheck --launch-timeout 0 --target-processes all \
python3 -m sglang.launch_server \
--model-path /data/hf_models/GLM-5.3-NVFP4 \
--tp 8 \
--mem-fraction-static ${MEMFRAC} \
--max-running-requests 16 \
--chunked-prefill-size ${CHUNK} \
--disable-shared-experts-fusion \
--moe-runner-backend flashinfer_cutlass \
--disable-flashinfer-autotune \
--reasoning-parser glm45 --tool-call-parser glm47 \
--enable-hierarchical-cache --hicache-ratio 3 \
${MTPARGS} \
${PAR} \
--host 0.0.0.0 --port 30000
echo "deployed(sanitizer memcheck, SYNC_MASK=0): par=[${PAR}] mtp=${MTPMODE}; waiting for health (up to 3600s)..."
for i in $(seq 10 10 3600); do
code=$(curl -s -o /dev/null -m3 -w '%{http_code}' http://127.0.0.1:30000/health 2>/dev/null)
if [ "$code" = "200" ]; then
echo "healthy after ${i}s"
docker logs glm53-nvfp4 2>&1 | grep -oE "max_total_num_tokens = [0-9]+" | head -1
exit 0
fi
if ! docker ps --format '{{.Names}}' | grep -q '^glm53-nvfp4$'; then
echo "CONTAINER DIED after ${i}s"; docker logs --tail 80 glm53-nvfp4 2>&1 | tail -40; exit 1
fi
sleep 10
done
echo "TIMEOUT waiting for health"; exit 1

View File

@ -0,0 +1,64 @@
#!/bin/bash
# Sanitizer 竞态触发 + 崩溃收割60.8, deploy_ppmtp_sanitize.sh 配套)
# 自带健康等待;健康后打 bench cc16 剖面杀手负载16x16384->512 random-ids,
# flush-cache 队列churn——09-06/07 档案中唯一可靠的竞态触发器);
# 容器死亡时收割全量 docker logs 到 /root/san_crash_s<seed>.log。
# 用法: nohup bash killer_cc16.sh [seed] > /root/killer_s<seed>.log 2>&1 &
SEED=${1:-6401}
URL=http://127.0.0.1:30000
echo "killer: waiting for health (up to 3600s)..."
HEALTHY=0
for i in $(seq 1 360); do
code=$(curl -s -o /dev/null -m5 -w '%{http_code}' $URL/health 2>/dev/null)
if [ "$code" = "200" ]; then HEALTHY=1; break; fi
if ! docker ps --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$'; then
echo "CONTAINER-DIED-BEFORE-HEALTH at $(date +%T)"
docker logs glm53-nvfp4 > /root/san_crash_prehealth_s${SEED}.log 2>&1
echo "log saved: /root/san_crash_prehealth_s${SEED}.log"
exit 2
fi
sleep 10
done
[ $HEALTHY -eq 1 ] || { echo "HEALTH-TIMEOUT"; exit 4; }
echo "healthy at $(date +%T)"
docker exec glm53-nvfp4 python3 -m sglang.bench_serving \
--backend sglang --host 127.0.0.1 --port 30000 \
--dataset-name random-ids --tokenizer /data/hf_models/GLM-5.3-NVFP4 \
--num-prompts 16 --random-input-len 16384 --random-output-len 512 \
--random-range-ratio 1.0 --max-concurrency 16 \
--temperature 0.0 --flush-cache --warmup-requests 1 --seed ${SEED} \
--output-file /data/hf_models/bs_results/san_cc16_s${SEED}.json \
> /root/san_bench_s${SEED}.log 2>&1 &
BPID=$!
echo "bench pid=$BPID started $(date +%T)"
BADH=0
while kill -0 $BPID 2>/dev/null; do
sleep 20
if ! docker ps --format '{{.Names}}' 2>/dev/null | grep -q '^glm53-nvfp4$'; then
echo "CONTAINER-DIED at $(date +%T)"
docker logs glm53-nvfp4 > /root/san_crash_s${SEED}.log 2>&1
echo "full log saved: /root/san_crash_s${SEED}.log ($(wc -l < /root/san_crash_s${SEED}.log) lines)"
grep -nE "=========|Invalid|invalid|illegal|CUDA_ERROR|out of bounds|Error Summ" /root/san_crash_s${SEED}.log | head -40
exit 2
fi
code=$(curl -s -o /dev/null -m5 -w '%{http_code}' $URL/health 2>/dev/null)
if [ "$code" != "200" ]; then
BADH=$((BADH+1))
if [ $BADH -ge 3 ]; then
echo "HEALTH-DEAD at $(date +%T)"
docker logs glm53-nvfp4 > /root/san_crash_s${SEED}.log 2>&1
echo "full log saved: /root/san_crash_s${SEED}.log ($(wc -l < /root/san_crash_s${SEED}.log) lines)"
grep -nE "=========|Invalid|invalid|illegal|CUDA_ERROR|out of bounds|Error Summ" /root/san_crash_s${SEED}.log | head -40
exit 3
fi
else
BADH=0
fi
done
wait $BPID; RC=$?
echo "BENCH-FINISHED rc=$RC $(date +%T) —— 竞态未触发(负载全程存活)"
tail -8 /root/san_bench_s${SEED}.log
exit $RC