diff --git a/python/sglang/srt/arg_groups/validation_hook.py b/python/sglang/srt/arg_groups/validation_hook.py index b9d807bc2..44241c633 100644 --- a/python/sglang/srt/arg_groups/validation_hook.py +++ b/python/sglang/srt/arg_groups/validation_hook.py @@ -49,9 +49,28 @@ def check_server_args(server_args: Any): ) if cfg.pp_size > 1: - assert cfg.disable_overlap_schedule and cfg.speculative_algorithm is None, ( - "Pipeline parallelism is not compatible with overlap schedule, speculative decoding" - ) + if get_platform().is_npu: + # NPU: allow PP + EAGLE speculative decoding + assert cfg.disable_overlap_schedule, ( + "Pipeline parallelism is not compatible with overlap schedule" + ) + if cfg.speculative_algorithm is not None: + assert ( + cfg.speculative_algorithm.upper() == "EAGLE" + and not cfg.enable_multi_layer_eagle + ), ( + "Pipeline parallelism currently only supports EAGLE " + "(non-multi-layer) speculative decoding" + ) + assert cfg.disaggregation_mode == "prefill", ( + "NPU PP + speculative decoding (MTP) is only supported " + "on prefill nodes (disaggregation-mode=prefill)" + ) + else: + # Non-NPU: PP + speculative decoding is not supported + assert cfg.disable_overlap_schedule and cfg.speculative_algorithm is None, ( + "Pipeline parallelism is not compatible with overlap schedule, speculative decoding" + ) assert cfg.min_free_slots_delay is None, ( "--min-free-slots-delay is not supported with pipeline " "parallelism: allocatable slots per microbatch are bounded by " diff --git a/python/sglang/srt/disaggregation/ascend/conn.py b/python/sglang/srt/disaggregation/ascend/conn.py index bc44b5b95..474a30fce 100644 --- a/python/sglang/srt/disaggregation/ascend/conn.py +++ b/python/sglang/srt/disaggregation/ascend/conn.py @@ -15,6 +15,7 @@ from sglang.srt.disaggregation.mooncake.conn import ( MooncakeKVReceiver, MooncakeKVSender, ) +from sglang.srt.distributed import get_pp_group from sglang.srt.utils.network import get_local_ip_auto logger = logging.getLogger(__name__) @@ -119,28 +120,37 @@ class AscendKVManager(MooncakeKVManager): # dst_kv_ptrs: k_data, v_data, index_k_data(optional) # state_type is accepted for parity with the common disaggregation path; # the NPU kv_buf_groups slicing below is state-type agnostic. - start_layer = self.kv_args.prefill_start_layer kv_buf_groups = getattr(self.kv_args, "kv_buf_groups", 1) - total_kv_layers = getattr(self.kv_args, "total_kv_layers", 0) + hidden_kv_layers = getattr(self.kv_args, "hidden_kv_layers", 0) + draft_kv_layers = getattr(self.kv_args, "draft_kv_layers", 0) src_layers = len(src_kv_ptrs) // kv_buf_groups - # When only speculative-algorithm is enabled for decode - # the KV has one more layer than prefill. - # The draft layer needs to be skipped. - dst_total_layers = ( - min(len(dst_kv_ptrs) // kv_buf_groups, total_kv_layers) - if total_kv_layers - else len(dst_kv_ptrs) // kv_buf_groups - ) - end_layer = start_layer + src_layers - if src_layers == dst_total_layers: + dst_layers = len(dst_kv_ptrs) // kv_buf_groups + if src_layers == dst_layers: sliced_dst_kv_ptrs = dst_kv_ptrs else: sliced_dst_kv_ptrs = [] + start_layer = self.kv_args.prefill_start_layer + transfer_draft_kv = get_pp_group().is_last_rank and draft_kv_layers + if transfer_draft_kv: + end_layer = start_layer + src_layers - draft_kv_layers + else: + end_layer = start_layer + src_layers + + # target kv for i in range(kv_buf_groups): - layer_offset = i * dst_total_layers + layer_offset = i * hidden_kv_layers sliced_dst_kv_ptrs.extend( dst_kv_ptrs[layer_offset + start_layer : layer_offset + end_layer] ) + # draft kv + if transfer_draft_kv: + for i in range(kv_buf_groups): + layer_offset = ( + i * draft_kv_layers + kv_buf_groups * hidden_kv_layers + ) + sliced_dst_kv_ptrs.extend( + dst_kv_ptrs[layer_offset : layer_offset + draft_kv_layers] + ) layers_current_pp_stage = len(src_kv_ptrs) return src_kv_ptrs, sliced_dst_kv_ptrs, layers_current_pp_stage diff --git a/python/sglang/srt/disaggregation/base/conn.py b/python/sglang/srt/disaggregation/base/conn.py index 2dea4485d..733128b3c 100644 --- a/python/sglang/srt/disaggregation/base/conn.py +++ b/python/sglang/srt/disaggregation/base/conn.py @@ -92,7 +92,9 @@ class KVArgs: # Only used of npu, for kv buf groups kv_buf_groups: int # Only used of npu, for decode total kv layers - total_kv_layers: int + hidden_kv_layers: int + # Only used of npu, for decode total kv layers + draft_kv_layers: int class KVPoll: diff --git a/python/sglang/srt/disaggregation/prefill.py b/python/sglang/srt/disaggregation/prefill.py index 9ecdb2775..13348fbf7 100644 --- a/python/sglang/srt/disaggregation/prefill.py +++ b/python/sglang/srt/disaggregation/prefill.py @@ -56,6 +56,7 @@ from sglang.srt.disaggregation.utils import ( prepare_abort, setup_state_kv_args, ) +from sglang.srt.distributed import get_pp_group from sglang.srt.environ import envs from sglang.srt.managers.schedule_batch import ( FINISH_ABORT, @@ -246,7 +247,11 @@ class PrefillBootstrapQueue: else getattr(self.token_to_kv_pool, "end_layer", None) ) - draft_kv_pool = self.draft_token_to_kv_pool if transfer_draft_cache else None + draft_kv_pool = ( + self.draft_token_to_kv_pool + if transfer_draft_cache and (not _is_npu or get_pp_group().is_last_rank) + else None + ) num_draft_entries = 0 if draft_kv_pool is not None: # We should also transfer draft model kv cache. The indices are diff --git a/python/sglang/srt/disaggregation/utils.py b/python/sglang/srt/disaggregation/utils.py index c9d8fb65e..84b58da33 100644 --- a/python/sglang/srt/disaggregation/utils.py +++ b/python/sglang/srt/disaggregation/utils.py @@ -1466,7 +1466,10 @@ def setup_state_kv_args( kv_args.kv_buf_groups = ( len(kv_args.kv_data_ptrs) // token_to_kv_pool.layer_num ) - kv_args.total_kv_layers = total_kv_layers + kv_args.hidden_kv_layers = total_kv_layers + kv_args.draft_kv_layers = ( + draft_token_to_kv_pool.layer_num if draft_token_to_kv_pool else 0 + ) else: append_state_component( kv_args, StateType.DSA, data_ptrs, data_lens, item_lens diff --git a/python/sglang/srt/hardware_backend/npu/modules/deepseek_v2_attention_mla_npu.py b/python/sglang/srt/hardware_backend/npu/modules/deepseek_v2_attention_mla_npu.py index da3f5643c..19b6a2690 100644 --- a/python/sglang/srt/hardware_backend/npu/modules/deepseek_v2_attention_mla_npu.py +++ b/python/sglang/srt/hardware_backend/npu/modules/deepseek_v2_attention_mla_npu.py @@ -447,7 +447,7 @@ def forward_dsa_prepare_npu( q_nope_out = q_nope_out.transpose(0, 1) - if m.layer_id == 0: + if m.layer_id == get_token_to_kv_pool().start_layer: m.rotary_emb.sin_cos_cache = m.rotary_emb.cos_sin_cache.index_select( 0, positions ) diff --git a/python/sglang/srt/layers/attention/dsa/dsa_npu_indexer.py b/python/sglang/srt/layers/attention/dsa/dsa_npu_indexer.py index 9329dad7a..1eb8ad655 100644 --- a/python/sglang/srt/layers/attention/dsa/dsa_npu_indexer.py +++ b/python/sglang/srt/layers/attention/dsa/dsa_npu_indexer.py @@ -159,7 +159,7 @@ class DSANPUIndexerMixin: k_pe = k_pe.unsqueeze(1) - if layer_id == 0: + if layer_id == get_token_to_kv_pool().start_layer: self.rotary_emb.sin_cos_cache = ( self.rotary_emb.cos_sin_cache.index_select(0, positions) ) diff --git a/python/sglang/srt/managers/overlap_utils.py b/python/sglang/srt/managers/overlap_utils.py index 6996b8e77..e43920cd4 100644 --- a/python/sglang/srt/managers/overlap_utils.py +++ b/python/sglang/srt/managers/overlap_utils.py @@ -608,7 +608,6 @@ class FutureMap: self.output_tokens_buf[indices] = payload.bonus_tokens.to( self.output_tokens_buf.dtype ) - if self.need_topk: self.topk_p_buf[indices] = payload.topk_p.to(self.topk_p_buf.dtype) self.topk_index_buf[indices] = payload.topk_index.to( diff --git a/python/sglang/srt/model_executor/model_runner_components/layer_setup.py b/python/sglang/srt/model_executor/model_runner_components/layer_setup.py index d421be458..c61b66133 100644 --- a/python/sglang/srt/model_executor/model_runner_components/layer_setup.py +++ b/python/sglang/srt/model_executor/model_runner_components/layer_setup.py @@ -5,6 +5,8 @@ from typing import TYPE_CHECKING, Any, NamedTuple import msgspec from torch import nn +from sglang.srt.utils import is_npu + if TYPE_CHECKING: from sglang.srt.configs.model_config import ModelConfig from sglang.srt.speculative.spec_info import SpeculativeAlgorithm @@ -156,12 +158,13 @@ def resolve_layer_indices( if loop_num > 1: num_effective_layers = num_effective_layers * loop_num - _assert_pp_mtp_compat( - model_has_mtp_layers=model_has_mtp_layers, - spec_algorithm=spec_algorithm, - num_effective_layers=num_effective_layers, - model_num_layers=model_num_layers, - ) + if not is_npu(): + _assert_pp_mtp_compat( + model_has_mtp_layers=model_has_mtp_layers, + spec_algorithm=spec_algorithm, + num_effective_layers=num_effective_layers, + model_num_layers=model_num_layers, + ) return ModelLayerInfo( start_layer=pp_range.start_layer, diff --git a/python/sglang/srt/models/deepseek_nextn.py b/python/sglang/srt/models/deepseek_nextn.py index a95dda44b..7d61f6f7c 100644 --- a/python/sglang/srt/models/deepseek_nextn.py +++ b/python/sglang/srt/models/deepseek_nextn.py @@ -39,7 +39,7 @@ from sglang.srt.layers.vocab_parallel_embedding import ( VocabParallelEmbedding, get_embedding_tp_kwargs, ) -from sglang.srt.model_executor.forward_batch_info import ForwardBatch +from sglang.srt.model_executor.forward_batch_info import ForwardBatch, PPProxyTensors from sglang.srt.models.deepseek_common.utils import enable_nextn_moe_bf16_cast_to_fp8 from sglang.srt.models.deepseek_v2 import DeepseekV2DecoderLayer, DeepseekV3ForCausalLM from sglang.srt.models.utils import WeightsMapper @@ -303,6 +303,7 @@ class DeepseekV3ForCausalLMNextN(DeepseekV3ForCausalLM): input_ids: torch.Tensor, positions: torch.Tensor, forward_batch: ForwardBatch, + pp_proxy_tensors: Optional[PPProxyTensors] = None, ) -> torch.Tensor: hidden_states = self.model(input_ids, positions, forward_batch) return self.logits_processor( diff --git a/python/sglang/srt/models/deepseek_v2.py b/python/sglang/srt/models/deepseek_v2.py index 31ee7d675..3d9438c91 100644 --- a/python/sglang/srt/models/deepseek_v2.py +++ b/python/sglang/srt/models/deepseek_v2.py @@ -2588,7 +2588,7 @@ class DeepseekV2Model(nn.Module): self.first_k_dense_replace = config.first_k_dense_replace self.pp_group = get_pp_group() - if self.pp_group.is_first_rank: + if self.pp_group.is_first_rank or (_is_npu and self.pp_group.is_last_rank): self.embed_tokens = VocabParallelEmbedding( config.vocab_size, config.hidden_size, diff --git a/python/sglang/srt/speculative/eagle_worker_v2.py b/python/sglang/srt/speculative/eagle_worker_v2.py index b0c79668a..192d57d82 100644 --- a/python/sglang/srt/speculative/eagle_worker_v2.py +++ b/python/sglang/srt/speculative/eagle_worker_v2.py @@ -41,7 +41,11 @@ from sglang.srt.model_executor.cuda_graph_config import ( Phase, check_cuda_graph_backend, ) -from sglang.srt.model_executor.forward_batch_info import CaptureHiddenMode, ForwardBatch +from sglang.srt.model_executor.forward_batch_info import ( + CaptureHiddenMode, + ForwardBatch, + PPProxyTensors, +) from sglang.srt.model_executor.forward_context import ForwardContext, forward_context from sglang.srt.model_executor.runner import ( DecodeCudaGraphRunner, @@ -1190,7 +1194,7 @@ class EAGLEWorkerV2(BaseSpecWorker): batch: ScheduleBatch, on_publish=None, grammar_barrier=None, - pp_proxy_tensors=None, + pp_proxy_tensors: Optional[PPProxyTensors] = None, ): if batch.forward_mode.is_extend() or batch.is_extend_in_batch: # Target prefill