Skip to content

vllm.models.kimi_k3.nvidia.model

Kimi-K3 multimodal model implementation for vLLM.

Classes:

KimiK3ForConditionalGeneration

Bases: Module, SupportsMultiModal, SupportsEncoderCudaGraph, SupportsPP, SupportsQuant, SupportsEagle3, HasInnerState, IsHybrid

Kimi-K3 model with Kimi-K2.5 vision and KimiLinear text.

Source code in vllm/models/kimi_k3/nvidia/model.py
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
1501
1502
1503
1504
1505
1506
1507
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
1564
1565
1566
1567
1568
1569
1570
1571
1572
1573
1574
1575
1576
1577
1578
1579
1580
1581
1582
1583
1584
1585
1586
1587
1588
1589
1590
1591
1592
1593
1594
1595
1596
1597
1598
1599
1600
1601
1602
1603
1604
1605
1606
1607
1608
1609
1610
1611
1612
1613
1614
1615
1616
1617
1618
1619
1620
1621
1622
1623
1624
1625
1626
1627
1628
1629
1630
1631
1632
1633
1634
1635
1636
1637
1638
1639
1640
1641
1642
1643
1644
1645
1646
1647
1648
1649
1650
1651
1652
1653
1654
1655
1656
1657
1658
1659
1660
1661
1662
1663
1664
1665
1666
1667
1668
1669
1670
1671
1672
1673
1674
1675
1676
1677
1678
1679
1680
1681
1682
1683
1684
1685
1686
1687
1688
1689
1690
1691
1692
1693
1694
1695
1696
1697
1698
1699
1700
1701
1702
1703
1704
1705
1706
1707
1708
1709
1710
1711
1712
1713
1714
1715
1716
1717
1718
1719
1720
1721
1722
1723
1724
1725
1726
1727
1728
1729
1730
1731
1732
1733
1734
1735
1736
1737
1738
1739
1740
1741
1742
1743
1744
1745
1746
1747
1748
1749
1750
1751
1752
1753
1754
1755
1756
1757
1758
1759
1760
1761
1762
1763
1764
1765
1766
1767
1768
1769
1770
1771
1772
1773
1774
1775
1776
1777
1778
1779
1780
1781
1782
1783
1784
1785
1786
1787
1788
1789
1790
1791
1792
1793
1794
1795
1796
1797
1798
1799
1800
1801
1802
1803
1804
1805
1806
1807
1808
1809
1810
1811
1812
1813
1814
1815
1816
1817
1818
1819
1820
1821
1822
1823
1824
1825
1826
1827
1828
1829
1830
1831
1832
1833
1834
1835
1836
1837
1838
1839
1840
1841
1842
1843
1844
1845
1846
1847
1848
1849
1850
1851
1852
1853
1854
1855
1856
1857
1858
1859
@MULTIMODAL_REGISTRY.register_processor(
    KimiK3MultiModalProcessor,
    info=KimiK3ProcessingInfo,
    dummy_inputs=KimiK3DummyInputsBuilder,
)
class KimiK3ForConditionalGeneration(
    nn.Module,
    SupportsMultiModal,
    SupportsEncoderCudaGraph,
    SupportsPP,
    SupportsQuant,
    SupportsEagle3,
    HasInnerState,
    IsHybrid,
):
    """Kimi-K3 model with Kimi-K2.5 vision and KimiLinear text."""

    supports_encoder_tp_data = True

    hf_to_vllm_mapper = WeightsMapper(
        orig_to_new_prefix={
            "language_model.layers.": "language_model.model.layers.",
            "mm_projector.proj.0": "mm_projector.linear_1",
            "mm_projector.proj.2": "mm_projector.linear_2",
        }
    )

    @classmethod
    def get_placeholder_str(cls, modality: str, i: int) -> str | None:
        if modality == "image":
            return "<|kimi_image_placeholder|>"
        raise ValueError(f"Unsupported modality: {modality}")

    def __init__(
        self,
        vllm_config: VllmConfig,
        prefix: str = "",
    ) -> None:
        super().__init__()
        model_config = vllm_config.model_config
        config: KimiK3Config = model_config.hf_config
        self.config = config
        self.model_config = model_config
        quant_config = vllm_config.quant_config

        multimodal_config = model_config.multimodal_config
        assert multimodal_config is not None
        self.use_data_parallel = is_vit_use_data_parallel(
            config.vision_config.num_attention_heads
        )
        self.hidden_size = config.text_config.hidden_size
        self.device = current_platform.current_device()

        with self._mark_tower_model(vllm_config, "image"):
            self.vision_tower = MoonViT3dPretrainedModel(
                config.vision_config,
                quant_config=self._maybe_ignore_quant_config(quant_config),
                prefix=maybe_prefix(prefix, "vision_tower"),
            )
            if self._maybe_ignore_quant_config(quant_config) is not None:
                self.vision_tower = self.vision_tower.to(device=self.device)
            else:
                self.vision_tower = self.vision_tower.to(
                    device=self.device, dtype=model_config.dtype
                )

            vision_attn = self.vision_tower.encoder.blocks[0].attn
            if vision_attn.is_flash_attn_backend and vision_attn._fa_version == 4:
                from vllm.models.kimi_k3.nvidia.ops.vision_fa4_warmup import (
                    KimiK3VisionFA4WarmupConfig,
                    register_kimi_k3_vision_fa4_warmup,
                )

                merge_height, merge_width = config.vision_config.merge_kernel_size
                mm_config = model_config.get_multimodal_config()
                assert mm_config is not None
                register_kimi_k3_vision_fa4_warmup(
                    KimiK3VisionFA4WarmupConfig(
                        num_heads=vision_attn.num_heads,
                        head_dim=vision_attn.head_size,
                        dtype=vision_attn.dtype,
                        max_batch_size=(
                            vllm_config.scheduler_config.max_num_seqs
                            * mm_config.get_limit_per_prompt("image")
                        ),
                        max_seqlen=(
                            vllm_config.scheduler_config.max_num_encoder_input_tokens
                            * merge_height
                            * merge_width
                        ),
                    )
                )

            self.mm_projector = KimiK25MultiModalProjector(
                config=config.vision_config,
                use_data_parallel=self.use_data_parallel,
                quant_config=self._maybe_ignore_quant_config(quant_config),
                prefix=maybe_prefix(prefix, "mm_projector"),
            )
            self.mm_projector = self.mm_projector.to(
                device=self.device, dtype=model_config.dtype
            )

        self.quant_config = quant_config
        with self._mark_language_model(vllm_config):
            self.language_model = init_vllm_registered_model(
                vllm_config=vllm_config,
                hf_config=config.text_config,
                prefix=maybe_prefix(prefix, "language_model"),
                architectures=["KimiLinearForCausalLM"],
            )
        self.make_empty_intermediate_tensors = (  # type: ignore[method-assign]
            self.language_model.make_empty_intermediate_tensors
        )
        self.media_placeholder: int = self.config.media_placeholder_token_id

    # -- SupportsEncoderCudaGraph protocol methods --

    def get_encoder_cudagraph_config(self):
        from vllm.v1.worker.encoder_cudagraph_defs import EncoderCudaGraphConfig

        return EncoderCudaGraphConfig(
            modalities=["image"],
            buffer_keys=[
                "pixel_values",
                "pos_embeds",
                "rope_freqs_cis",
                "cu_seqlens",
                "max_seqlen",
                "sequence_lengths",
                "merge_gather_idx",
            ],
            out_hidden_size=self.hidden_size,
        )

    def get_encoder_cudagraph_budget_range(
        self, vllm_config: VllmConfig
    ) -> tuple[int, int]:
        min_budget = 64
        max_budget = min(
            vllm_config.scheduler_config.max_num_batched_tokens,
            self.model_config.max_model_len,
        )
        return min_budget, max_budget

    @staticmethod
    def _get_grid_thws(mm_kwargs: dict[str, Any]) -> list[list[int]]:
        grid_thws = mm_kwargs["grid_thws"]
        if not isinstance(grid_thws, list):
            grid_thws = grid_thws.tolist()
        return grid_thws

    @staticmethod
    def _get_pixel_values(mm_kwargs: dict[str, Any]) -> torch.Tensor:
        pixel_values = mm_kwargs["pixel_values"]
        if isinstance(pixel_values, list):
            pixel_values = torch.cat(pixel_values)
        if pixel_values.ndim in (3, 5):
            pixel_values = pixel_values.reshape(
                pixel_values.shape[0] * pixel_values.shape[1],
                *pixel_values.shape[2:],
            )
        return pixel_values

    def get_encoder_cudagraph_item_specs(self, mm_kwargs: dict[str, Any]):
        from vllm.v1.worker.encoder_cudagraph_defs import EncoderItemSpec

        kh, kw = self.config.vision_config.merge_kernel_size
        return [
            EncoderItemSpec(
                input_size=t * h * w,
                output_tokens=(h // kh) * (w // kw),
            )
            for t, h, w in self._get_grid_thws(mm_kwargs)
        ]

    def select_encoder_cudagraph_items(
        self, mm_kwargs: dict[str, Any], indices: list[int]
    ) -> dict[str, Any]:
        grid_thws = self._get_grid_thws(mm_kwargs)
        pixel_values = self._get_pixel_values(mm_kwargs)
        source_grid = mm_kwargs["grid_thws"]

        if not indices:
            empty_grid = (
                source_grid[:0] if isinstance(source_grid, torch.Tensor) else []
            )
            return {"pixel_values": pixel_values[:0], "grid_thws": empty_grid}

        patch_counts = [t * h * w for t, h, w in grid_thws]
        offsets = [0]
        for count in patch_counts:
            offsets.append(offsets[-1] + count)
        selected_pixel_values = torch.cat(
            [pixel_values[offsets[i] : offsets[i + 1]] for i in indices]
        )
        grid_device = (
            source_grid.device if isinstance(source_grid, torch.Tensor) else None
        )
        selected_grid = torch.tensor(
            [grid_thws[i] for i in indices],
            dtype=torch.long,
            device=grid_device,
        )
        return {"pixel_values": selected_pixel_values, "grid_thws": selected_grid}

    def prepare_encoder_cudagraph_capture_inputs(
        self,
        token_budget: int,
        max_batch_size: int,
        max_frames_per_batch: int,
        device: torch.device,
        dtype: torch.dtype,
        path: str = "default",
    ):
        from vllm.v1.worker.encoder_cudagraph_defs import (
            EncoderCudaGraphCaptureInputs,
        )

        kh, kw = self.config.vision_config.merge_kernel_size
        per_item_output = (token_budget + max_batch_size - 1) // max_batch_size
        rope = self.vision_tower.encoder.rope_2d
        max_output_width = rope.max_width // kw
        max_output_height = rope.max_height // kh
        output_width = min(math.ceil(math.sqrt(per_item_output)), max_output_width)
        output_height = (per_item_output + output_width - 1) // output_width
        if output_height > max_output_height:
            output_height = max_output_height
            output_width = (per_item_output + output_height - 1) // output_height
        if output_width > max_output_width:
            raise ValueError(
                f"Encoder CUDA graph budget {token_budget} exceeds K3 RoPE "
                f"capacity for max_batch_size={max_batch_size}"
            )
        grid_thws = [
            [1, output_height * kh, output_width * kw] for _ in range(max_batch_size)
        ]

        patch_size: int | tuple[int, int] = self.config.vision_config.patch_size
        if isinstance(patch_size, int):
            patch_size = (patch_size, patch_size)
        total_patches = sum(t * h * w for t, h, w in grid_thws)
        pixel_values = torch.randn(
            total_patches,
            3,
            patch_size[0],
            patch_size[1],
            device=device,
            dtype=dtype,
        )
        metadata = self.vision_tower.prepare_encoder_cudagraph_metadata(
            grid_thws,
            max_batch_size=max_batch_size,
            max_seqlen_override=max(
                token_budget * kh * kw,
                max(t * h * w for t, h, w in grid_thws),
            ),
            device=device,
        )
        return EncoderCudaGraphCaptureInputs(
            values=metadata | {"pixel_values": pixel_values}
        )

    def prepare_encoder_cudagraph_replay_buffers(
        self,
        mm_kwargs: dict[str, Any],
        max_batch_size: int,
        max_frames_per_batch: int,
        path: str = "default",
    ):
        from vllm.v1.worker.encoder_cudagraph_defs import (
            EncoderCudaGraphReplayBuffers,
        )

        pixel_values = self._get_pixel_values(mm_kwargs)
        metadata = self.vision_tower.prepare_encoder_cudagraph_metadata(
            self._get_grid_thws(mm_kwargs),
            max_batch_size=max_batch_size,
            device=pixel_values.device,
        )
        return EncoderCudaGraphReplayBuffers(
            values=metadata | {"pixel_values": pixel_values}
        )

    def _project_encoder_features(self, image_features: torch.Tensor) -> torch.Tensor:
        projector_dtype = next(self.mm_projector.parameters()).dtype
        if image_features.dtype != projector_dtype:
            image_features = image_features.to(projector_dtype)
        output = self.mm_projector(image_features)
        return output.reshape(-1, output.shape[-1])

    def encoder_cudagraph_forward(
        self,
        values: dict[str, torch.Tensor],
        path: str = "default",
    ) -> torch.Tensor:
        pixel_values = values.pop("pixel_values")
        image_features = self.vision_tower(pixel_values, None, encoder_metadata=values)
        return self._project_encoder_features(image_features)

    def encoder_eager_forward(
        self,
        mm_kwargs: dict[str, Any],
        path: str = "default",
    ) -> torch.Tensor:
        image_features = self.vision_tower(
            self._get_pixel_values(mm_kwargs).to(
                next(self.vision_tower.parameters()).dtype
            ),
            self._get_grid_thws(mm_kwargs),
        )
        return self._project_encoder_features(torch.cat(image_features))

    def _maybe_ignore_quant_config(
        self, quant_config: QuantizationConfig | None
    ) -> QuantizationConfig | None:
        if isinstance(quant_config, compressed_tensors.CompressedTensorsConfig):
            return None
        return quant_config

    def _parse_and_validate_media_input(
        self, **kwargs: object
    ) -> KimiK25MediaPixelInputs | None:
        pixel_values = kwargs.pop("pixel_values", None)
        grid_thws = kwargs.pop("grid_thws", None)
        if pixel_values is None:
            return None

        if isinstance(pixel_values, list):
            pixel_values = torch.cat(cast(list[torch.Tensor], pixel_values), dim=0)
        if not isinstance(pixel_values, torch.Tensor):
            raise TypeError(
                "pixel_values must be a tensor or a list of tensors, "
                f"got {type(pixel_values)}"
            )

        if len(pixel_values.shape) == 5 or len(pixel_values.shape) == 3:
            pixel_values = pixel_values.reshape(
                pixel_values.shape[0] * pixel_values.shape[1], *pixel_values.shape[2:]
            )

        target_dtype = next(self.vision_tower.parameters()).dtype
        pixel_values = pixel_values.to(target_dtype)
        assert isinstance(grid_thws, torch.Tensor), (
            f"expect grid_thws to be a tensor, got {type(grid_thws)}"
        )
        grid_thws = grid_thws.reshape(-1, grid_thws.shape[-1])
        assert grid_thws.ndim == 2 and grid_thws.size(1) == 3, (
            f"unexpected shape for grid_thws: {grid_thws.shape}"
        )

        return KimiK25MediaPixelInputs(
            type="pixel_values",
            pixel_values=pixel_values,
            grid_thws=grid_thws,
        )

    def _process_media_input(
        self, media_input: KimiK25MediaPixelInputs
    ) -> list[torch.Tensor]:
        media_features = vision_tower_forward(
            self.vision_tower,
            media_input["pixel_values"],
            media_input["grid_thws"],
            mm_projector=self.mm_projector,
            use_data_parallel=self.use_data_parallel,
        )
        return media_features

    def embed_multimodal(self, **kwargs: object) -> NestedTensors | None:
        media_input = self._parse_and_validate_media_input(**kwargs)
        if media_input is None:
            return None
        return self._process_media_input(media_input)

    def forward(  # type: ignore[override]
        self,
        input_ids: torch.Tensor,
        positions: torch.Tensor,
        intermediate_tensors: IntermediateTensors | None = None,
        inputs_embeds: torch.Tensor | None = None,
        **kwargs: object,
    ) -> torch.Tensor | IntermediateTensors | tuple[torch.Tensor, list[torch.Tensor]]:
        if intermediate_tensors is not None:
            inputs_embeds = None
        return self.language_model(
            input_ids=input_ids,
            positions=positions,
            intermediate_tensors=intermediate_tensors,
            inputs_embeds=inputs_embeds,
        )

    def compute_logits(self, hidden_states: torch.Tensor, **kwargs) -> torch.Tensor:
        return self.language_model.compute_logits(hidden_states)

    def copy_inputs_before_cuda_graphs(self, input_buffers, **kwargs):
        return self.language_model.mamba_cache.copy_inputs_before_cuda_graphs(
            input_buffers, **kwargs
        )

    def get_seqlen_agnostic_capture_inputs(self, batch_size: int):
        return self.language_model.mamba_cache.get_seqlen_agnostic_capture_inputs(
            batch_size
        )

    @classmethod
    def get_mamba_state_dtype_from_config(cls, vllm_config: VllmConfig):
        text_config = vllm_config.model_config.hf_config.text_config
        temp_vllm_config = vllm_config.with_hf_config(text_config)
        return KimiLinearForCausalLM.get_mamba_state_dtype_from_config(temp_vllm_config)

    @classmethod
    def get_mamba_state_shape_from_config(cls, vllm_config: VllmConfig):
        text_config = vllm_config.model_config.hf_config.text_config
        temp_vllm_config = vllm_config.with_hf_config(text_config)
        return KimiLinearForCausalLM.get_mamba_state_shape_from_config(temp_vllm_config)

    @classmethod
    def get_mamba_state_copy_func(cls):
        return KimiLinearForCausalLM.get_mamba_state_copy_func()

    def load_weights(self, weights: Iterable[tuple[str, torch.Tensor]]):
        loader = AutoWeightsLoader(self)
        return loader.load_weights(weights, mapper=self.hf_to_vllm_mapper)

KimiK3MegaMoEExperts

Bases: DeepseekV4MegaMoEExperts

Kimi K3 adapter for the DeepGEMM MegaMoE kernel.

Source code in vllm/models/kimi_k3/nvidia/model.py
class KimiK3MegaMoEExperts(DeepseekV4MegaMoEExperts):
    """Kimi K3 adapter for the DeepGEMM MegaMoE kernel."""

    _kimi_symm_buffer_cache: dict[tuple[object, ...], object] = {}
    _synchronized_ep_groups: set[tuple[int, int]] = set()

    def __init__(
        self,
        *args,
        activation: str,
        activation_beta: float | None,
        activation_linear_beta: float | None,
        **kwargs,
    ):
        super().__init__(*args, **kwargs)
        self.activation = activation
        self.activation_beta = activation_beta
        self.activation_linear_beta = activation_linear_beta

    def synchronize_first_launch(self) -> None:
        ep_group = get_ep_group()
        device = torch.accelerator.current_device_index()
        key = (id(ep_group.cpu_group), device)
        if key in self._synchronized_ep_groups:
            return
        torch.accelerator.synchronize()
        torch.distributed.barrier(group=ep_group.cpu_group)
        self._synchronized_ep_groups.add(key)

    def finalize_weights(self) -> None:
        if self._transformed_l1_weights is not None:
            return

        self._check_runtime_supported()
        from vllm.utils.deep_gemm import _import_deep_gemm

        deep_gemm = _import_deep_gemm()
        w13_scale = deep_gemm.transform_sf_into_required_layout(
            self._ue8m0_uint8_to_float(self.w13_weight_scale.data).contiguous(),
            2 * self.intermediate_size,
            self.hidden_size,
            (1, 32),
            self.num_local_experts,
        )
        w2_scale = deep_gemm.transform_sf_into_required_layout(
            self._ue8m0_uint8_to_float(self.w2_weight_scale.data).contiguous(),
            self.hidden_size,
            self.intermediate_size,
            (1, 32),
            self.num_local_experts,
        )
        self._transformed_l1_weights, self._transformed_l2_weights = (
            deep_gemm.transform_weights_for_mega_moe(
                (self.w13_weight.data.view(torch.int8).contiguous(), w13_scale),
                (self.w2_weight.data.view(torch.int8).contiguous(), w2_scale),
                activation=self.activation,
            )
        )
        self.w13_weight = None
        self.w13_weight_scale = None
        self.w2_weight = None
        self.w2_weight_scale = None

    def get_symm_buffer(self):
        from vllm.utils.deep_gemm import _import_deep_gemm

        deep_gemm = _import_deep_gemm()
        group = get_ep_group().device_group
        device = torch.accelerator.current_device_index()
        key = (
            id(group),
            device,
            self.num_experts,
            self.max_num_tokens,
            self.top_k,
            self.hidden_size,
            self.intermediate_size,
            self.activation,
        )
        symm_buffer = self._kimi_symm_buffer_cache.get(key)
        if symm_buffer is None:
            symm_buffer = deep_gemm.get_symm_buffer_for_mega_moe(
                group,
                self.num_experts,
                self.max_num_tokens,
                self.top_k,
                self.hidden_size,
                self.intermediate_size,
                activation=self.activation,
            )
            self._kimi_symm_buffer_cache[key] = symm_buffer
        return symm_buffer

    def forward(
        self,
        hidden_states: torch.Tensor,
        topk_weights: torch.Tensor,
        topk_ids: torch.Tensor,
        *,
        activation_clamp: float | None,
        fast_math: bool = True,
    ) -> torch.Tensor:
        self.synchronize_first_launch()
        if hidden_states.shape[0] > self.max_num_tokens:
            raise ValueError(
                f"Kimi K3 MegaMoE got {hidden_states.shape[0]} tokens, "
                f"but its symmetric buffer supports {self.max_num_tokens}."
            )
        y = torch.empty_like(hidden_states, dtype=torch.bfloat16)
        from vllm.utils.deep_gemm import _import_deep_gemm

        deep_gemm = _import_deep_gemm()
        symm_buffer = self.get_symm_buffer()
        num_tokens = hidden_states.shape[0]
        is_padding = None
        if envs.VLLM_MOE_SKIP_PADDING and is_forward_context_available():
            is_padding = get_forward_context().is_padding
            if is_padding is not None:
                is_padding = is_padding[:num_tokens]

        eplb_state = self.eplb_state
        if eplb_state.logical_to_physical_map is not None:
            assert eplb_state.expert_load_view is not None
            assert eplb_state.logical_replica_count is not None
            assert eplb_state.should_record_tensor is not None
            if is_padding is not None:
                topk_ids = torch.where(is_padding.unsqueeze(1), -1, topk_ids)
            topk_ids = eplb_map_to_physical_and_record(
                topk_ids=topk_ids,
                expert_load_view=eplb_state.expert_load_view,
                logical_to_physical_map=eplb_state.logical_to_physical_map,
                logical_replica_count=eplb_state.logical_replica_count,
                record_enabled=eplb_state.should_record_tensor,
                num_unpadded_tokens=eplb_state.num_unpadded_tokens_tensors[
                    dbo_current_ubatch_id()
                ]
                if eplb_state.num_unpadded_tokens_tensors is not None
                else None,
            )

        prepare_megamoe_inputs(
            hidden_states,
            topk_weights,
            topk_ids,
            symm_buffer.x[:num_tokens],
            symm_buffer.x_sf[:num_tokens],
            symm_buffer.topk_idx[:num_tokens],
            symm_buffer.topk_weights[:num_tokens],
            is_padding=is_padding,
        )
        self.finalize_weights()
        assert self._transformed_l1_weights is not None
        assert self._transformed_l2_weights is not None
        deep_gemm.fp8_fp4_mega_moe(
            y,
            self._transformed_l1_weights,
            self._transformed_l2_weights,
            symm_buffer,
            activation_clamp=activation_clamp,
            activation=self.activation,
            activation_beta=self.activation_beta,
            activation_linear_beta=self.activation_linear_beta,
            fast_math=fast_math,
        )
        return y

KimiMoE

Bases: Module

Source code in vllm/models/kimi_k3/nvidia/model.py
class KimiMoE(nn.Module):
    def __init__(
        self,
        config: KimiLinearConfig,
        vllm_config: VllmConfig,
        quant_config: QuantizationConfig | None = None,
        prefix: str = "",
        layer_idx: int = 0,
        use_sequence_parallel: bool = False,
    ):
        super().__init__()
        hidden_size = config.hidden_size
        moe_intermediate_size = config.moe_intermediate_size
        num_experts = config.num_experts
        num_experts_per_token = config.num_experts_per_token
        assert moe_intermediate_size is not None
        assert num_experts is not None
        assert num_experts_per_token is not None
        moe_renormalize = config.moe_renormalize
        routed_expert_hidden_size = config.routed_expert_hidden_size
        self.use_latent_moe = routed_expert_hidden_size is not None
        self.moe_hidden_size = (
            routed_expert_hidden_size
            if routed_expert_hidden_size is not None
            else hidden_size
        )
        self.latent_moe_use_norm = config.latent_moe_use_norm
        self.tp_size = get_tensor_model_parallel_world_size()
        self.routed_scaling_factor = config.routed_scaling_factor
        self.moe_renormalize = moe_renormalize
        self.use_grouped_topk = config.use_grouped_topk
        self.num_expert_group = config.num_expert_group
        self.topk_group = config.topk_group
        self.moe_router_activation_func = config.moe_router_activation_func
        self.num_shared_experts = config.num_shared_experts
        self.layer_idx = layer_idx
        self.use_mega_moe = (
            vllm_config.kernel_config.moe_backend == "deep_gemm_mega_moe"
        )
        if self.use_mega_moe and not vllm_config.parallel_config.enable_expert_parallel:
            raise NotImplementedError(
                "Kimi K3 MegaMoE requires expert parallel. Enable it with "
                "--enable-expert-parallel."
            )
        if self.use_mega_moe and config.hidden_act != "situ":
            raise ValueError("Kimi K3 MegaMoE requires SITU activation.")
        if self.use_mega_moe and not self.use_latent_moe:
            raise ValueError("Kimi K3 MegaMoE requires latent MoE projections.")
        if self.use_mega_moe and not self.use_grouped_topk:
            raise ValueError("Kimi K3 MegaMoE requires grouped top-k routing.")
        if self.use_mega_moe and (self.num_expert_group != 1 or self.topk_group != 1):
            raise NotImplementedError(
                "Kimi K3 MegaMoE currently requires one expert group."
            )
        self.padded_moe_intermediate_size = moe_intermediate_size
        min_moe_intermediate_per_partition = getattr(
            config, "min_moe_intermediate_per_partition", 256
        )
        if self.tp_size > 1:
            moe_intermediate_per_partition = moe_intermediate_size // self.tp_size
            if moe_intermediate_per_partition < min_moe_intermediate_per_partition:
                self.padded_moe_intermediate_size = (
                    min_moe_intermediate_per_partition * self.tp_size
                )
        activation_situ_beta = (
            config.activation_situ_beta if config.hidden_act == "situ" else None
        )
        activation_situ_linear_beta = (
            config.activation_situ_linear_beta if config.hidden_act == "situ" else None
        )

        # Route with fp32 logits for numerically stable expert selection.
        self.gate = GateLinear(
            input_size=hidden_size,
            output_size=num_experts,
            bias=False,
            out_dtype=torch.float32,
            prefix=f"{prefix}.gate",
        )

        self.gate.e_score_correction_bias = nn.Parameter(
            torch.empty(num_experts, dtype=torch.float32)
        )

        if self.num_shared_experts is not None:
            shared_intermediate_size = moe_intermediate_size * self.num_shared_experts
            self.shared_experts = KimiMLP(
                hidden_size=config.hidden_size,
                intermediate_size=shared_intermediate_size,
                hidden_act=config.hidden_act,
                quant_config=quant_config,
                reduce_results=False,
                use_sequence_parallel=use_sequence_parallel,
                prefix=f"{prefix}.shared_experts",
                activation_situ_beta=activation_situ_beta,
                activation_situ_linear_beta=activation_situ_linear_beta,
            )
        else:
            self.shared_experts = None

        self.routed_expert_down_proj: ReplicatedLinear | None
        self.routed_expert_norm: RMSNorm | None
        self.routed_expert_up_proj: ReplicatedLinear | None
        self.routed_output_transform: KimiRoutedOutputTransform | None
        if self.use_latent_moe:
            self.routed_expert_down_proj = ReplicatedLinear(
                hidden_size,
                self.moe_hidden_size,
                bias=False,
                quant_config=None,
                prefix=f"{prefix}.routed_expert_down_proj",
            )
            self.routed_expert_norm = (
                RMSNorm(self.moe_hidden_size, eps=config.rms_norm_eps)
                if self.latent_moe_use_norm
                else None
            )
            # Replicated up-proj: the full weight lives on every rank and
            # produces the full hidden dim locally. This lets LatentMoERunner
            # fuse the latent and shared reductions into a single all-reduce
            # (concat the two partials, reduce once), then run the up-proj and
            # shared add locally with no further collective.
            self.routed_expert_up_proj = ReplicatedLinear(
                self.moe_hidden_size,
                hidden_size,
                bias=False,
                quant_config=None,
                prefix=f"{prefix}.routed_expert_up_proj",
            )

            self.routed_output_transform = KimiRoutedOutputTransform(
                self.routed_expert_norm, self.routed_expert_up_proj
            )
            # Auxiliary CUDA stream to overlap the router gate with the routed
            # down projection on decode-sized batches (gated by
            # VLLM_ROUTED_DOWN_PROJ_STREAM_TOKEN_THRESHOLD).
            self._down_proj_stream: torch.cuda.Stream | None = aux_stream()
            self._down_proj_events = (torch.cuda.Event(), torch.cuda.Event())
        else:
            self.routed_expert_down_proj = None
            self.routed_expert_norm = None
            self.routed_expert_up_proj = None
            self.routed_output_transform = None

        if self.use_mega_moe:
            ep_group = get_ep_group()
            ep_size = ep_group.world_size
            ep_rank = ep_group.rank_in_group
            if num_experts % ep_size != 0:
                raise ValueError(
                    f"Kimi K3 num_experts={num_experts} must be divisible by "
                    f"EP size {ep_size}."
                )
            num_local_experts = num_experts // ep_size
            self.experts = KimiK3MegaMoEExperts(
                vllm_config,
                num_experts=num_experts,
                num_local_experts=num_local_experts,
                experts_start_idx=ep_rank * num_local_experts,
                top_k=num_experts_per_token,
                hidden_size=self.moe_hidden_size,
                intermediate_size=self.padded_moe_intermediate_size,
                prefix=f"{prefix}.experts",
                activation="situ",
                activation_beta=activation_situ_beta,
                activation_linear_beta=activation_situ_linear_beta,
            )
        else:
            enable_tail_fusion = envs.VLLM_ENABLE_K3_LATENT_MOE_TAIL_FUSION
            self.experts = FusedMoE(
                shared_experts=self.shared_experts,
                num_experts=num_experts,
                top_k=num_experts_per_token,
                hidden_size=self.moe_hidden_size,
                intermediate_size=self.padded_moe_intermediate_size,
                activation=config.hidden_act,
                activation_situ_beta=activation_situ_beta,
                activation_situ_linear_beta=activation_situ_linear_beta,
                renormalize=moe_renormalize,
                quant_config=quant_config,
                use_grouped_topk=config.use_grouped_topk,
                num_expert_group=config.num_expert_group,
                topk_group=config.topk_group,
                prefix=f"{prefix}.experts",
                scoring_func=config.moe_router_activation_func,
                e_score_correction_bias=self.gate.e_score_correction_bias,
                routed_scaling_factor=self.routed_scaling_factor,
                # Down projection runs outside FusedMoE so it can overlap the
                # router gate on the aux stream (see forward()); the original
                # hidden states are passed to forward() as shared_experts_input
                # so shared experts still see the untransformed input.
                routed_input_transform=None,
                routed_output_transform=self.routed_output_transform,
                is_sequence_parallel=use_sequence_parallel,
                runner_cls=LatentMoERunner if self.use_latent_moe else None,
                runner_args=(
                    {"enable_k3_latent_moe_tail_fusion": enable_tail_fusion}
                    if self.use_latent_moe
                    else None
                ),
            )
        if self.padded_moe_intermediate_size != moe_intermediate_size:
            w13_weight = getattr(self.experts, "w13_weight", None)
            if w13_weight is None:
                w13_weight = getattr(self.experts, "w13_weight_packed", None)
            w2_weight = getattr(self.experts, "w2_weight", None)
            if w2_weight is None:
                w2_weight = getattr(self.experts, "w2_weight_packed", None)
            if w13_weight is not None:
                w13_weight.data.zero_()
            if w2_weight is not None:
                w2_weight.data.zero_()
            self.experts.moe_config.intermediate_size_per_partition_unpadded = (
                moe_intermediate_size // self.tp_size
            )

    def _maybe_overlap_router_and_down_proj(
        self, hidden_states: torch.Tensor
    ) -> tuple[torch.Tensor, torch.Tensor, torch.Tensor | None]:
        """Compute the routed-expert down projection alongside the router,
        overlapping them on separate CUDA streams when latent MoE is enabled.

        The router gate and the down projection both read ``hidden_states``, so
        the gate runs on the default stream and the down projection on the aux
        stream, joined via ``maybe_execute_in_parallel``. For MegaMoE the
        grouped top-k selection consumes only the gate logits, so it also runs
        on the default stream and overlaps the down projection.

        Returns:
            ``(routed_hidden_states, router_output, topk_ids)``.
            ``routed_hidden_states`` is the down-projected latent (or the
            original ``hidden_states`` when latent MoE is disabled). For MegaMoE
            ``router_output`` holds the grouped top-k weights and ``topk_ids``
            the selected experts; otherwise ``router_output`` holds the raw gate
            logits and ``topk_ids`` is ``None``.
        """

        def _router(
            hidden_states: torch.Tensor,
        ) -> tuple[torch.Tensor, torch.Tensor | None]:
            router_logits, _ = self.gate(hidden_states)
            if not self.use_mega_moe:
                return router_logits, None
            return fused_grouped_topk(
                hidden_states=hidden_states,
                gating_output=router_logits,
                topk=self.experts.top_k,
                renormalize=self.moe_renormalize,
                e_score_correction_bias=self.gate.e_score_correction_bias.data,
                num_expert_group=self.num_expert_group,
                topk_group=self.topk_group,
                scoring_func=self.moe_router_activation_func,
                routed_scaling_factor=self.routed_scaling_factor,
            )

        down_proj = self.routed_expert_down_proj
        if down_proj is None:
            router_output, topk_ids = _router(hidden_states)
            return hidden_states, router_output, topk_ids
        num_tokens = hidden_states.shape[0]
        (router_output, topk_ids), (routed_hidden_states, _) = (
            maybe_execute_in_parallel(
                lambda: _router(hidden_states),
                lambda: down_proj(hidden_states),
                self._down_proj_events[0],
                self._down_proj_events[1],
                self._down_proj_stream
                if num_tokens <= envs.VLLM_ROUTED_DOWN_PROJ_STREAM_TOKEN_THRESHOLD
                else None,
            )
        )
        return routed_hidden_states, router_output, topk_ids

    def forward(self, hidden_states: torch.Tensor) -> torch.Tensor:
        num_tokens, hidden_size = hidden_states.shape
        hidden_states = hidden_states.view(-1, hidden_size)
        # Overlap the gate with the routed down projection; the returned hidden
        # states are already down-projected. Keep the original ``hidden_states``
        # for the shared experts.
        routed_hidden_states, router_output, topk_ids = (
            self._maybe_overlap_router_and_down_proj(hidden_states)
        )
        if self.use_mega_moe:
            assert self.routed_output_transform is not None
            assert topk_ids is not None
            final_hidden_states = self.experts(
                routed_hidden_states,
                router_output,
                topk_ids,
                activation_clamp=None,
            )
            final_hidden_states = self.routed_output_transform(final_hidden_states)
            if self.shared_experts is not None:
                final_hidden_states = final_hidden_states + self.shared_experts(
                    hidden_states
                )
        else:
            # Routed experts consume the down-projected latent; shared experts
            # (inside FusedMoE) get the original hidden states via
            # shared_experts_input.
            final_hidden_states = self.experts(
                hidden_states=routed_hidden_states,
                router_logits=router_output,
                shared_experts_input=hidden_states,
            )
        return final_hidden_states.view(num_tokens, hidden_size)

_maybe_overlap_router_and_down_proj(hidden_states)

Compute the routed-expert down projection alongside the router, overlapping them on separate CUDA streams when latent MoE is enabled.

The router gate and the down projection both read hidden_states, so the gate runs on the default stream and the down projection on the aux stream, joined via maybe_execute_in_parallel. For MegaMoE the grouped top-k selection consumes only the gate logits, so it also runs on the default stream and overlaps the down projection.

Returns:

  • Tensor

    (routed_hidden_states, router_output, topk_ids).

  • Tensor

    routed_hidden_states is the down-projected latent (or the

  • Tensor | None

    original hidden_states when latent MoE is disabled). For MegaMoE

  • tuple[Tensor, Tensor, Tensor | None]

    router_output holds the grouped top-k weights and topk_ids

  • tuple[Tensor, Tensor, Tensor | None]

    the selected experts; otherwise router_output holds the raw gate

  • tuple[Tensor, Tensor, Tensor | None]

    logits and topk_ids is None.

Source code in vllm/models/kimi_k3/nvidia/model.py
def _maybe_overlap_router_and_down_proj(
    self, hidden_states: torch.Tensor
) -> tuple[torch.Tensor, torch.Tensor, torch.Tensor | None]:
    """Compute the routed-expert down projection alongside the router,
    overlapping them on separate CUDA streams when latent MoE is enabled.

    The router gate and the down projection both read ``hidden_states``, so
    the gate runs on the default stream and the down projection on the aux
    stream, joined via ``maybe_execute_in_parallel``. For MegaMoE the
    grouped top-k selection consumes only the gate logits, so it also runs
    on the default stream and overlaps the down projection.

    Returns:
        ``(routed_hidden_states, router_output, topk_ids)``.
        ``routed_hidden_states`` is the down-projected latent (or the
        original ``hidden_states`` when latent MoE is disabled). For MegaMoE
        ``router_output`` holds the grouped top-k weights and ``topk_ids``
        the selected experts; otherwise ``router_output`` holds the raw gate
        logits and ``topk_ids`` is ``None``.
    """

    def _router(
        hidden_states: torch.Tensor,
    ) -> tuple[torch.Tensor, torch.Tensor | None]:
        router_logits, _ = self.gate(hidden_states)
        if not self.use_mega_moe:
            return router_logits, None
        return fused_grouped_topk(
            hidden_states=hidden_states,
            gating_output=router_logits,
            topk=self.experts.top_k,
            renormalize=self.moe_renormalize,
            e_score_correction_bias=self.gate.e_score_correction_bias.data,
            num_expert_group=self.num_expert_group,
            topk_group=self.topk_group,
            scoring_func=self.moe_router_activation_func,
            routed_scaling_factor=self.routed_scaling_factor,
        )

    down_proj = self.routed_expert_down_proj
    if down_proj is None:
        router_output, topk_ids = _router(hidden_states)
        return hidden_states, router_output, topk_ids
    num_tokens = hidden_states.shape[0]
    (router_output, topk_ids), (routed_hidden_states, _) = (
        maybe_execute_in_parallel(
            lambda: _router(hidden_states),
            lambda: down_proj(hidden_states),
            self._down_proj_events[0],
            self._down_proj_events[1],
            self._down_proj_stream
            if num_tokens <= envs.VLLM_ROUTED_DOWN_PROJ_STREAM_TOKEN_THRESHOLD
            else None,
        )
    )
    return routed_hidden_states, router_output, topk_ids