Skip to content
Draft
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
41 changes: 1 addition & 40 deletions src/mcore_bridge/config/model_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -269,47 +269,8 @@ class ModelConfig(TransformerConfig):
num_labels: Optional[int] = None
mlp_padding_free: bool = False

_mindspeed_defaults_cache = None

def _augment_mindspeed_defaults(self):
if not is_torch_npu_available():
return

if ModelConfig._mindspeed_defaults_cache is None:
defaults = {}
try:
import mindspeed.features_manager as mfm
import sys
from argparse import ArgumentParser
from mindspeed.arguments import process_args

original_features = list(mfm.FEATURES_LIST)
full_features = mfm.create_features_list()
mfm.FEATURES_LIST.clear()
mfm.FEATURES_LIST.extend(full_features)
try:
parser = ArgumentParser()
process_args(parser)
# Parse args from sys.argv
args, _ = parser.parse_known_args([])
defaults = vars(args)
finally:
mfm.FEATURES_LIST.clear()
mfm.FEATURES_LIST.extend(original_features)
except Exception as e:
logger.warning(f'Failed to get MindSpeed defaults, which may cause issues on NPU: {e}')
defaults = {}
ModelConfig._mindspeed_defaults_cache = defaults

for name, value in ModelConfig._mindspeed_defaults_cache.items():
if not hasattr(self, name):
setattr(self, name, value)
elif hasattr(self, name) and getattr(self, name) is None and value is not None:
setattr(self, name, value)

def __post_init__(self):
from mcore_bridge.model import get_mcore_model_type, get_model_meta
self._augment_mindspeed_defaults()
self._format_config()
if self.experimental_attention_variant is not None:
require_version('megatron-core>=0.16.0.dev',
Expand Down Expand Up @@ -402,7 +363,7 @@ def _check_npu(self):
required_ep = (num_experts + MAX_NPU_EXPERTS_PER_EP - 1) // MAX_NPU_EXPERTS_PER_EP
if expert_model_parallel_size < required_ep:
logger.warning(f'{">" * 20} WARNING {"<" * 20}\n'
f'MindSpeed on NPU supports up to {MAX_NPU_EXPERTS_PER_EP} experts per EP group. '
f'NPU grouped matmul supports up to {MAX_NPU_EXPERTS_PER_EP} experts per EP group. '
f'num_experts={num_experts}, '
f'expert_model_parallel_size={expert_model_parallel_size}. '
f'Please set expert_model_parallel_size (EP) to {required_ep} '
Expand Down
94 changes: 23 additions & 71 deletions src/mcore_bridge/tuners/lora.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
import torch.nn.functional as F
import warnings
from contextlib import contextmanager, nullcontext
from importlib import metadata
from megatron.core import parallel_state
from megatron.core.dist_checkpointing.mapping import ShardedStateDict
from megatron.core.extensions.transformer_engine import (TEColumnParallelGroupedLinear, TEColumnParallelLinear,
Expand All @@ -33,36 +32,9 @@

mcore_016 = version.parse(megatron.core.__version__) >= version.parse('0.16.0rc0')
peft_019 = version.parse(peft.__version__) >= version.parse('0.19.0')
MINDSPEED_015 = version.parse('0.15.0')


def _get_mindspeed_version():
try:
return version.parse(metadata.version('mindspeed'))
except metadata.PackageNotFoundError:
return None
except Exception:
return None


def _use_legacy_npu_local_linear() -> bool:
if not is_torch_npu_available():
return False
mindspeed_version = _get_mindspeed_version()
if mindspeed_version is None:
# Fall back to the conservative path when the version is unknown so we
# do not force an older NPU stack onto the 0.15 TE semantics.
return True
return mindspeed_version < MINDSPEED_015


def _build_local_te_linear(input_size: int, output_size: int, bias: bool, **kwargs):
if _use_legacy_npu_local_linear():
return nn.Linear(
in_features=input_size,
out_features=output_size,
bias=bias,
)
local_kwargs = dict(kwargs)
local_kwargs.pop('tp_group', None)
return TELinear(
Expand All @@ -76,19 +48,27 @@ def _build_local_te_linear(input_size: int, output_size: int, bias: bool, **kwar


def _get_tensor_parallel_group_for_lora(base_layer):
"""Resolve the tensor-parallel group across TE and MindSpeed TE variants.

Megatron's TE layers expose ``tp_group`` directly, but MindSpeed 0.15.x
replaces some TE classes (for example
``MindSpeedTELayerNormColumnParallelLinear``) with implementations that keep
the same tensor-parallel semantics under ``parallel_group`` instead. LoRA
still needs to forward the right group into the newly created parallel
adapter layers, otherwise adapter injection fails before training starts.
"""
"""Resolve the tensor-parallel group from the standard TE/MCore contract."""
tp_group = getattr(base_layer, 'tp_group', None)
if tp_group is not None:
return tp_group
return getattr(base_layer, 'parallel_group', None)
return getattr(base_layer, '_tp_group', None)


def _forward_npu_layernorm_column(base_layer, x, args, kwargs):
"""Run TENPU LayerNormLinear once and retain its normalized activation."""
original_return_layernorm_output = base_layer.return_layernorm_output
base_layer.return_layernorm_output = True
try:
out = base_layer(x, *args, **kwargs)
finally:
base_layer.return_layernorm_output = original_return_layernorm_output

if base_layer.te_return_bias:
result, bias, layernorm_output = out
else:
(result, layernorm_output), bias = out
return result, bias, layernorm_output


class LoraParallelLinear(MegatronModule, LoraLayer):
Expand Down Expand Up @@ -224,10 +204,8 @@ def update_layer(self, adapter_name, r, *, lora_alpha, **kwargs):
lora_b = _build_local_te_linear(r, self.out_features, lora_bias, **kwargs)
lora_a.parallel_mode = self.base_layer.parallel_mode # fix moe_shared_expert_overlap
else:
if is_torch_npu_available():
out_features = self.out_features
else:
out_features = self.out_features * self.tp_size
# PEFT reports local features; MCore's constructor expects the global size.
out_features = self.out_features * self.tp_size
if self.is_grouped:
if is_torch_npu_available():
lora_a = NpuGroupedLoraLinear(
Expand Down Expand Up @@ -380,37 +358,11 @@ def forward(self, x: torch.Tensor, *args: Any, **kwargs: Any):
self.base_layer.return_layernorm_output = False
result, bias = self.base_layer(x, *args, **kwargs)
else:
self.base_layer.return_layernorm_output = True
if is_torch_npu_available():
# NPU: base_layer only returns (output, bias); it does not expose the LayerNorm/RMSNorm output.
inp = x # Keep the original pre-norm input.
result, bias = self.base_layer(inp, *args, **kwargs)

# Key: For LoRA we need the same "x" as in the non-NPU branch, i.e. the post-norm activation
# (LayerNorm/RMSNorm output, which is the actual input to the fused linear).
if hasattr(self.base_layer, 'config') and (hasattr(self.base_layer, '_layernorm')
or hasattr(self.base_layer, '_rmsnorm')):
norm_type = getattr(self.base_layer.config, 'normalization', None)

if norm_type == 'LayerNorm':
if not hasattr(self.base_layer, '_layernorm'):
raise RuntimeError(
'NPU LoRA path expects base_layer to provide `_layernorm`, but it is missing. '
'Cannot reconstruct the post-LayerNorm activation for LoRA.')
x = self.base_layer._layernorm(inp)
else:
# Default to RMSNorm path when normalization is not LayerNorm.
if not hasattr(self.base_layer, '_rmsnorm'):
raise RuntimeError(
'NPU LoRA path expects base_layer to provide `_rmsnorm`, but it is missing. '
'Cannot reconstruct the post-RMSNorm activation for LoRA.')
x = self.base_layer._rmsnorm(inp)
else:
raise RuntimeError('NPU LoRA path requires base_layer to expose post-norm activations '
'(LayerNorm/RMSNorm output). Expected base_layer to have `config` '
'and either `_layernorm` or `_rmsnorm`. '
f'Got base_layer type: {type(self.base_layer)}. ')
result, bias, x = _forward_npu_layernorm_column(
self.base_layer, x, args, kwargs)
else:
self.base_layer.return_layernorm_output = True
(result, x), bias = self.base_layer(x, *args, **kwargs)
elif isinstance(self.base_layer, (TELinear, TEGroupedLinear)):
result, bias = self.base_layer(x, *args, **kwargs)
Expand Down
18 changes: 7 additions & 11 deletions src/mcore_bridge/tuners/npu_lora.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,6 @@
_GMM_GROUP_LIST_IS_EXPERT_SIZES = 1


def _is_mindspeed_grouped_linear(base_layer) -> bool:
if not (is_torch_npu_available() and isinstance(base_layer, TEGroupedLinear)):
return False
return type(base_layer).__module__.startswith('mindspeed.')


def _has_moe_local_expert_grouping(base_layer) -> bool:
config = getattr(base_layer, 'config', None)
num_moe_experts = getattr(config, 'num_moe_experts', None)
Expand All @@ -41,16 +35,16 @@ def is_expert_layer(base_layer) -> bool:
is_expert = getattr(base_layer, 'is_expert', None)
if is_expert is not None:
return bool(is_expert)
if _is_mindspeed_grouped_linear(base_layer):
if is_torch_npu_available() and isinstance(base_layer, TEGroupedLinear):
if getattr(base_layer, 'explicit_expert_comm', False):
return True
has_moe_local_expert_grouping = _has_moe_local_expert_grouping(base_layer)
if getattr(base_layer, 'expert_parallel', False):
return has_moe_local_expert_grouping
# MindSpeedTEGroupedLinear receives is_expert but does not keep it as
# an attribute. When EP/ETP does not trigger explicit expert comm,
# fall back to the TEGroupedMLP invariant: one grouped slot per local
# expert.
# MCore's TEGroupedLinear requires is_expert=True at construction but
# does not retain that argument. With EP=ETP=1 there is no explicit
# expert communication, so use the TEGroupedMLP invariant: one grouped
# GEMM slot per local expert.
return has_moe_local_expert_grouping
return False

Expand Down Expand Up @@ -137,6 +131,8 @@ def sharded_state_dict(
param,
f'{key_prefix}{param_name}',
prepend_offsets=new_sharded_offsets,
tp_group=parallel_state.get_expert_tensor_parallel_group(),
dp_cp_group=parallel_state.get_expert_data_parallel_group(),
)
sharded_state_dict[f'{prefix}{local_name}'] = self._set_expert_replica_id(sharded_tensor)

Expand Down
Loading