Skip to content
KernelIndex
Search⌘K

submission 745202

Danishlynx · python · License unknown

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
AMD MXFP4 MoEsuite of 7 cases
AMD Instinct MI355X
69.9µs
#1 of 782
2026-04-06

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.

fp4tile_m=32, tile_n=256, tile_k=128, a_dtype="fp4", b_dtype="fp4", out_dtype="bf16",
split-ksrc = src.replace('split_k=ksplit', 'split_k=2')

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, functools
os.environ['HIP_FORCE_DEV_KERNARG'] = '1'
⋯ 20 unchanged lines
if '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 API
if 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 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']:
- 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 check
init_path = os.path.join(flydsl_dst, '__init__.py')
with open(init_path) as f:
init_src = f.read()
⋯ 1 unchanged lines
with open(init_path, 'w') as f:
f.write(init_src)
- # Copy MLIR kernel generators
+ # Copy kernels
kernels_src = os.path.join(flydsl_src, 'kernels')
kernels_dst = os.path.join(flydsl_dst, 'kernels')
if os.path.exists(kernels_src):
⋯ 2 unchanged lines
if 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 override
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))
⋯ 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_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 # 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_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 _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_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)
- _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.stage1
def _s1w(*aa, **kk):
kk['splitk'] = 2
return _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_g2sc
import sys as _s
- print("[v913] FlyDSL BOTH stages patched", file=_s.stderr)
+ print("[v917] FlyDSL BOTH stages patched", file=_s.stderr)
'''
src += tail
with open(_FMOE_PATH, 'w') as f:
⋯ 6 unchanged lines
from 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: pass
for 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