Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
48 commits
Select commit Hold shift + click to select a range
9a4fb8f
add mooncake connector v2 base class
nwpu-zxr Jul 30, 2026
70f89c9
add base scheduler worker
nwpu-zxr Jul 30, 2026
b6befac
add base worker
nwpu-zxr Jul 31, 2026
18948f6
add basic worker
nwpu-zxr Aug 2, 2026
ccbcb67
add register_kv_caches
nwpu-zxr Aug 3, 2026
c1d0327
add scheduler send/recv thread
nwpu-zxr Aug 3, 2026
9b8b847
add full block ids
nwpu-zxr Aug 3, 2026
918bc2e
add full-pp transfermetadata
nwpu-zxr Aug 3, 2026
1b246d5
add worker basic function
nwpu-zxr Aug 4, 2026
08b07b8
add get_remote_metadata
nwpu-zxr Aug 4, 2026
6186c97
add transfer metadata
nwpu-zxr Aug 5, 2026
36e87c7
add merge transfermetadata
nwpu-zxr Aug 6, 2026
bddcaec
simplify metadata
nwpu-zxr Aug 6, 2026
12cd4d6
fix ruff format
nwpu-zxr Aug 6, 2026
b867c23
add kernel block sizes
nwpu-zxr Aug 7, 2026
f8deee9
add remote_tp_layout
nwpu-zxr Aug 7, 2026
25e3744
change metadatagroups
nwpu-zxr Aug 8, 2026
414aea8
add layerwise generate tp_layout
nwpu-zxr Aug 8, 2026
75e338c
add compute block ids
nwpu-zxr Aug 8, 2026
013a439
add dcp=1 block ids compute
nwpu-zxr Aug 9, 2026
77d3cf8
add dcp>1 block ids compute
nwpu-zxr Aug 9, 2026
17d43c6
add compute address and transfer
nwpu-zxr Aug 9, 2026
c2dcc20
add mamba address and transfer
nwpu-zxr Aug 9, 2026
ab6c0e4
add dcp with different block size
nwpu-zxr Aug 10, 2026
2cfcee3
merge dcp with same/different block size code
nwpu-zxr Aug 10, 2026
24572e9
add threadpool
nwpu-zxr Aug 10, 2026
6414193
some bugfix
nwpu-zxr Aug 12, 2026
d94f206
change connector path
nwpu-zxr Aug 12, 2026
0d87a1b
add __init__
nwpu-zxr Aug 12, 2026
f5c182d
fix import
nwpu-zxr Aug 12, 2026
351d6c0
fix ruff
nwpu-zxr Aug 13, 2026
489bfcb
fix ruff
nwpu-zxr Aug 16, 2026
cab1d3b
fix mypy
nwpu-zxr Aug 21, 2026
5d54a5d
add reuseable zmq_context for pull scheduler
nwpu-zxr Aug 21, 2026
1bd6553
add mooncake connector v2 UT
nwpu-zxr Aug 21, 2026
e414ca6
fix UT
nwpu-zxr Aug 21, 2026
f000e29
add more UT
nwpu-zxr Aug 21, 2026
1f971c0
fix ut ruff
nwpu-zxr Aug 21, 2026
930c862
add base class UT
nwpu-zxr Aug 21, 2026
cbbb4e3
support kvpp
nwpu-zxr Aug 22, 2026
da5aa12
add comments & fix mypy
nwpu-zxr Aug 22, 2026
b9ba3e4
support p-kvpp with d-dcp
nwpu-zxr Aug 22, 2026
30b6753
fix indexerspec block size
nwpu-zxr Aug 28, 2026
a8af7c4
add block merge
nwpu-zxr Sep 7, 2026
d07686a
fix ruff
nwpu-zxr Sep 7, 2026
3bd67f3
chore(mooncake): add connector V2 runtime logs
wangxiaoteng888 Sep 7, 2026
5f12fba
support new kv_cache_tensors
nwpu-zxr Sep 9, 2026
7de7708
test(mooncake): patch module objects to avoid CPU suite import failures
wangxiaoteng888 Sep 9, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions tests/ut/kv_offload/mooncake_v2/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
"""Unit tests for the refactored Mooncake connector."""
177 changes: 177 additions & 0 deletions tests/ut/kv_offload/mooncake_v2/helpers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
# SPDX-License-Identifier: Apache-2.0
"""Shared builders for refactored Mooncake connector tests."""

from types import SimpleNamespace
from typing import Any
from unittest.mock import MagicMock

import torch
from vllm.v1.kv_cache_interface import FullAttentionSpec, MambaSpec, SlidingWindowSpec

from vllm_ascend.core.kv_cache_interface import AscendSFAIndexerCacheSpec
from vllm_ascend.distributed.kv_transfer.kv_p2p.mooncake.metadata import (
MooncakePPTransferMetadata,
MooncakeTPTransferMetadata,
MooncakeTransferMetadata,
MooncakeTransferMetadataGroups,
)


def make_full_spec(block_size: int = 16, num_kv_heads: int = 1) -> FullAttentionSpec:
return FullAttentionSpec(
block_size=block_size,
num_kv_heads=num_kv_heads,
head_size=8,
head_size_v=8,
dtype=torch.float16,
)


def make_sliding_spec(block_size: int = 16) -> SlidingWindowSpec:
return SlidingWindowSpec(
block_size=block_size,
num_kv_heads=1,
head_size=8,
dtype=torch.float16,
sliding_window=64,
)


def make_mamba_spec(block_size: int = 16) -> MambaSpec:
return MambaSpec(
block_size=block_size,
shapes=((3, 16), (2, 4, 4)),
dtypes=(torch.float16, torch.float16),
)


def make_sfa_indexer_spec(
block_size: int = 16,
replication_size: int = 2,
) -> AscendSFAIndexerCacheSpec:
return AscendSFAIndexerCacheSpec(
block_size=block_size,
num_kv_heads=1,
head_size=8,
dtype=torch.float16,
sfa_dcp_replicated_indexer_size=replication_size,
)


def make_transfer_metadata(
*,
engine_id: str = "engine-p",
te_rpc_port: int = 9000,
local_ip: str = "10.0.0.1",
handshake_port: int = 5000,
layer_names: list[str] | None = None,
group_indices: list[int] | None = None,
layer_block_sizes: list[int] | None = None,
base_addrs: list[list[int]] | None = None,
block_strides: list[list[int]] | None = None,
block_lens: list[list[int]] | None = None,
block_shapes: list[list[tuple[int, ...]]] | None = None,
block_size_scales: list[list[int]] | None = None,
) -> MooncakeTransferMetadata:
layer_names = layer_names or ["model.layers.0.self_attn"]
num_layers = len(layer_names)
return MooncakeTransferMetadata(
engine_id=engine_id,
te_rpc_port=te_rpc_port,
block_size=16,
num_blocks=32,
layer_names=layer_names,
layer_block_sizes=layer_block_sizes or [16] * num_layers,
group_indices=group_indices or [0] * num_layers,
kv_caches_base_addr=base_addrs or [[1000 + index * 1000] for index in range(num_layers)],
block_strides=block_strides or [[128] for _ in range(num_layers)],
block_lens=block_lens or [[128] for _ in range(num_layers)],
block_shapes=block_shapes or [[(1, 16, 4)] for _ in range(num_layers)],
block_size_scales=block_size_scales or [[1] for _ in range(num_layers)],
local_ip=local_ip,
handshake_port=handshake_port,
)


def make_pp_metadata(
*,
layer_names: list[str] | None = None,
layer_block_sizes: list[int] | None = None,
block_shapes: list[list[tuple[int, ...]]] | None = None,
block_strides: list[list[int]] | None = None,
block_lens: list[list[int]] | None = None,
block_size_scales: list[list[int]] | None = None,
tp_base_addrs: dict[int, list[list[int]]] | None = None,
tp_layer_indices: dict[int, list[int]] | None = None,
) -> MooncakePPTransferMetadata:
layer_names = layer_names or ["model.layers.0.self_attn"]
num_layers = len(layer_names)
tp_base_addrs = tp_base_addrs or {0: [[5000 + index * 1000] for index in range(num_layers)]}
tp_layer_indices = tp_layer_indices or {}
return MooncakePPTransferMetadata(
block_size=16,
num_blocks=32,
layer_names=layer_names,
layer_block_sizes=layer_block_sizes or [16] * num_layers,
group_indices=[0] * num_layers,
block_strides=block_strides or [[128] for _ in range(num_layers)],
block_lens=block_lens or [[128] for _ in range(num_layers)],
block_shapes=block_shapes or [[(1, 16, 4)] for _ in range(num_layers)],
block_size_scales=block_size_scales or [[1] for _ in range(num_layers)],
metadata_by_tp_rank={
tp_rank: MooncakeTPTransferMetadata(
te_rpc_port=9000 + tp_rank,
layer_indices=tp_layer_indices.get(tp_rank, list(range(num_layers))),
kv_caches_base_addr=base_addrs,
local_ip=f"10.0.0.{tp_rank + 1}",
handshake_port=5000 + tp_rank,
)
for tp_rank, base_addrs in tp_base_addrs.items()
},
)


def make_metadata_groups(
*,
engine_id: str = "engine-p",
tp_size: int = 1,
use_kv_pp: bool = False,
pp_metadata: MooncakePPTransferMetadata | None = None,
) -> MooncakeTransferMetadataGroups:
return MooncakeTransferMetadataGroups(
engine_id=engine_id,
scheduler_host="10.0.0.10",
scheduler_port=6000,
pp_size=1,
pcp_size=1,
dcp_size=1,
tp_size=tp_size,
use_kv_pp=use_kv_pp,
metadata_by_pp_rank={0: pp_metadata or make_pp_metadata()},
)


def make_request(**overrides: Any) -> SimpleNamespace:
values: dict[str, Any] = {
"request_id": "request-0",
"prompt_token_ids": list(range(32)),
"prompt_embeds": None,
"num_prompt_tokens": 32,
"kv_transfer_params": {},
"status": "running",
"output_token_ids": [100],
"_all_token_ids": list(range(32)),
"max_tokens": 16,
}
values.update(overrides)
return SimpleNamespace(**values)


def make_blocks(
unhashed: tuple[list[int], ...] = ([10, 11],),
full: tuple[list[int], ...] = ([1, 2, 10, 11],),
) -> MagicMock:
blocks = MagicMock()
blocks.get_unhashed_block_ids_all_groups.return_value = unhashed
blocks.get_block_ids.return_value = full
return blocks
162 changes: 162 additions & 0 deletions tests/ut/kv_offload/mooncake_v2/test_base_scheduler.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
# SPDX-License-Identifier: Apache-2.0

from types import SimpleNamespace
from unittest.mock import MagicMock

import pytest
from vllm.v1.kv_cache_interface import KVCacheGroupSpec, UniformTypeKVCacheSpecs

from vllm_ascend.distributed.kv_transfer.kv_p2p.mooncake import base_scheduler
from vllm_ascend.distributed.kv_transfer.kv_p2p.mooncake.base_scheduler import (
MooncakeBaseConnectorScheduler,
)

from .helpers import make_full_spec, make_mamba_spec, make_sliding_spec


def make_scheduler_config(
*,
is_consumer: bool = True,
is_producer: bool = False,
pcp_size: int = 1,
) -> SimpleNamespace:
return SimpleNamespace(
kv_transfer_config=SimpleNamespace(
is_kv_consumer=is_consumer,
is_kv_producer=is_producer,
kv_role="kv_consumer" if is_consumer else "kv_producer",
kv_port=6000,
),
cache_config=SimpleNamespace(block_size=16),
speculative_config=SimpleNamespace(num_speculative_tokens=3),
parallel_config=SimpleNamespace(
pipeline_parallel_size=2,
tensor_parallel_size=4,
prefill_context_parallel_size=pcp_size,
decode_context_parallel_size=2,
data_parallel_size=3,
data_parallel_rank=1,
),
model_config=SimpleNamespace(hf_config=SimpleNamespace(compress_ratios=None)),
)


def test_base_scheduler_initializes_parallel_layout_and_control_port(
monkeypatch: pytest.MonkeyPatch,
) -> None:
config = make_scheduler_config()
group = KVCacheGroupSpec(layer_names=["layer.0"], kv_cache_spec=make_full_spec())
kv_cache_config = SimpleNamespace(kv_cache_groups=[group])
ascend_config = object()
init_ascend_config = MagicMock()
monkeypatch.setattr(
base_scheduler,
"init_ascend_config",
init_ascend_config,
)
monkeypatch.setattr(
base_scheduler,
"get_ascend_config",
MagicMock(return_value=ascend_config),
)
monkeypatch.setattr(
base_scheduler,
"get_ip",
MagicMock(return_value="10.0.0.1"),
)

scheduler = MooncakeBaseConnectorScheduler(config, "engine-d", kv_cache_config) # type: ignore[arg-type]

assert scheduler.engine_id == "engine-d"
assert scheduler.block_size == 16
assert scheduler.num_speculative_tokens == 3
assert scheduler.ascend_config is ascend_config
assert scheduler.side_channel_host == "10.0.0.1"
assert scheduler.max_device_id == 24
assert scheduler.side_channel_port == 6025
assert scheduler.group_block_size == [16]
assert scheduler.group_unique_specs == [[group.kv_cache_spec]]
assert scheduler.need_truncate is False
init_ascend_config.assert_called_once_with(config)


@pytest.mark.parametrize(("is_consumer", "is_producer"), [(False, False), (True, True)])
def test_base_scheduler_rejects_invalid_transfer_roles(
is_consumer: bool,
is_producer: bool,
) -> None:
config = make_scheduler_config(is_consumer=is_consumer, is_producer=is_producer)

with pytest.raises(ValueError, match="exactly one KV transfer role"):
MooncakeBaseConnectorScheduler(config, "engine", SimpleNamespace()) # type: ignore[arg-type]


def test_base_scheduler_rejects_unsupported_pcp(monkeypatch: pytest.MonkeyPatch) -> None:
config = make_scheduler_config(pcp_size=2)
monkeypatch.setattr(
base_scheduler,
"init_ascend_config",
MagicMock(),
)
monkeypatch.setattr(
base_scheduler,
"get_ascend_config",
MagicMock(),
)
monkeypatch.setattr(
base_scheduler,
"get_ip",
MagicMock(return_value="10.0.0.1"),
)

with pytest.raises(AssertionError, match="prefill context parallel size 1"):
MooncakeBaseConnectorScheduler(config, "engine", SimpleNamespace()) # type: ignore[arg-type]


def test_base_scheduler_expands_unique_uniform_specs_in_layer_order() -> None:
full = make_full_spec()
sliding = make_sliding_spec()
uniform = UniformTypeKVCacheSpecs(
block_size=16,
kv_cache_specs={"layer.0": full, "layer.1": full, "layer.2": sliding},
)
group = KVCacheGroupSpec(
layer_names=["layer.1", "layer.0", "layer.2"],
kv_cache_spec=uniform,
)

assert MooncakeBaseConnectorScheduler._get_group_unique_specs(group) == [full, sliding]


def test_base_scheduler_transfer_blocks_handles_empty_and_state_without_speculation() -> None:
scheduler = MooncakeBaseConnectorScheduler.__new__(MooncakeBaseConnectorScheduler)
scheduler.pcp_size = 1
scheduler.dcp_size = 1
scheduler.num_speculative_tokens = 0
scheduler.group_block_size = [16]
scheduler.group_unique_specs = [[make_mamba_spec()]]

assert scheduler._get_transfer_block_ids((), prompt_len=32) == ()
assert scheduler._get_transfer_block_ids(([10, 11],), prompt_len=32) == ([10, 11],)


def test_base_scheduler_abstract_contract_and_legacy_metadata_delegation() -> None:
scheduler = MooncakeBaseConnectorScheduler.__new__(MooncakeBaseConnectorScheduler)

with pytest.raises(NotImplementedError):
scheduler.on_new_request(MagicMock())
with pytest.raises(NotImplementedError):
scheduler.update_connector_output(MagicMock())
with pytest.raises(NotImplementedError):
scheduler.get_num_new_matched_tokens(MagicMock(), 0)
with pytest.raises(NotImplementedError):
scheduler.update_state_after_alloc(MagicMock(), MagicMock(), 0)
with pytest.raises(NotImplementedError):
scheduler.build_connector_meta(MagicMock())
with pytest.raises(NotImplementedError):
scheduler.request_finished(MagicMock(), ())

scheduler.set_xfer_handshake_metadata_from_workers = MagicMock() # type: ignore[method-assign]
metadata = {0: MagicMock()}
scheduler.set_xfer_handshake_metadata(metadata)
scheduler.set_xfer_handshake_metadata_from_workers.assert_called_once_with(metadata)
Loading
Loading