|
| 1 | +import unittest |
| 2 | +from unittest.mock import MagicMock, patch |
| 3 | + |
| 4 | +import torch |
| 5 | +from vllm.config import VllmConfig |
| 6 | +from vllm.worker.model_runner import ModelInputForGPUWithSamplingMetadata |
| 7 | + |
| 8 | +from vllm_ascend.distributed.kv_transfer.simple_buffer import SimpleBuffer |
| 9 | +from vllm_ascend.distributed.kv_transfer.simple_connector import \ |
| 10 | + SimpleConnector |
| 11 | +from vllm_ascend.distributed.kv_transfer.simple_pipe import SimplePipe |
| 12 | + |
| 13 | + |
| 14 | +class TestSimpleConnector(unittest.TestCase): |
| 15 | + |
| 16 | + def setUp(self): |
| 17 | + self.mock_pipe = MagicMock(spec=SimplePipe) |
| 18 | + self.mock_buffer = MagicMock(spec=SimpleBuffer) |
| 19 | + |
| 20 | + patcher = patch( |
| 21 | + 'vllm_ascend.distributed.kv_transfer.simple_buffer.SimpleBuffer') |
| 22 | + self.addCleanup(patcher.stop) |
| 23 | + self.MockSimpleBuffer = patcher.start() |
| 24 | + self.MockSimpleBuffer.return_value = self.mock_buffer |
| 25 | + |
| 26 | + def _create_mock_config(self, kv_role): |
| 27 | + mock_config = MagicMock() |
| 28 | + mock_config.kv_role = "kv_producer" |
| 29 | + mock_config.kv_connector_extra_config = { |
| 30 | + "prefill_device_ips": ["127.0.0.1"], |
| 31 | + "decode_device_ips": ["127.0.0.1"], |
| 32 | + "llmdatadist_comm_port": 26000, |
| 33 | + "http_port": 8000, |
| 34 | + "proxy_ip": "127.0.0.1", |
| 35 | + "proxy_port": "8000", |
| 36 | + "port": 5500 |
| 37 | + } |
| 38 | + mock_config.kv_port = 5500 |
| 39 | + self.mock_config = MagicMock(spec=VllmConfig) |
| 40 | + self.mock_config.kv_transfer_config.is_kv_producer = True |
| 41 | + self.mock_config.model_config.hf_config.hidden_size = 128 |
| 42 | + self.mock_config.model_config.hf_config.num_attention_heads = 8 |
| 43 | + self.mock_config.model_config.hf_config.num_key_value_heads = 8 |
| 44 | + self.mock_config.model_config.hf_config.qk_rope_head_dim = 16 |
| 45 | + self.mock_config.model_config.hf_config.kv_lora_rank = 16 |
| 46 | + self.mock_config.model_config.is_deepseek_mla = True |
| 47 | + # 模拟 parallel_config |
| 48 | + self.mock_config.parallel_config = MagicMock() |
| 49 | + self.mock_config.parallel_config.tensor_parallel_size = 1 |
| 50 | + self.mock_config.parallel_config.get_num_layers.return_value = 4 |
| 51 | + |
| 52 | + if kv_role == "kv_producer": |
| 53 | + self.mock_config.kv_transfer_config.kv_role = "kv_producer" |
| 54 | + else: |
| 55 | + self.mock_config.kv_transfer_config.kv_role = "kv_consumer" |
| 56 | + return mock_config |
| 57 | + |
| 58 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimplePipe') |
| 59 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimpleBuffer') |
| 60 | + @patch('llm_datadist.LLMDataDist') |
| 61 | + def test_select_init(self, mock_pipe, mock_buffer, MockLLMDataDist): |
| 62 | + """Test select method when buffer retrieval succeeds.""" |
| 63 | + connector = SimpleConnector( |
| 64 | + rank=0, |
| 65 | + local_rank=0, |
| 66 | + config=self._create_mock_config("kv_producer")) |
| 67 | + assert connector.producer_data_pipe is not None |
| 68 | + assert connector.producer_buffer is not None |
| 69 | + mock_data_dist = MockLLMDataDist.return_value |
| 70 | + mock_data_dist.init.return_value = None |
| 71 | + |
| 72 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimplePipe') |
| 73 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimpleBuffer') |
| 74 | + @patch('llm_datadist.LLMDataDist') |
| 75 | + def test_select_select(self, mock_pipe, mock_buffer, MockLLMDataDist): |
| 76 | + |
| 77 | + connector = SimpleConnector( |
| 78 | + rank=0, |
| 79 | + local_rank=0, |
| 80 | + config=self._create_mock_config("kv_consumer")) |
| 81 | + connector.consumer_data_pipe = mock_pipe |
| 82 | + connector.consumer_buffer = mock_buffer |
| 83 | + assert connector.consumer_data_pipe is not None |
| 84 | + assert connector.consumer_buffer is not None |
| 85 | + input_tokens = torch.tensor([1, 2, 3]) |
| 86 | + roi = torch.tensor([True, True, True]) |
| 87 | + req_id = "test_req" |
| 88 | + connector.select(input_tokens, roi, req_id) |
| 89 | + |
| 90 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimplePipe') |
| 91 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimpleBuffer') |
| 92 | + @patch('llm_datadist.LLMDataDist') |
| 93 | + def test_insert(self, mock_pipe, mock_buffer, MockLLMDataDist): |
| 94 | + """Test insert operation""" |
| 95 | + connector = SimpleConnector( |
| 96 | + rank=0, |
| 97 | + local_rank=0, |
| 98 | + config=self._create_mock_config("kv_producer")) |
| 99 | + |
| 100 | + connector.producer_buffer = mock_buffer |
| 101 | + |
| 102 | + input_tokens = torch.randint(0, 1000, (5, )) |
| 103 | + roi = torch.ones_like(input_tokens, dtype=torch.bool) |
| 104 | + keys = torch.randn(3, 5, 1, 96) |
| 105 | + values = torch.randn(3, 5, 1, 96) |
| 106 | + hidden = torch.randn(5, 768) |
| 107 | + req_id = "test_req" |
| 108 | + |
| 109 | + connector.insert(input_tokens, roi, keys, values, hidden, req_id) |
| 110 | + |
| 111 | + mock_buffer.insert.assert_called_once_with(input_tokens, roi, keys, |
| 112 | + values, hidden, req_id) |
| 113 | + |
| 114 | + @patch.object(SimpleConnector, 'insert') |
| 115 | + @patch('torch.distributed.get_rank', return_value=0) |
| 116 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimplePipe') |
| 117 | + @patch('vllm_ascend.distributed.kv_transfer.simple_connector.SimpleBuffer') |
| 118 | + @patch('llm_datadist.LLMDataDist') |
| 119 | + def test_send_kv_caches_and_hidden_states(self, mock_pipe, mock_buffer, |
| 120 | + MockLLMDataDist, mock_insert, |
| 121 | + mock_rank): |
| 122 | + """Test sending KV caches and hidden states""" |
| 123 | + connector = SimpleConnector( |
| 124 | + rank=0, |
| 125 | + local_rank=0, |
| 126 | + config=self._create_mock_config("kv_producer")) |
| 127 | + |
| 128 | + mock_model_executable = MagicMock() |
| 129 | + mock_model_executable.model.start_layer = 0 |
| 130 | + mock_model_executable.model.end_layer = 3 |
| 131 | + |
| 132 | + mock_model_input = MagicMock(spec=ModelInputForGPUWithSamplingMetadata) |
| 133 | + mock_model_input.input_tokens = torch.randint(0, 1000, (10, )) |
| 134 | + mock_model_input.attn_metadata.seq_lens = [5, 5] |
| 135 | + mock_model_input.attn_metadata.slot_mapping = torch.randint( |
| 136 | + 0, 100, (10, )) |
| 137 | + mock_model_input.attn_metadata.num_prefill_tokens = 10 |
| 138 | + mock_model_input.request_ids_to_seq_ids = {"req1": [0], "req2": [1]} |
| 139 | + |
| 140 | + kv_caches = [torch.randn(2, 100, 1, 96) for _ in range(3)] |
| 141 | + |
| 142 | + hidden_states = torch.randn(10, 768) |
| 143 | + |
| 144 | + connector.send_kv_caches_and_hidden_states(mock_model_executable, |
| 145 | + mock_model_input, kv_caches, |
| 146 | + hidden_states) |
0 commit comments