submission 745202
Danishlynx · python · License unknown
Kernel source · 308 lines ↓holds 1 record
Use it
Vendorable · source mirrored · license unknownView source →
No package. Vendor the mirrored source: 308 lines, June 9 Researcher Reciprocity License v1.0.
submission.py
curl "https://kernelindex.com/api/v1/implementations/kernelbot-amd-moe-mxfp4-745202?include=source"interfacepython
Compatibility
measured onAMD Instinct MI355X
declared hardwareAMD Instinct MI355X
architecturesgfx950
dtypesbf16, fp32, fp8_e8m0, int32, mxfp4
Benchmark evidence
1 measurement across 1 GPU, fastest first.
Operation / workload
Hardware
Latency
Rank
Observed
Reported · How evidence levels are derived →
Source and license
sourceavailable
revision digestsha256:da2cce193fd4b8eb51602768f2903828bdced958b1d6d6d6de159bb55f7a11cc
license declaredunknown
license concludedunknown
authorsDanishlynx
imported2026-08-15
Techniques
Extracted from the mirrored source by pattern, never inferred. Each row cites its line.
Kernel source
submission.py308 lines
# /// script
# leaderboard = "amd-moe-mxfp4"
# ///
"""v917: v913 fixed to survive AITER repo updates.
Key fix: clone with depth=50 and find a commit that has _get_compiled_stage1
(the old API). The repo was updated on ~April 4-5 and renamed the functions.
Falls back to new API with runtime patch if needed."""
import os, sys, torch, subprocess, shutil, re, functools
os.environ['HIP_FORCE_DEV_KERNARG'] = '1'
os.environ['AITER_ENABLE_VSKIP'] = '1'
os.environ['AMD_DIRECT_DISPATCH'] = '1'
os.environ['GPU_MAX_HW_QUEUES'] = '8'
os.environ['HSA_ENABLE_SDMA'] = '0'
os.environ['HIP_LAUNCH_BLOCKING'] = '0'
os.environ['AMD_SERIALIZE_KERNEL'] = '0'
os.environ['GPU_SINGLE_ALLOC_PERCENT'] = '100'
os.environ['HSA_SVM_GUARD_PAGES'] = '0'
_AITER_DIR = '/home/runner/aiter'
_NEW_DIR = '/tmp/aiter_fresh'
_FLYDSL_DIR = '/tmp/flydsl_new'
# ========== STEP 1: Install flydsl 0.1.1 ==========
os.makedirs(_FLYDSL_DIR, exist_ok=True)
subprocess.run([sys.executable, '-m', 'pip', 'install', 'flydsl==0.1.1',
'--target', _FLYDSL_DIR, '--upgrade', '--break-system-packages'],
capture_output=True, text=True, timeout=120)
sys.path.insert(0, _FLYDSL_DIR)
for key in list(sys.modules.keys()):
if 'flydsl' in key:
del sys.modules[key]
import flydsl
print(f"[v917] flydsl {flydsl.__version__}", file=sys.stderr)
# ========== STEP 2: Setup AITER — pin to working commit ==========
subprocess.run(['git', 'checkout', '--', '.'], capture_output=True, cwd=_AITER_DIR, timeout=10)
# Clone with enough depth to find old API
if not os.path.exists(os.path.join(_NEW_DIR, 'aiter')):
subprocess.run(['git', 'clone', '--depth=50', 'https://github.com/ROCm/aiter.git', _NEW_DIR],
capture_output=True, text=True, timeout=120)
# Find a commit where _get_compiled_stage1 exists in moe_kernels.py
_mk_path_new = os.path.join(_NEW_DIR, 'aiter', 'ops', 'flydsl', 'moe_kernels.py')
_found_old_api = False
# Check HEAD first
if os.path.exists(_mk_path_new):
with open(_mk_path_new) as f:
_mk_content = f.read()
if '_get_compiled_stage1' in _mk_content:
_found_old_api = True
print("[v917] HEAD has old API (_get_compiled_stage1)", file=sys.stderr)
if not _found_old_api:
# Search older commits
print("[v917] HEAD has new API, searching older commits...", file=sys.stderr)
r = subprocess.run(['git', 'log', '--oneline', '-50', '--', 'aiter/ops/flydsl/moe_kernels.py'],
capture_output=True, text=True, cwd=_NEW_DIR, timeout=10)
print(f"[v917] Commits touching moe_kernels.py:\n{r.stdout[:500]}", file=sys.stderr)
# Try each commit until we find one with old API
r = subprocess.run(['git', 'log', '--format=%H', '-50'],
capture_output=True, text=True, cwd=_NEW_DIR, timeout=10)
commits = r.stdout.strip().split('\n')
for commit in commits[1:]: # skip HEAD (already checked)
if not commit.strip():
continue
# Checkout and check
subprocess.run(['git', 'checkout', commit.strip()], capture_output=True, cwd=_NEW_DIR, timeout=10)
if os.path.exists(_mk_path_new):
with open(_mk_path_new) as f:
_c = f.read()
if '_get_compiled_stage1' in _c:
_found_old_api = True
print(f"[v917] Found old API in commit {commit.strip()[:12]}", file=sys.stderr)
break
if not _found_old_api:
print("[v917] WARNING: Could not find old API in any commit!", file=sys.stderr)
# Copy newer fused_moe.py (main dispatch file)
shutil.copy2(os.path.join(_NEW_DIR, 'aiter', 'fused_moe.py'),
os.path.join(_AITER_DIR, 'aiter', 'fused_moe.py'))
# Copy flydsl files
flydsl_src = os.path.join(_NEW_DIR, 'aiter', 'ops', 'flydsl')
flydsl_dst = os.path.join(_AITER_DIR, 'aiter', 'ops', 'flydsl')
for fname in ['moe_kernels.py', 'utils.py']:
src_p = os.path.join(flydsl_src, fname)
dst_p = os.path.join(flydsl_dst, fname)
if os.path.exists(src_p):
shutil.copy2(src_p, dst_p)
# Patch version check
init_path = os.path.join(flydsl_dst, '__init__.py')
with open(init_path) as f:
init_src = f.read()
init_src = re.sub(r'raise ImportError\([^)]*\)', 'pass # version check bypassed', init_src)
with open(init_path, 'w') as f:
f.write(init_src)
# Copy kernels
kernels_src = os.path.join(flydsl_src, 'kernels')
kernels_dst = os.path.join(flydsl_dst, 'kernels')
if os.path.exists(kernels_src):
os.makedirs(kernels_dst, exist_ok=True)
for fname in os.listdir(kernels_src):
if fname.endswith('.py'):
shutil.copy2(os.path.join(kernels_src, fname), os.path.join(kernels_dst, fname))
# ========== STEP 3: Source patches (same as v913) ==========
_FMOE_PATH = os.path.join(_AITER_DIR, 'aiter', 'fused_moe.py')
with open(_FMOE_PATH, 'r') as f:
src = f.read()
src = src.replace('split_k=ksplit', 'split_k=2')
src = src.replace('splitk=ksplit', 'splitk=2')
bm_marker = 'def get_block_size_M(token, topk, expert, inter_dim):'
if bm_marker in src:
dec_search = src.rfind('@functools.lru_cache', 0, src.index(bm_marker))
start = dec_search if dec_search >= 0 else src.index(bm_marker)
lines_bm = src[src.index(bm_marker):].split('\n')
end_offset = 0; found_body = False
for i, line in enumerate(lines_bm):
if i == 0: end_offset += len(line) + 1; continue
stripped = line.strip()
if stripped == '' or stripped.startswith('#'): end_offset += len(line) + 1; continue
if found_body and (not line.startswith(' ') and not line.startswith('\t') and stripped != ''): break
if line.startswith(' ') or line.startswith('\t'): found_body = True
end_offset += len(line) + 1
end = src.index(bm_marker) + end_offset
new_bm = '''@functools.lru_cache(maxsize=2048)
def get_block_size_M(token, topk, expert, inter_dim):
_ov = {(512,9,33,512): 128}
key = (token, topk, expert, inter_dim)
if key in _ov: return _ov[key]
cu_num = get_cu_num(); tileN = 128; tgN = (inter_dim + tileN - 1) // tileN
tmp = []
for el in [32, 64, 128]:
mnt = token * topk + expert * el - topk; tg = tgN * (mnt + el - 1) // el
tmp.append(((tg + cu_num - 1) // cu_num, cu_num - tg % cu_num, el))
return sorted(tmp, key=lambda x: x[:2])[0][-1]
'''
src = src[:start] + new_bm + src[end:]
# Append FlyDSL stage1 + stage2 wrappers (same as v913)
tail = r'''
# v917: FlyDSL BOTH stages (v913 logic, pinned AITER)
def _v917_fly_s1(hidden_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w1_scale=None, a1_scale=None, sorted_weights=None, **_kw):
from aiter.ops.flydsl.moe_kernels import flydsl_moe_stage1
flydsl_moe_stage1(a=hidden_states, w1=w1, sorted_token_ids=sorted_token_ids,
sorted_expert_ids=sorted_expert_ids, num_valid_ids=num_valid_ids, out=out, topk=topk,
tile_m=32, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",
act="silu", w1_scale=w1_scale, a1_scale=a1_scale, sorted_weights=sorted_weights)
return out
def _v917_fly_s2_direct(inter_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w2_scale=None, a2_scale=None, sorted_weights=None, **_kw):
from aiter.ops.flydsl.moe_kernels import flydsl_moe_stage2
flydsl_moe_stage2(inter_states=inter_states, w2=w2, sorted_token_ids=sorted_token_ids,
sorted_expert_ids=sorted_expert_ids, num_valid_ids=num_valid_ids, out=out, topk=topk,
tile_m=16, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",
mode="direct", w2_scale=w2_scale, a2_scale=a2_scale, sorted_weights=sorted_weights)
def _v917_fly_s2_atomic(inter_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w2_scale=None, a2_scale=None, sorted_weights=None, **_kw):
from aiter.ops.flydsl.moe_kernels import flydsl_moe_stage2
flydsl_moe_stage2(inter_states=inter_states, w2=w2, sorted_token_ids=sorted_token_ids,
sorted_expert_ids=sorted_expert_ids, num_valid_ids=num_valid_ids, out=out, topk=topk,
tile_m=16, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",
mode="atomic", w2_scale=w2_scale, a2_scale=a2_scale, sorted_weights=sorted_weights)
_v917_orig_g2sc = get_2stage_cfgs
try: _v917_orig_g2sc.cache_clear()
except: pass
_v917_fly_ready = [False]
_v917_s1_ready = [False]
@functools.lru_cache(maxsize=2048)
def _v917_g2sc(*a, **kw):
md = _v917_orig_g2sc(*a, **kw)
if len(a) >= 5:
E_val = a[4]; bs_val = a[1]
if E_val <= 100 and bs_val >= 512:
md.use_non_temporal_load = True
if _v917_s1_ready[0] and not md.run_1stage and md.stage1 is not None:
md.stage1 = functools.partial(_v917_fly_s1)
elif not md.run_1stage and md.stage1 is not None:
_orig_s1 = md.stage1
def _s1w(*aa, **kk):
kk['splitk'] = 2
return _orig_s1(*aa, **kk)
md.stage1 = _s1w
if _v917_fly_ready[0] and not md.run_1stage and md.stage2 is not None:
E_val = a[4] if len(a) >= 5 else 0
md.stage2 = functools.partial(_v917_fly_s2_atomic if E_val > 100 else _v917_fly_s2_direct)
return md
get_2stage_cfgs = _v917_g2sc
import sys as _s
print("[v917] FlyDSL BOTH stages patched", file=_s.stderr)
'''
src += tail
with open(_FMOE_PATH, 'w') as f:
f.write(src)
# ========== STEP 4: Import ==========
from aiter import ActivationType, QuantType
from aiter.fused_moe import fused_moe
import aiter.fused_moe as _fmoe
from task import input_t, output_t
# ========== STEP 5: Pre-compile FlyDSL ==========
# Try both old and new API names
import aiter.ops.flydsl.moe_kernels as _mk
_compile_s1 = getattr(_mk, '_get_compiled_stage1', None) or getattr(_mk, 'compile_flydsl_moe_stage1', None)
_compile_s2 = getattr(_mk, '_get_compiled_stage2', None) or getattr(_mk, 'compile_flydsl_moe_stage2', None)
print(f"[v917] S1 compile fn: {_compile_s1.__name__ if _compile_s1 else 'NONE'}", file=sys.stderr)
print(f"[v917] S2 compile fn: {_compile_s2.__name__ if _compile_s2 else 'NONE'}", file=sys.stderr)
if _compile_s1:
# Detect which API we're using by checking parameter names
import inspect
_s1_params = inspect.signature(_compile_s1).parameters
_use_old_api = 'doweight' in _s1_params # old API uses 'doweight', new uses 'doweight_stage1'
for inter, exp in [(256, 257), (512, 33), (2048, 33)]:
try:
if _use_old_api:
_compile_s1(model_dim=7168, inter_dim=inter, experts=exp, topk=9,
tile_m=32, tile_n=256, tile_k=128, doweight=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", act="silu")
else:
_compile_s1(model_dim=7168, inter_dim=inter, experts=exp, topk=9,
tile_m=32, tile_n=256, tile_k=256, doweight_stage1=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", act="silu")
print(f"[v917] Stage1 compiled: inter={inter} E={exp}", file=sys.stderr)
except Exception as e:
print(f"[v917] Stage1 FAILED inter={inter} E={exp}: {str(e)[:100]}", file=sys.stderr)
_fmoe._v917_s1_ready[0] = True
print("[v917] FlyDSL STAGE1 ENABLED!", file=sys.stderr)
if _compile_s2:
_s2_params = inspect.signature(_compile_s2).parameters
_s2_old_api = 'doweight' in _s2_params
try:
if _s2_old_api:
_compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,
tile_m=16, tile_n=256, tile_k=128, doweight=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True)
else:
_compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,
tile_m=16, tile_n=256, tile_k=256, doweight_stage2=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16")
except: pass
try:
if _s2_old_api:
_compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,
tile_m=16, tile_n=256, tile_k=128, doweight=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True, mode="atomic")
else:
_compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,
tile_m=16, tile_n=256, tile_k=256, doweight_stage2=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", mode="atomic")
except: pass
for d in [512, 2048]:
try:
if _s2_old_api:
_compile_s2(model_dim=7168, inter_dim=d, experts=33, topk=9,
tile_m=16, tile_n=256, tile_k=128, doweight=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True)
else:
_compile_s2(model_dim=7168, inter_dim=d, experts=33, topk=9,
tile_m=16, tile_n=256, tile_k=256, doweight_stage2=True,
a_dtype="fp4", b_dtype="fp4", out_dtype="bf16")
except: pass
_fmoe._v917_fly_ready[0] = True
print("[v917] FlyDSL STAGE2 ENABLED!", file=sys.stderr)
# Per-shape warmup
_warmed = set()
def _warmup(data):
(hs, guw, dw, guws, dws, guw_sh, dw_sh, guws_sh, dws_sh, tw, ti, cfg) = data
key = (cfg["bs"], cfg["n_routed_experts"], cfg["d_expert"])
if key in _warmed: return
_warmed.add(key)
try:
_ = fused_moe(hs[:2], guw_sh, dw_sh, tw[:2], ti[:2],
activation=ActivationType.Silu, quant_type=QuantType.per_1x32,
w1_scale=guws_sh, w2_scale=dws_sh,
hidden_pad=cfg["d_hidden_pad"]-cfg["d_hidden"],
intermediate_pad=cfg["d_expert_pad"]-cfg["d_expert"])
torch.cuda.synchronize()
except: pass
def custom_kernel(data: input_t) -> output_t:
_warmup(data)
(hs, guw, dw, guws, dws, guw_sh, dw_sh, guws_sh, dws_sh, tw, ti, cfg) = data
return fused_moe(hs, guw_sh, dw_sh, tw, ti,
activation=ActivationType.Silu, quant_type=QuantType.per_1x32,
w1_scale=guws_sh, w2_scale=dws_sh,
hidden_pad=cfg["d_hidden_pad"]-cfg["d_hidden"],
intermediate_pad=cfg["d_expert_pad"]-cfg["d_expert"])
scrolls · 308 lines total
Source code from GPU Mode and the KernelBot dataset · June 9 Researcher Reciprocity License v1.0
Changes from previous submission
Against this author's previous submission submission 714495.
# /// script# leaderboard = "amd-moe-mxfp4"# ///- """v913: FlyDSL for BOTH stages — the breakthrough.- flydsl 0.1.1 fixes scf.yield bug. FlyDSL stage1 replaces CK stage1.- Expected: -10 to -20μs (122→102-112μs)."""+ """v917: v913 fixed to survive AITER repo updates.+ Key fix: clone with depth=50 and find a commit that has _get_compiled_stage1+ (the old API). The repo was updated on ~April 4-5 and renamed the functions.+ Falls back to new API with runtime patch if needed."""import os, sys, torch, subprocess, shutil, re, functoolsos.environ['HIP_FORCE_DEV_KERNARG'] = '1'⋯ 20 unchanged linesif 'flydsl' in key:del sys.modules[key]import flydsl- print(f"[v913] flydsl {flydsl.__version__}", file=sys.stderr)+ print(f"[v917] flydsl {flydsl.__version__}", file=sys.stderr)- # ========== STEP 2: Setup AITER ==========+ # ========== STEP 2: Setup AITER — pin to working commit ==========subprocess.run(['git', 'checkout', '--', '.'], capture_output=True, cwd=_AITER_DIR, timeout=10)+ # Clone with enough depth to find old APIif not os.path.exists(os.path.join(_NEW_DIR, 'aiter')):- subprocess.run(['git', 'clone', '--depth=1', 'https://github.com/ROCm/aiter.git', _NEW_DIR],+ subprocess.run(['git', 'clone', '--depth=50', 'https://github.com/ROCm/aiter.git', _NEW_DIR],capture_output=True, text=True, timeout=120)- # Copy newer fused_moe.py+ # Find a commit where _get_compiled_stage1 exists in moe_kernels.py+ _mk_path_new = os.path.join(_NEW_DIR, 'aiter', 'ops', 'flydsl', 'moe_kernels.py')+ _found_old_api = False++ # Check HEAD first+ if os.path.exists(_mk_path_new):+ with open(_mk_path_new) as f:+ _mk_content = f.read()+ if '_get_compiled_stage1' in _mk_content:+ _found_old_api = True+ print("[v917] HEAD has old API (_get_compiled_stage1)", file=sys.stderr)++ if not _found_old_api:+ # Search older commits+ print("[v917] HEAD has new API, searching older commits...", file=sys.stderr)+ r = subprocess.run(['git', 'log', '--oneline', '-50', '--', 'aiter/ops/flydsl/moe_kernels.py'],+ capture_output=True, text=True, cwd=_NEW_DIR, timeout=10)+ print(f"[v917] Commits touching moe_kernels.py:\n{r.stdout[:500]}", file=sys.stderr)++ # Try each commit until we find one with old API+ r = subprocess.run(['git', 'log', '--format=%H', '-50'],+ capture_output=True, text=True, cwd=_NEW_DIR, timeout=10)+ commits = r.stdout.strip().split('\n')+ for commit in commits[1:]: # skip HEAD (already checked)+ if not commit.strip():+ continue+ # Checkout and check+ subprocess.run(['git', 'checkout', commit.strip()], capture_output=True, cwd=_NEW_DIR, timeout=10)+ if os.path.exists(_mk_path_new):+ with open(_mk_path_new) as f:+ _c = f.read()+ if '_get_compiled_stage1' in _c:+ _found_old_api = True+ print(f"[v917] Found old API in commit {commit.strip()[:12]}", file=sys.stderr)+ break+ if not _found_old_api:+ print("[v917] WARNING: Could not find old API in any commit!", file=sys.stderr)++ # Copy newer fused_moe.py (main dispatch file)shutil.copy2(os.path.join(_NEW_DIR, 'aiter', 'fused_moe.py'),os.path.join(_AITER_DIR, 'aiter', 'fused_moe.py'))- # Copy safe flydsl files (NOT gemm_kernels.py which needs flydsl.expr)+ # Copy flydsl filesflydsl_src = os.path.join(_NEW_DIR, 'aiter', 'ops', 'flydsl')flydsl_dst = os.path.join(_AITER_DIR, 'aiter', 'ops', 'flydsl')for fname in ['moe_kernels.py', 'utils.py']:- shutil.copy2(os.path.join(flydsl_src, fname), os.path.join(flydsl_dst, fname))+ src_p = os.path.join(flydsl_src, fname)+ dst_p = os.path.join(flydsl_dst, fname)+ if os.path.exists(src_p):+ shutil.copy2(src_p, dst_p)- # Patch __init__.py version check+ # Patch version checkinit_path = os.path.join(flydsl_dst, '__init__.py')with open(init_path) as f:init_src = f.read()⋯ 1 unchanged lineswith open(init_path, 'w') as f:f.write(init_src)- # Copy MLIR kernel generators+ # Copy kernelskernels_src = os.path.join(flydsl_src, 'kernels')kernels_dst = os.path.join(flydsl_dst, 'kernels')if os.path.exists(kernels_src):⋯ 2 unchanged linesif fname.endswith('.py'):shutil.copy2(os.path.join(kernels_src, fname), os.path.join(kernels_dst, fname))- # ========== STEP 3: Source patches ==========+ # ========== STEP 3: Source patches (same as v913) ==========_FMOE_PATH = os.path.join(_AITER_DIR, 'aiter', 'fused_moe.py')with open(_FMOE_PATH, 'r') as f:src = f.read()src = src.replace('split_k=ksplit', 'split_k=2')src = src.replace('splitk=ksplit', 'splitk=2')- # Minimal block_m overridebm_marker = 'def get_block_size_M(token, topk, expert, inter_dim):'if bm_marker in src:dec_search = src.rfind('@functools.lru_cache', 0, src.index(bm_marker))⋯ 22 unchanged lines'''src = src[:start] + new_bm + src[end:]- # Append FlyDSL stage1 + stage2 wrappers+ # Append FlyDSL stage1 + stage2 wrappers (same as v913)tail = r'''- # v913: FlyDSL BOTH stages- def _v913_fly_s1(hidden_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w1_scale=None, a1_scale=None, sorted_weights=None, **_kw):+ # v917: FlyDSL BOTH stages (v913 logic, pinned AITER)+ def _v917_fly_s1(hidden_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w1_scale=None, a1_scale=None, sorted_weights=None, **_kw):from aiter.ops.flydsl.moe_kernels import flydsl_moe_stage1flydsl_moe_stage1(a=hidden_states, w1=w1, sorted_token_ids=sorted_token_ids,sorted_expert_ids=sorted_expert_ids, num_valid_ids=num_valid_ids, out=out, topk=topk,tile_m=32, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",act="silu", w1_scale=w1_scale, a1_scale=a1_scale, sorted_weights=sorted_weights)- return out # flydsl writes in-place, but caller expects return value+ return out- def _v913_fly_s2_direct(inter_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w2_scale=None, a2_scale=None, sorted_weights=None, **_kw):+ def _v917_fly_s2_direct(inter_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w2_scale=None, a2_scale=None, sorted_weights=None, **_kw):from aiter.ops.flydsl.moe_kernels import flydsl_moe_stage2flydsl_moe_stage2(inter_states=inter_states, w2=w2, sorted_token_ids=sorted_token_ids,sorted_expert_ids=sorted_expert_ids, num_valid_ids=num_valid_ids, out=out, topk=topk,tile_m=16, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",mode="direct", w2_scale=w2_scale, a2_scale=a2_scale, sorted_weights=sorted_weights)- def _v913_fly_s2_atomic(inter_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w2_scale=None, a2_scale=None, sorted_weights=None, **_kw):+ def _v917_fly_s2_atomic(inter_states, w1, w2, sorted_token_ids, sorted_expert_ids, num_valid_ids, out, topk, w2_scale=None, a2_scale=None, sorted_weights=None, **_kw):from aiter.ops.flydsl.moe_kernels import flydsl_moe_stage2flydsl_moe_stage2(inter_states=inter_states, w2=w2, sorted_token_ids=sorted_token_ids,sorted_expert_ids=sorted_expert_ids, num_valid_ids=num_valid_ids, out=out, topk=topk,tile_m=16, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",mode="atomic", w2_scale=w2_scale, a2_scale=a2_scale, sorted_weights=sorted_weights)- _v913_orig_g2sc = get_2stage_cfgs- try: _v913_orig_g2sc.cache_clear()+ _v917_orig_g2sc = get_2stage_cfgs+ try: _v917_orig_g2sc.cache_clear()except: pass- _v913_fly_ready = [False]- _v913_s1_ready = [False]+ _v917_fly_ready = [False]+ _v917_s1_ready = [False]@functools.lru_cache(maxsize=2048)- def _v913_g2sc(*a, **kw):- md = _v913_orig_g2sc(*a, **kw)+ def _v917_g2sc(*a, **kw):+ md = _v917_orig_g2sc(*a, **kw)if len(a) >= 5:E_val = a[4]; bs_val = a[1]if E_val <= 100 and bs_val >= 512:md.use_non_temporal_load = True- # Replace stage1 with FlyDSL (if compiled)- if _v913_s1_ready[0] and not md.run_1stage and md.stage1 is not None:- md.stage1 = functools.partial(_v913_fly_s1)+ if _v917_s1_ready[0] and not md.run_1stage and md.stage1 is not None:+ md.stage1 = functools.partial(_v917_fly_s1)elif not md.run_1stage and md.stage1 is not None:_orig_s1 = md.stage1def _s1w(*aa, **kk):kk['splitk'] = 2return _orig_s1(*aa, **kk)md.stage1 = _s1w- # Replace stage2 with FlyDSL- if _v913_fly_ready[0] and not md.run_1stage and md.stage2 is not None:+ if _v917_fly_ready[0] and not md.run_1stage and md.stage2 is not None:E_val = a[4] if len(a) >= 5 else 0- md.stage2 = functools.partial(_v913_fly_s2_atomic if E_val > 100 else _v913_fly_s2_direct)+ md.stage2 = functools.partial(_v917_fly_s2_atomic if E_val > 100 else _v917_fly_s2_direct)return md- get_2stage_cfgs = _v913_g2sc+ get_2stage_cfgs = _v917_g2scimport sys as _s- print("[v913] FlyDSL BOTH stages patched", file=_s.stderr)+ print("[v917] FlyDSL BOTH stages patched", file=_s.stderr)'''src += tailwith open(_FMOE_PATH, 'w') as f:⋯ 6 unchanged linesfrom task import input_t, output_t# ========== STEP 5: Pre-compile FlyDSL ==========- try:- from aiter.ops.flydsl.moe_kernels import _get_compiled_stage1, _get_compiled_stage2+ # Try both old and new API names+ import aiter.ops.flydsl.moe_kernels as _mk- # Stage1 for all shapes+ _compile_s1 = getattr(_mk, '_get_compiled_stage1', None) or getattr(_mk, 'compile_flydsl_moe_stage1', None)+ _compile_s2 = getattr(_mk, '_get_compiled_stage2', None) or getattr(_mk, 'compile_flydsl_moe_stage2', None)++ print(f"[v917] S1 compile fn: {_compile_s1.__name__ if _compile_s1 else 'NONE'}", file=sys.stderr)+ print(f"[v917] S2 compile fn: {_compile_s2.__name__ if _compile_s2 else 'NONE'}", file=sys.stderr)++ if _compile_s1:+ # Detect which API we're using by checking parameter names+ import inspect+ _s1_params = inspect.signature(_compile_s1).parameters+ _use_old_api = 'doweight' in _s1_params # old API uses 'doweight', new uses 'doweight_stage1'+for inter, exp in [(256, 257), (512, 33), (2048, 33)]:try:- _get_compiled_stage1(model_dim=7168, inter_dim=inter, experts=exp, topk=9,- tile_m=32, tile_n=256, tile_k=128, doweight=True,- a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", act="silu")- print(f"[v913] Stage1 compiled: inter={inter} E={exp}", file=sys.stderr)+ if _use_old_api:+ _compile_s1(model_dim=7168, inter_dim=inter, experts=exp, topk=9,+ tile_m=32, tile_n=256, tile_k=128, doweight=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", act="silu")+ else:+ _compile_s1(model_dim=7168, inter_dim=inter, experts=exp, topk=9,+ tile_m=32, tile_n=256, tile_k=256, doweight_stage1=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", act="silu")+ print(f"[v917] Stage1 compiled: inter={inter} E={exp}", file=sys.stderr)except Exception as e:- print(f"[v913] Stage1 FAILED inter={inter} E={exp}: {str(e)[:100]}", file=sys.stderr)+ print(f"[v917] Stage1 FAILED inter={inter} E={exp}: {str(e)[:100]}", file=sys.stderr)- _fmoe._v913_s1_ready[0] = True- print("[v913] FlyDSL STAGE1 ENABLED!", file=sys.stderr)+ _fmoe._v917_s1_ready[0] = True+ print("[v917] FlyDSL STAGE1 ENABLED!", file=sys.stderr)- # Stage2 for all shapes- _get_compiled_stage2(model_dim=7168, inter_dim=256, experts=257, topk=9,- tile_m=16, tile_n=256, tile_k=128, doweight=True,- a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True)+ if _compile_s2:+ _s2_params = inspect.signature(_compile_s2).parameters+ _s2_old_api = 'doweight' in _s2_params+try:- _get_compiled_stage2(model_dim=7168, inter_dim=256, experts=257, topk=9,- tile_m=16, tile_n=256, tile_k=128, doweight=True,- a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True, mode="atomic")+ if _s2_old_api:+ _compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,+ tile_m=16, tile_n=256, tile_k=128, doweight=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True)+ else:+ _compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,+ tile_m=16, tile_n=256, tile_k=256, doweight_stage2=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16")except: pass+ try:+ if _s2_old_api:+ _compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,+ tile_m=16, tile_n=256, tile_k=128, doweight=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True, mode="atomic")+ else:+ _compile_s2(model_dim=7168, inter_dim=256, experts=257, topk=9,+ tile_m=16, tile_n=256, tile_k=256, doweight_stage2=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", mode="atomic")+ except: passfor d in [512, 2048]:- _get_compiled_stage2(model_dim=7168, inter_dim=d, experts=33, topk=9,- tile_m=16, tile_n=256, tile_k=128, doweight=True,- a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True)- _fmoe._v913_fly_ready[0] = True- print("[v913] FlyDSL STAGE2 ENABLED!", file=sys.stderr)- except Exception as e:- print(f"[v913] Pre-compile failed: {e}", file=sys.stderr)- import traceback; traceback.print_exc(file=sys.stderr)+ try:+ if _s2_old_api:+ _compile_s2(model_dim=7168, inter_dim=d, experts=33, topk=9,+ tile_m=16, tile_n=256, tile_k=128, doweight=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16", accumulate=True)+ else:+ _compile_s2(model_dim=7168, inter_dim=d, experts=33, topk=9,+ tile_m=16, tile_n=256, tile_k=256, doweight_stage2=True,+ a_dtype="fp4", b_dtype="fp4", out_dtype="bf16")+ except: pass+ _fmoe._v917_fly_ready[0] = True+ print("[v917] FlyDSL STAGE2 ENABLED!", file=sys.stderr)# Per-shape warmup_warmed = set()
scrolls · 297 diff lines total
Best evidence level for this revision: reported
JSON