File size: 14,711 Bytes
6e668dc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project

"""
Integration tests for MultiprocExecutor at the executor level.
This test directly tests the executor without going through the LLM interface,
focusing on executor initialization, RPC calls, and distributed execution.
"""

import multiprocessing
import os
import socket

from tests.utils import multi_gpu_test
from vllm.config import VllmConfig
from vllm.engine.arg_utils import EngineArgs
from vllm.v1.core.sched.output import SchedulerOutput
from vllm.v1.executor.multiproc_executor import MultiprocExecutor

MODEL = "facebook/opt-125m"


def create_vllm_config(
    tensor_parallel_size: int = 1,
    pipeline_parallel_size: int = 1,
    max_model_len: int = 256,
    gpu_memory_utilization: float = 0.3,
    distributed_executor_backend: str = "mp",
    nnodes: int = 1,
    node_rank: int = 0,
    master_port: int = 0,
) -> VllmConfig:
    """Create a VllmConfig for testing using EngineArgs."""
    engine_args = EngineArgs(
        model=MODEL,
        tensor_parallel_size=tensor_parallel_size,
        pipeline_parallel_size=pipeline_parallel_size,
        max_model_len=max_model_len,
        gpu_memory_utilization=gpu_memory_utilization,
        distributed_executor_backend=distributed_executor_backend,
        enforce_eager=True,
    )
    vllm_config = engine_args.create_engine_config()

    # Override distributed node settings if needed
    if nnodes > 1 or node_rank > 0:
        vllm_config.parallel_config.nnodes = nnodes
        vllm_config.parallel_config.node_rank = node_rank
        vllm_config.parallel_config.master_port = master_port
    if nnodes > 1:
        vllm_config.parallel_config.disable_custom_all_reduce = True

    return vllm_config


def create_test_scheduler_output(num_requests: int = 1) -> SchedulerOutput:
    """Create a minimal SchedulerOutput for testing."""
    # This is a simplified version - in practice you'd need proper
    # SchedulerOutput construction based on the actual vLLM v1 API
    return SchedulerOutput(
        scheduled_new_reqs=[],
        scheduled_resumed_reqs=[],
        scheduled_running_reqs=[],
        num_scheduled_tokens={},
        total_num_scheduled_tokens=0,
    )


def test_multiproc_executor_initialization():
    """Test that MultiprocExecutor can be initialized with proper config."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=1,
        pipeline_parallel_size=1,
    )

    # Create executor - this should initialize workers
    executor = MultiprocExecutor(vllm_config=vllm_config)

    # Verify executor properties
    assert executor.world_size == 1, "World size should be 1 for single GPU"
    assert executor.local_world_size == 1, "Local world size should be 1"
    assert hasattr(executor, "workers"), "Executor should have workers"
    assert len(executor.workers) == 1, "Should have 1 worker for single GPU"

    # Clean up
    executor.shutdown()


@multi_gpu_test(num_gpus=2)
def test_multiproc_executor_initialization_tensor_parallel():
    """Test MultiprocExecutor initialization with tensor parallelism."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=2,
        pipeline_parallel_size=1,
    )

    # Create executor
    executor = MultiprocExecutor(vllm_config=vllm_config)

    # Verify executor properties
    assert executor.world_size == 2, "World size should be 2 for TP=2"
    assert executor.local_world_size == 2, "Local world size should be 2"
    assert len(executor.workers) == 2, "Should have 2 workers for TP=2"

    # Verify output rank calculation
    output_rank = executor._get_output_rank()
    assert output_rank == 0, "Output rank should be 0 for TP=2, PP=1"

    # Clean up
    executor.shutdown()


@multi_gpu_test(num_gpus=2)
def test_multiproc_executor_collective_rpc():
    """Test collective RPC calls to all workers."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=2,
        pipeline_parallel_size=1,
    )

    # Create executor
    executor = MultiprocExecutor(vllm_config=vllm_config)

    try:
        # Test check_health RPC - should work without errors
        executor.check_health()

        # Test that RPC works correctly
        # Note: We're just testing that the RPC mechanism works,
        # not testing actual model execution here
        assert not executor.is_failed, "Executor should not be in failed state"

    finally:
        # Clean up
        executor.shutdown()


def test_multiproc_executor_failure_callback():
    """Test failure callback registration and invocation."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=1,
        pipeline_parallel_size=1,
    )

    executor = MultiprocExecutor(vllm_config=vllm_config)

    try:
        # Test callback registration
        callback_invoked = []

        def test_callback():
            callback_invoked.append(True)

        # Register callback
        executor.register_failure_callback(test_callback)

        # Callback should not be invoked yet
        assert len(callback_invoked) == 0, "Callback should not be invoked immediately"

        # Simulate failure
        executor.is_failed = True

        # Register another callback - should be invoked immediately
        executor.register_failure_callback(test_callback)
        assert len(callback_invoked) == 1, (
            "Callback should be invoked when executor is failed"
        )

    finally:
        # Clean up
        executor.shutdown()


@multi_gpu_test(num_gpus=2)
def test_multiproc_executor_worker_monitor():
    """Test that worker monitor is set up correctly."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=2,
        pipeline_parallel_size=1,
    )

    executor = MultiprocExecutor(vllm_config=vllm_config)

    try:
        # Verify all worker processes are alive
        for worker in executor.workers:
            assert worker.proc.is_alive(), f"Worker rank {worker.rank} should be alive"

        # Verify executor is not in failed state
        assert not executor.is_failed, "Executor should not be in failed state"

    finally:
        # Clean up
        executor.shutdown()

        # After shutdown, workers should be terminated
        import time

        time.sleep(0.5)  # Give processes time to terminate
        for worker in executor.workers:
            assert not worker.proc.is_alive(), (
                f"Worker rank {worker.rank} should terminate after shutdown"
            )


@multi_gpu_test(num_gpus=2)
def test_multiproc_executor_get_response_message_queues():
    """Test message queue retrieval for different ranks."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=2,
        pipeline_parallel_size=1,
    )

    executor = MultiprocExecutor(vllm_config=vllm_config)

    try:
        # Get all message queues
        all_queues = executor.get_response_mqs()
        assert len(all_queues) == 2, "Should have 2 message queues for 2 workers"

        # Get message queue for specific rank
        rank0_queue = executor.get_response_mqs(unique_reply_rank=0)
        assert len(rank0_queue) == 1, "Should have 1 message queue for rank 0"

        rank1_queue = executor.get_response_mqs(unique_reply_rank=1)
        assert len(rank1_queue) == 1, "Should have 1 message queue for rank 1"

    finally:
        # Clean up
        executor.shutdown()


def test_multiproc_executor_shutdown_cleanup():
    """Test that shutdown properly cleans up resources."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=1,
        pipeline_parallel_size=1,
    )

    executor = MultiprocExecutor(vllm_config=vllm_config)

    # Verify executor is set up
    assert hasattr(executor, "workers"), "Executor should have workers"
    assert len(executor.workers) > 0, "Should have at least one worker"

    # Shutdown
    executor.shutdown()

    # Verify cleanup
    import time

    time.sleep(0.5)  # Give processes time to terminate

    for worker in executor.workers:
        assert not worker.proc.is_alive(), "Worker processes should be terminated"

    # Verify shutdown event is set
    assert executor.shutdown_event.is_set(), "Shutdown event should be set"

    # Multiple shutdowns should be safe (idempotent)
    executor.shutdown()
    executor.shutdown()


@multi_gpu_test(num_gpus=4)
def test_multiproc_executor_pipeline_parallel():
    """Test MultiprocExecutor with pipeline parallelism."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=2,
        pipeline_parallel_size=2,
    )

    executor = MultiprocExecutor(vllm_config=vllm_config)

    try:
        # Verify executor properties
        assert executor.world_size == 4, "World size should be 4 for TP=2, PP=2"
        assert len(executor.workers) == 4, "Should have 4 workers"

        # Verify output rank calculation
        # For TP=2, PP=2: output should be from the last PP stage (ranks 2-3)
        # Specifically rank 2 (first rank of last PP stage)
        output_rank = executor._get_output_rank()
        assert output_rank == 2, "Output rank should be 2 (first rank of last PP stage)"

        # Verify max_concurrent_batches for pipeline parallel
        assert executor.max_concurrent_batches == 2, (
            "Max concurrent batches should equal PP size"
        )

    finally:
        # Clean up
        executor.shutdown()


def test_multiproc_executor_properties():
    """Test various executor properties and configurations."""
    vllm_config = create_vllm_config(
        tensor_parallel_size=1,
        pipeline_parallel_size=1,
    )

    executor = MultiprocExecutor(vllm_config=vllm_config)

    try:
        # Test supports_pp property
        assert MultiprocExecutor.supports_pp is True, (
            "MultiprocExecutor should support pipeline parallelism"
        )

        # Test world_size calculation
        assert executor.world_size == (
            executor.parallel_config.tensor_parallel_size
            * executor.parallel_config.pipeline_parallel_size
        ), "World size should equal TP * PP"

        # Test local_world_size calculation
        assert executor.local_world_size == (
            executor.parallel_config.world_size // executor.parallel_config.nnodes
        ), "Local world size should be world_size / nnodes"

    finally:
        # Clean up
        executor.shutdown()


@multi_gpu_test(num_gpus=4)
def test_multiproc_executor_multi_node():
    """
    Test MultiprocExecutor with multi-node configuration.
    This simulates 2 nodes with TP=4:
    - Node 0 (rank 0): Uses GPUs 0,1 (CUDA_VISIBLE_DEVICES=0,1) with TP=2
    - Node 1 (rank 1): Uses GPUs 2,3 (CUDA_VISIBLE_DEVICES=2,3) with TP=2
    Total world_size = 4, nnodes = 2
    """
    with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
        s.bind(("", 0))
        port = s.getsockname()[1]
    # symm_mem does not work for simulating multi instance in single node
    os.environ["VLLM_ALLREDUCE_USE_SYMM_MEM"] = "0"

    def run_node(node_rank: int, result_queue: multiprocessing.Queue, port: int):
        """Run a single node's executor."""
        executor = None
        try:
            # Set CUDA_VISIBLE_DEVICES for this node
            if node_rank == 0:
                os.environ["CUDA_VISIBLE_DEVICES"] = "0,1"
            else:
                os.environ["CUDA_VISIBLE_DEVICES"] = "2,3"

            # Create config for this node
            vllm_config = create_vllm_config(
                tensor_parallel_size=4,  # Total TP across all nodes
                pipeline_parallel_size=1,
                nnodes=2,  # 2 nodes
                node_rank=node_rank,
                master_port=port,  # same port
            )

            # Create executor for this node
            executor = MultiprocExecutor(vllm_config=vllm_config)

            # Verify node-specific properties
            assert executor.world_size == 4, (
                f"World size should be 4 on node {node_rank}"
            )
            assert executor.local_world_size == 2, (
                f"Local world size should be 2 on node {node_rank}"
            )
            assert len(executor.workers) == 2, (
                f"Should have 2 local workers on node {node_rank}"
            )

            # Verify worker ranks are correct for this node
            expected_ranks = [node_rank * 2, node_rank * 2 + 1]
            actual_ranks = sorted([w.rank for w in executor.workers])
            assert actual_ranks == expected_ranks, (
                f"Node {node_rank} should have workers "
                f"with ranks {expected_ranks}, got {actual_ranks}"
            )
            # Verify all workers are alive
            for worker in executor.workers:
                assert worker.proc.is_alive(), (
                    f"Worker rank {worker.rank} should be alive on node {node_rank}"
                )
            # executor.gen
            # Put success result in queue BEFORE shutdown to avoid hanging
            result_queue.put({"node": node_rank, "success": True})
            import time

            time.sleep(2)
            executor.shutdown()
        except Exception as e:
            # Put failure result in queue
            result_queue.put({"node": node_rank, "success": False, "error": str(e)})
            raise e
        finally:
            if executor is not None:
                executor.shutdown()

    # Create a queue to collect results from both processes
    result_queue: multiprocessing.Queue[dict[str, int | bool]] = multiprocessing.Queue()

    # Start both node processes
    processes = []
    for node_rank in range(2):
        p = multiprocessing.Process(
            target=run_node,
            args=(node_rank, result_queue, port),
            name=f"Node{node_rank}",
        )
        p.start()
        processes.append(p)

    # Wait for both processes to complete
    all_completed = True
    for p in processes:
        p.join(timeout=60)
        if p.is_alive():
            p.terminate()
            p.join(timeout=20)
            if p.is_alive():
                p.kill()
                p.join()
            all_completed = False

    # Check results from both nodes
    results: list[dict[str, int | bool]] = []
    while len(results) < 2:
        try:
            result = result_queue.get(timeout=1)
            results.append(result)
        except Exception:
            pass
    assert all_completed, "Not all processes completed successfully"
    assert len(results) == 2, f"Expected 2 results, got {len(results)}"
    assert results[0]["success"], f"Node 0 failed: {results[0]}"
    assert results[1]["success"], f"Node 1 failed: {results[1]}"