Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions lmdeploy/turbomind/model_loader.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,13 +50,23 @@ def _bind_runtime(self):
model_tp = ParallelGroup(ec.attn_tp_size * ec.attn_cp_size,
[mc.model_tp_rank(g) for g in range(self.gpu_count)])

# Dense (non-expert) FFN TP: node-local — one node's ranks within
# the comm domain (shard index = inner_rank % domain_size,
# inner_rank = ep_rank * mlp_tp_size + mlp_tp_rank).
dense_size = min(mlp_tp.size * ep.size, self.gpu_count)
dense_tp = ParallelGroup(
dense_size,
[(e * mlp_tp.size + m) % dense_size
for e, m in zip(ep.ranks, mlp_tp.ranks)])

self.model.bind_runtime(
ctx=ctx,
root_handles=[mc.root(g) for g in range(self.gpu_count)],
attn_tp=attn_tp,
mlp_tp=mlp_tp,
ep=ep,
model_tp=model_tp,
dense_tp=dense_tp,
)

def export(self):
Expand Down
19 changes: 11 additions & 8 deletions lmdeploy/turbomind/models/glm4_moe_lite.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,14 +100,14 @@ def attn(self, pfx):
# FFN / MoE factories
# ------------------------------------------------------------------

def ffn(self, pfx, inter_size, is_expert=False):
def ffn(self, pfx, inter_size, is_expert=False, *, tp):
w1, w3, w2 = [self._linear(pfx + f'{x}_proj') for x in ('gate', 'up', 'down')]

cfg = self._ffn_cfg.clone()
cfg.inter_size = inter_size
cfg.is_expert = is_expert

m = FfnBuilder(cfg, self._ctx, tp=self._mlp_tp)
m = FfnBuilder(cfg, self._ctx, tp=tp)
m.add_ffn(w1, w2, w3)
return m.build()

Expand All @@ -124,13 +124,15 @@ def moe(self, pfx):
experts = ModuleListBuilder(ModuleListConfig(), self._ctx)
for e in m.range(cfg.expert_num):
experts[e] = self.ffn(pfx + 'experts' + e,
self.cfg.moe_intermediate_size, is_expert=True)
self.cfg.moe_intermediate_size,
is_expert=True, tp=self._mlp_tp)
m.experts = experts.build()

shared = self.ffn(pfx + 'shared_experts',
self.cfg.intermediate_size * self.cfg.n_shared_experts)
m.shared = self.ffn(pfx + 'shared_experts',
self.cfg.intermediate_size * self.cfg.n_shared_experts,
tp=self._mlp_tp)

return m.build(), shared
return m.build()

def layers(self, pfx):
layers = ModuleListBuilder(ModuleListConfig(), self._ctx)
Expand All @@ -140,8 +142,9 @@ def layers(self, pfx):
d.attention = self.attn(p + 'self_attn')
d.ffn_norm = self.norm(p + 'post_attention_layernorm')
if self.cfg.mlp_layer_types[i] == 'sparse':
d.moe_ffn, d.feed_forward = self.moe(p + 'mlp')
d.moe_ffn = self.moe(p + 'mlp')
else:
d.feed_forward = self.ffn(p + 'mlp', self.cfg.intermediate_size)
d.feed_forward = self.ffn(p + 'mlp', self.cfg.intermediate_size,
tp=self._dense_tp)
layers[i] = d.build()
return layers.build()
13 changes: 8 additions & 5 deletions lmdeploy/turbomind/models/interns2_mobius.py
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,12 @@ def build_meta_moe_layer(text_model, pfx, layer_id: int):

m = MoeBuilder(cfg, text_model._ctx, ep=text_model._ep)
m.add_gate('shared_gate', text_model._linear(pfx + 'shared_expert_gate'))
# Attach the shared expert before build(): post-build child assignment
# raises, and the C++ MoeWeight owns it as the `shared` child.
m.shared = text_model.ffn(
pfx + 'shared_expert',
text_model.cfg.shared_expert_intermediate_size,
tp=text_model._mlp_tp)
moe = m.build()

pack = text_model._meta_pack_modules[layer_id % text_model._n_meta_groups]
Expand All @@ -126,10 +132,7 @@ def build_meta_moe_layer(text_model, pfx, layer_id: int):
with text_model._ctx.devices[i]:
moe_h.set_meta_pack(pack_h)

shared = text_model.ffn(
pfx + 'shared_expert',
text_model.cfg.shared_expert_intermediate_size)
return moe, shared
return moe


class InternS2MobiusTextModel(Qwen3_5TextModel):
Expand Down Expand Up @@ -178,7 +181,7 @@ def layers(self, pfx):
d.linear_attn = self.linear_attn(p + 'linear_attn')
else:
d.attention = self.attn(p + 'self_attn')
d.moe_ffn, d.feed_forward = build_meta_moe_layer(self, p + 'mlp', i)
d.moe_ffn = build_meta_moe_layer(self, p + 'mlp', i)
d.attention_norm = self.norm(p + 'input_layernorm', zero_centered=True)
d.ffn_norm = self.norm(p + 'post_attention_layernorm', zero_centered=True)
layers[i] = d.build()
Expand Down
3 changes: 2 additions & 1 deletion lmdeploy/turbomind/models/internvl.py
Original file line number Diff line number Diff line change
Expand Up @@ -433,14 +433,15 @@ def __init__(self, cfg: PretrainedConfig, *, resolver,
raise ValueError(f'InternVL TurboMind vision architecture {arch!r} is not supported.')

def bind_runtime(self, *, ctx, root_handles,
attn_tp, mlp_tp, ep, model_tp):
attn_tp, mlp_tp, ep, model_tp, dense_tp):
self.text_model.bind_runtime(
ctx=ctx,
root_handles=root_handles,
attn_tp=attn_tp,
mlp_tp=mlp_tp,
ep=ep,
model_tp=model_tp,
dense_tp=dense_tp,
)
if self.vision_model is not None:
vision_ctx = Context(
Expand Down
18 changes: 10 additions & 8 deletions lmdeploy/turbomind/models/qwen2.py
Original file line number Diff line number Diff line change
Expand Up @@ -96,14 +96,14 @@ def reorder(x):

return m.build()

def ffn(self, pfx, inter_size, is_expert=False):
def ffn(self, pfx, inter_size, is_expert=False, *, tp):
w1, w3, w2 = [self._linear(pfx + f'{x}_proj') for x in ('gate', 'up', 'down')]

cfg = self._ffn_cfg.clone()
cfg.inter_size = inter_size
cfg.is_expert = is_expert

m = FfnBuilder(cfg, self._ctx, tp=self._mlp_tp)
m = FfnBuilder(cfg, self._ctx, tp=tp)
m.add_ffn(w1, w2, w3)
return m.build()

Expand All @@ -118,14 +118,15 @@ def moe(self, pfx):
for e in m.range(self.cfg.num_experts):
experts[e] = self.ffn(pfx + 'experts' + e,
self.cfg.moe_intermediate_size,
is_expert=True)
is_expert=True, tp=self._mlp_tp)
m.experts = experts.build()

m.add_gate('shared_gate', self._linear(pfx + 'shared_expert_gate'))
shared = self.ffn(pfx + 'shared_expert',
self.cfg.shared_expert_intermediate_size)
m.shared = self.ffn(pfx + 'shared_expert',
self.cfg.shared_expert_intermediate_size,
tp=self._mlp_tp)

return m.build(), shared
return m.build()

# ------------------------------------------------------------------
# layers() — layer dispatch loop
Expand All @@ -137,9 +138,10 @@ def layers(self, pfx):
d = DecoderLayerBuilder(DecoderLayerConfig(), self._ctx)
d.attention = self.attn(p + 'self_attn')
if self._n_experts > 0:
d.moe_ffn, d.feed_forward = self.moe(p + 'mlp')
d.moe_ffn = self.moe(p + 'mlp')
else:
d.feed_forward = self.ffn(p + 'mlp', self.cfg.intermediate_size)
d.feed_forward = self.ffn(p + 'mlp', self.cfg.intermediate_size,
tp=self._mlp_tp)
d.attention_norm = self.norm(p + 'input_layernorm')
d.ffn_norm = self.norm(p + 'post_attention_layernorm')
layers[i] = d.build()
Expand Down
3 changes: 2 additions & 1 deletion lmdeploy/turbomind/models/qwen2_vl.py
Original file line number Diff line number Diff line change
Expand Up @@ -413,14 +413,15 @@ def __init__(self, cfg, *, resolver, vision_resolver=None,
self._vision_data_type = vision_data_type

def bind_runtime(self, *, ctx, root_handles,
attn_tp, mlp_tp, ep, model_tp):
attn_tp, mlp_tp, ep, model_tp, dense_tp):
self.text_model.bind_runtime(
ctx=ctx,
root_handles=root_handles,
attn_tp=attn_tp,
mlp_tp=mlp_tp,
ep=ep,
model_tp=model_tp,
dense_tp=dense_tp,
)
if self.vision_model is not None:
vision_ctx = Context(
Expand Down
19 changes: 14 additions & 5 deletions lmdeploy/turbomind/models/qwen3.py
Original file line number Diff line number Diff line change
Expand Up @@ -97,13 +97,13 @@ def reorder(x):

return m.build()

def ffn(self, pfx, is_expert=False):
def ffn(self, pfx, is_expert=False, *, tp):
w1, w3, w2 = [self._linear(pfx + f'{x}_proj') for x in ('gate', 'up', 'down')]

cfg = self._ffn_cfg.clone()
cfg.is_expert = is_expert

m = FfnBuilder(cfg, self._ctx, tp=self._mlp_tp)
m = FfnBuilder(cfg, self._ctx, tp=tp)
m.add_ffn(w1, w2, w3)
return m.build()

Expand All @@ -115,21 +115,30 @@ def moe(self, pfx):

experts = ModuleListBuilder(ModuleListConfig(), self._ctx)
for e in m.range(self.cfg.num_experts):
experts[e] = self.ffn(pfx + 'experts' + e, is_expert=True)
experts[e] = self.ffn(pfx + 'experts' + e, is_expert=True,
tp=self._mlp_tp)
m.experts = experts.build()

return m.build()

def layers(self, pfx):
mlp_only = set(getattr(self.cfg, 'mlp_only_layers', None) or [])
# Dense layers in a MoE model shard node-locally (_dense_tp); a
# pure-dense model keeps the historical mlp_tp sharding. Keyed on
# whether any layer actually builds MoE — the same fact the C++
# decoder keys the dense reduce group on.
any_moe = self._n_experts > 0 and any(
i not in mlp_only for i in range(self.cfg.num_hidden_layers))
dense_tp = self._dense_tp if any_moe else self._mlp_tp
layers = ModuleListBuilder(ModuleListConfig(), self._ctx)
for i, p in pfx.slices(0, self.cfg.num_hidden_layers):
d = DecoderLayerBuilder(DecoderLayerConfig(), self._ctx)
d.attention_norm = self.norm(p + 'input_layernorm')
d.attention = self.attn(p + 'self_attn')
d.ffn_norm = self.norm(p + 'post_attention_layernorm')
if self._n_experts:
if self._n_experts and i not in mlp_only:
Comment thread
lzhangzz marked this conversation as resolved.
d.moe_ffn = self.moe(p + 'mlp')
else:
d.feed_forward = self.ffn(p + 'mlp')
d.feed_forward = self.ffn(p + 'mlp', tp=dense_tp)
layers[i] = d.build()
return layers.build()
23 changes: 15 additions & 8 deletions lmdeploy/turbomind/models/qwen3_5.py
Original file line number Diff line number Diff line change
Expand Up @@ -204,7 +204,7 @@ def linear_attn(self, pfx):
# FFN / MoE factories
# ------------------------------------------------------------------

def ffn(self, pfx, inter_size, is_expert=False):
def ffn(self, pfx, inter_size, is_expert=False, *, tp):
try:
w1, w3, w2 = [self._linear(pfx + f'{x}_proj')
for x in ('gate', 'up', 'down')]
Expand All @@ -215,7 +215,7 @@ def ffn(self, pfx, inter_size, is_expert=False):
cfg.inter_size = inter_size
cfg.is_expert = is_expert

m = FfnBuilder(cfg, self._ctx, tp=self._mlp_tp)
m = FfnBuilder(cfg, self._ctx, tp=tp)
m.add_ffn(w1, w2, w3)
return m.build()

Expand All @@ -234,9 +234,13 @@ def moe(self, pfx):
m.experts = experts.build()

m.add_gate('shared_gate', self._linear(pfx + 'shared_expert_gate'))
shared = self.ffn(pfx + 'shared_expert', self.cfg.shared_expert_intermediate_size)
shared = self.ffn(pfx + 'shared_expert',
self.cfg.shared_expert_intermediate_size,
tp=self._mlp_tp)
if shared is not None:
m.shared = shared

return m.build(), shared
return m.build()

def _packed_moe_ffn(self, experts_pfx, expert_idx, inter_size):
w1, w2, w3 = read_packed_moe_expert(
Expand All @@ -254,7 +258,8 @@ def _packed_moe_ffn(self, experts_pfx, expert_idx, inter_size):

def _moe_expert_ffn(self, experts_pfx, expert_idx, inter_size):
expert_pfx = experts_pfx + expert_idx
return (self.ffn(expert_pfx, inter_size, is_expert=True)
return (self.ffn(expert_pfx, inter_size, is_expert=True,
tp=self._mlp_tp)
or self._packed_moe_ffn(experts_pfx, expert_idx, inter_size))

# ------------------------------------------------------------------
Expand All @@ -270,9 +275,10 @@ def layers(self, pfx):
else:
d.attention = self.attn(p + 'self_attn')
if self._n_experts > 0:
d.moe_ffn, d.feed_forward = self.moe(p + 'mlp')
d.moe_ffn = self.moe(p + 'mlp')
else:
d.feed_forward = self.ffn(p + 'mlp', self.cfg.intermediate_size)
d.feed_forward = self.ffn(p + 'mlp', self.cfg.intermediate_size,
tp=self._mlp_tp)
d.attention_norm = self.norm(
p + 'input_layernorm',
zero_centered=True,
Expand Down Expand Up @@ -602,14 +608,15 @@ def __init__(self, cfg: Qwen3_5Config | Qwen3_5MoeConfig, *, resolver,
self._vision_data_type = vision_data_type

def bind_runtime(self, *, ctx, root_handles,
attn_tp, mlp_tp, ep, model_tp):
attn_tp, mlp_tp, ep, model_tp, dense_tp):
self.text_model.bind_runtime(
ctx=ctx,
root_handles=root_handles,
attn_tp=attn_tp,
mlp_tp=mlp_tp,
ep=ep,
model_tp=model_tp,
dense_tp=dense_tp,
)

if self.vision_model is not None:
Expand Down
6 changes: 5 additions & 1 deletion lmdeploy/turbomind/text_model.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,13 +41,17 @@ def _vocab_size(self) -> int:
return self.cfg.vocab_size

def bind_runtime(self, *, ctx, root_handles,
attn_tp, mlp_tp, ep, model_tp):
attn_tp, mlp_tp, ep, model_tp, dense_tp):
self._ctx = ctx
self._root_handles = root_handles
self._attn_tp = attn_tp
self._mlp_tp = mlp_tp
self._ep = ep
self._model_tp = model_tp
# TP group for dense layers inside MoE models — sharded node-locally
# by ModelLoader (C++ reduce group: d_node_group). Pure-dense models
# use the mlp_tp group.
self._dense_tp = dense_tp

def _linear(self, pfx: Prefix, *,
optional: bool = False) -> Linear | None:
Expand Down
24 changes: 15 additions & 9 deletions src/turbomind/comm/nccl/deepep/moe_a2a_utils.cu
Original file line number Diff line number Diff line change
Expand Up @@ -184,16 +184,23 @@ void invokeMoeA2AMapping(int* f2n,
template<class T, int vec_size>
__global__ void MoeA2ASharedCombineKernel(T* output, //
const T* routed,
const T* shared, // this rank's shared FFN rows, or nullptr
const float* shared_scales,
int hidden_dim,
float shared_scale)
int hidden_dim)
{
const int token_idx = blockIdx.x;

output += (int64_t)token_idx * hidden_dim;
routed += (int64_t)token_idx * hidden_dim;

if (shared_scales) {
// Nullability lives in the kernel, not the caller: this kernel is also
// the routed write-back (the Store below is unconditional), so the
// caller must stay unconditional and pass nullptr when there is no
// shared expert.
const T* shared_row = shared ? shared + (int64_t)token_idx * hidden_dim : nullptr;

float shared_scale = 1.f;
if (shared_scales) { // null for gate-less shared experts
shared_scale *= fdividef(1.f, 1.f + expf(-__ldg(shared_scales + token_idx)));
}

Expand All @@ -203,21 +210,20 @@ __global__ void MoeA2ASharedCombineKernel(T* output, //
Vec routed_vec;
Load(routed_vec, routed + i);
auto result = cast<float>(routed_vec);

if (shared_scale != 0.f) {
if (shared_row) {
Vec shared_vec;
Load(shared_vec, output + i);
Load(shared_vec, shared_row + i); // was: in-place read from output
using namespace ops;
result = result + cast<float>(shared_vec) * shared_scale;
}
Store(output + i, cast<T>(result));
Store(output + i, cast<T>(result)); // ALWAYS executed — the routed write-back
}
}

void invokeMoeA2ASharedCombine(core::Tensor& output,
const core::Tensor& routed,
const core::Tensor& shared,
const float* shared_scales,
float shared_scale,
cudaStream_t stream)
{
const int tokens = output.shape(0);
Expand All @@ -231,7 +237,7 @@ void invokeMoeA2ASharedCombine(core::Tensor& output,
constexpr int vec_size = 16 / sizeof(T);
constexpr int block_dim = 256;
MoeA2ASharedCombineKernel<T, vec_size><<<tokens, block_dim, 0, stream>>>(
output.data<T>(), routed.data<T>(), shared_scales, hidden_dim, shared_scale);
output.data<T>(), routed.data<T>(), shared.data_or((T*)nullptr), shared_scales, hidden_dim);
TM_CUDA_CHECK(cudaGetLastError());
};

Expand Down
Loading
Loading