For the complete documentation index, see llms.txt. Markdown versions of all pages are available by appending .md to any URL (e.g. /get-started.md).
Mojo struct
Struct_mega_ffn_ep_combine_send
struct Struct_mega_ffn_ep_combine_send
MOGG wrapper for the MegaFFN MoE FFN with the EP combine send fused in.
Same kernel as mo.composite.mega_ffn_nvfp4, with the L2 epilogue's peer
scatter-send and its pair-space arrival signal switched on. That makes the
send the only consumer of the FFN's output, so the local store and its TMA
descriptor are both elided and this op has NO tensor result: the down
projection's (max_recv_tokens, hidden_size) staging buffer -- the largest
single activation in the MoE region -- is never allocated.
The model emits this op directly rather than a graph-compiler pattern rewriting into it. The send only pays for itself at prefill-scale tokens per expert, so it is a per-model choice rather than a universal fusion rule, and the caller here is the one place that knows a combine WAIT is downstream and therefore that the arrival signal must be published.
A consumer still has to run ep.combine_wait afterwards to drain the peer
buffers and reduce; that op reads only the symmetric receive buffers and
never took the FFN's output, which is what makes dropping it possible.
NVFP4 only, unlike its sibling: the send needs a per-expert input scale, and E8M0 cannot carry one.
Implemented traits
Methods
execute
static def execute[a_type: DType, b_type: DType, scales_type: DType, //, combine_dtype: DType, hidden_size: Int, top_k: Int, n_experts: Int, max_token_per_rank: Int, n_gpus_per_node: Int, n_nodes: Int, clamp_activation: Bool, target: StringSpan[ImmStaticOrigin]](arrival_count: ManagedTensorSlice[IOSpec[_, _].MutableInput, static_spec=arrival_count.static_spec], atomic_counters: ManagedTensorSlice[IOSpec[_, _].MutableInput, static_spec=atomic_counters.static_spec], hidden_states: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=hidden_states.static_spec], gate_up_weight: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_weight.static_spec], gate_up_a_scales: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_a_scales.static_spec], gate_up_b_scales: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_b_scales.static_spec], down_weight: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=down_weight.static_spec], down_b_scales: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=down_b_scales.static_spec], expert_start_indices: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=expert_start_indices.static_spec], expert_ids: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=expert_ids.static_spec], a_scale_offsets: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=a_scale_offsets.static_spec], gate_up_expert_scales: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_expert_scales.static_spec], down_expert_scales: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=down_expert_scales.static_spec], c_input_scales: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=c_input_scales.static_spec], src_info: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=src_info.static_spec], recv_ptrs: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=recv_ptrs.static_spec], recv_count_ptrs: ManagedTensorSlice[IOSpec[_, _].Input, static_spec=recv_count_ptrs.static_spec], estimated_total_m: UInt32, num_active_experts: UInt32, context: DeviceContext)
Runs the fused MoE FFN and sends its output straight to the peers.
Parameters:
- a_type (
DType): Token / activation element type (inferred). - b_type (
DType): Weight element type (inferred). - scales_type (
DType): Block scale-factor dtype (inferred). - combine_dtype (
DType): Payload dtype of the combine phase; also the dtype of the output tensor this op declines to allocate. - hidden_size (
Int): Model hidden dimension, the send's row width. - top_k (
Int): Experts each token routes to; the receive buffer's second dimension. - n_experts (
Int): GLOBAL expert count across ranks. Distinct from the kernel'snum_experts, which is this rank's LOCAL count and comes from the weight shape; the sync-counter layout is keyed on the global count. - max_token_per_rank (
Int): Receive-buffer capacity per rank, in tokens. - n_gpus_per_node (
Int): GPUs per node; the number of live pointer-table entries. - n_nodes (
Int): Physical node count. - clamp_activation (
Bool):Trueselects the clamped (swigluoai) L1 activation; the dispatch supplies its canonical constants. - target (
StringSpan[ImmStaticOrigin]): Target GPU device.
Args:
- arrival_count (
ManagedTensorSlice[IOSpec[_, _].MutableInput, static_spec=arrival_count.static_spec]): Persistent cross-CTA pool-slot buffer, zeroed once at setup. Also holds the send's own per-expert state -- the arrival-signal election marker and the L2 completion counts -- in the padding of the slots the pool protocol never touches, which is why this op needs no per-launch-zeroed operand of its own. - atomic_counters (
ManagedTensorSlice[IOSpec[_, _].MutableInput, static_spec=atomic_counters.static_spec]): EP sync counters for this device. The send reads the combine-async region for its destination resolve and the rank-completion counter that follows it. - hidden_states (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=hidden_states.static_spec]): Dispatched tokens(max_recv_tokens, K1 // 2). - gate_up_weight (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_weight.static_spec]): Pre-permuted gate/up weights(E, N1, K1 // 2). - gate_up_a_scales (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_a_scales.static_spec]): Token A-scale tile for L1. - gate_up_b_scales (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_b_scales.static_spec]): Pre-permuted W13 scale tile. - down_weight (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=down_weight.static_spec]): Down-projection weights(E, N2, moe_dim // 2). - down_b_scales (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=down_b_scales.static_spec]): W2 scale tile. - expert_start_indices (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=expert_start_indices.static_spec]): Per-expert prefix-sum token offsets. - expert_ids (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=expert_ids.static_spec]): Active local expert IDs (-1= masked). - a_scale_offsets (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=a_scale_offsets.static_spec]): Per-expert scale-block offsets. - gate_up_expert_scales (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=gate_up_expert_scales.static_spec]): L1 per-expert output scaling. - down_expert_scales (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=down_expert_scales.static_spec]): L2 per-expert output scaling. - c_input_scales (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=c_input_scales.static_spec]): L1 per-expert input scale for the NVFP4 re-quant. - src_info (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=src_info.static_spec]): Per-row dispatch metadata the send resolves destinations from; the high bits carry the one-based source rank. - recv_ptrs (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=recv_ptrs.static_spec]): Peer combine receive-buffer addresses, one per rank. - recv_count_ptrs (
ManagedTensorSlice[IOSpec[_, _].Input, static_spec=recv_count_ptrs.static_spec]): Peer receive-count addresses the arrival signal releases, one per rank. - estimated_total_m (
UInt32): Estimated total non-padded tokens; the tile regime gate's numerator. - num_active_experts (
UInt32): Active expert slots; the kernel walks one shared expert list for both legs. - context (
DeviceContext): Device context.