test_router_e2e_with_vllm.py 23 KB
Newer Older
1
# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2
# SPDX-License-Identifier: Apache-2.0
3

4
5
import asyncio

6
7
8
9
# Timing notes (measured locally):
# - GPU-1 subset (`-m "gpu_1 and not gpu_2"`): 130.43s total for 3 tests.
# These tests load a real model and can be slow/flaky when GPU resources are contended,
# so we set explicit pytest timeouts to fail fast on hangs (see per-test markers below).
10
import json
11
12
13
14
import logging
import os
from typing import Any, Dict, Optional

15
import aiohttp
16
17
import pytest

18
19
20
21
22
from tests.router.e2e_harness import (
    ManagedEngineProcessMixin,
    run_basic_router_test,
    run_indexers_sync_test,
    run_router_decisions_test,
23
)
24
25
26
27
28
from tests.router.helper import (
    generate_random_suffix,
    get_kv_indexer_command,
    wait_for_indexer_workers_active,
)
29
from tests.utils.constants import DefaultPort
30
from tests.utils.managed_process import ManagedProcess
31
from tests.utils.port_utils import allocate_ports, deallocate_ports
32
33
34
35

logger = logging.getLogger(__name__)

MODEL_NAME = "TinyLlama/TinyLlama-1.1B-Chat-v1.0"
36
37
38

pytestmark = [
    pytest.mark.e2e,
39
    pytest.mark.router,
40
41
42
    pytest.mark.vllm,
    pytest.mark.model(MODEL_NAME),
]
43
44
45
SPEEDUP_RATIO = 10.0
BLOCK_SIZE = 16

46
47
48
49
50
51
52
53
54
55
# Shared vLLM configuration for all tests
# gpu_memory_utilization limits actual VRAM allocation (required for multi-worker on same GPU)
VLLM_ARGS: Dict[str, Any] = {
    "block_size": BLOCK_SIZE,
    "model": MODEL_NAME,
    "gpu_memory_utilization": 0.4,  # Limit VRAM allocation per worker
    "max_model_len": 1024,  # Limit context length to reduce KV cache size
    "enforce_eager": True,  # Disable CUDA graphs for faster startup & lower memory
}

56
57
58
59
60
61
62
VLLM_ARGS_NO_BLOCK_SIZE: Dict[str, Any] = {
    "model": MODEL_NAME,
    "gpu_memory_utilization": 0.4,  # Limit VRAM allocation per worker
    "max_model_len": 1024,  # Limit context length to reduce KV cache size
    "enforce_eager": True,  # Disable CUDA graphs for faster startup & lower memory
}

63

64
class VLLMProcess(ManagedEngineProcessMixin):
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
    """Manages vLLM workers using dynamo.vllm (HTTP API + KV events).

    This is a drop-in replacement for MockerProcess that uses real vLLM workers.
    The key difference: dynamo.vllm automatically handles:
    - HTTP API serving
    - KV cache event publishing (ZMQ → NATS bridge)
    - Integration with dynamo.frontend router
    """

    def __init__(
        self,
        request,
        vllm_args: Optional[Dict[str, Any]] = None,
        num_workers: int = 2,
        single_gpu: bool = False,
        data_parallel_size: Optional[int] = None,
81
82
        request_plane: str = "tcp",
        store_backend: str = "etcd",
83
        durable_kv_events: bool = False,
84
85
        standalone_indexer: bool = False,
        zmq_replay: bool = False,
86
87
88
89
90
91
92
    ):
        """Initialize vLLM workers with dynamo integration.

        Args:
            request: pytest request fixture for log directory
            vllm_args: Configuration dict with keys:
                - model: Model name/path (default: TinyLlama-1.1B)
93
94
                - gpu_memory_utilization: Fraction of GPU memory to allocate (optional)
                - num_gpu_blocks_override: Cap on number of KV cache blocks (optional)
95
                - max_model_len: Maximum sequence length (optional)
96
                - enforce_eager: Disable CUDA graphs (default: False)
97
            num_workers: Number of vLLM worker processes
98
            single_gpu: If True, all workers share GPU 0
99
            data_parallel_size: If set, enables data parallelism with this many ranks (num_workers must equal data_parallel_size)
100
101
            request_plane: Request plane to use ("nats", "tcp", or "http"). Defaults to "tcp".
            store_backend: Storage backend to use ("etcd" or "file"). Defaults to "etcd".
102
            durable_kv_events: If True, use JetStream for durable KV events. Defaults to False (NATS Core mode).
103
104
105
106
107
108
109
        """
        # Generate unique namespace for isolation
        namespace_suffix = generate_random_suffix()
        self.namespace = f"test-namespace-{namespace_suffix}"
        self.component_name = "backend"
        self.endpoint = f"dyn://{self.namespace}.{self.component_name}.generate"
        self.num_workers = num_workers
110
        self.data_parallel_size = data_parallel_size
111
        self.worker_processes = []
112
113
        self.worker_id_to_zmq_ports: dict[int, dict[int, str]] = {}
        self._worker_id_to_replay_ports: dict[int, dict[int, str]] = {}
114
        self.store_backend = store_backend
115
116
117
118
119
120
121
122
        self._request = request
        self._request_plane = request_plane
        self._standalone_indexer = standalone_indexer
        self._zmq_replay = zmq_replay
        self._standalone_indexer_port: Optional[int] = None
        self._standalone_indexer_b_port: Optional[int] = None
        self._indexer_process: Optional[ManagedProcess] = None
        self._indexer_b_process: Optional[ManagedProcess] = None
123

124
125
126
127
128
        # Dynamically allocate unique system, KV event, and NIXL side-channel
        # ports (one of each per worker) to avoid conflicts in parallel test runs.
        self._system_ports = allocate_ports(num_workers, DefaultPort.SYSTEM1.value)
        self._kv_event_ports = allocate_ports(num_workers, DefaultPort.SYSTEM1.value)
        self._nixl_ports = allocate_ports(num_workers, DefaultPort.SYSTEM1.value)
129
130
131
132
133
134
135
136
137
138
139
        self._replay_ports = (
            allocate_ports(num_workers, DefaultPort.SYSTEM1.value)
            if standalone_indexer and zmq_replay
            else []
        )
        self._indexer_ports = (
            allocate_ports(2, DefaultPort.SYSTEM1.value) if standalone_indexer else []
        )
        if standalone_indexer:
            self._standalone_indexer_port = self._indexer_ports[0]
            self._standalone_indexer_b_port = self._indexer_ports[1]
140
141
        request.addfinalizer(
            lambda: deallocate_ports(
142
143
144
145
146
                self._system_ports
                + self._kv_event_ports
                + self._nixl_ports
                + self._replay_ports
                + self._indexer_ports
147
148
149
            )
        )

150
151
152
153
        if vllm_args is None:
            vllm_args = {}

        model = vllm_args.get("model", MODEL_NAME)
154
155
        gpu_memory_utilization = vllm_args.get("gpu_memory_utilization")
        num_gpu_blocks_override = vllm_args.get("num_gpu_blocks_override")
156
        max_model_len = vllm_args.get("max_model_len")
157
        enforce_eager = vllm_args.get("enforce_eager", False)
158
159

        self.model_name = model
160
        self.block_size = vllm_args.get("block_size", BLOCK_SIZE)
161
162
163
164
165
166

        # Create vLLM worker processes
        # Matches test.sh behavior:
        # - When data_parallel_size is set, launch one process per DP rank
        # - Each process gets --data-parallel-rank and --data-parallel-size
        # - Each process runs on its own GPU via CUDA_VISIBLE_DEVICES
167
        # - --kv-transfer-config enables KV cache transfer between ranks
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186

        for worker_idx in range(num_workers):
            # Calculate GPU device for this process
            if single_gpu:
                # Force all processes to GPU 0 (for single-GPU testing)
                gpu_device = "0"
            elif data_parallel_size is not None:
                # Worker sees dp_rank GPUs (each DP rank gets its own GPU)
                worker_start_gpu = worker_idx * data_parallel_size
                gpu_device = ",".join(
                    str(i)
                    for i in range(
                        worker_start_gpu, worker_start_gpu + data_parallel_size
                    )
                )
            else:
                # No DP; worker sees one GPU
                gpu_device = str(worker_idx)

187
188
189
190
            command = ["python3", "-m", "dynamo.vllm", "--model", model]

            if "block_size" in vllm_args:
                command.extend(["--block-size", str(vllm_args["block_size"])])
191

192
193
194
195
196
197
198
199
200
201
            # Disable CUDA graphs for faster startup & lower memory
            if enforce_eager:
                command.append("--enforce-eager")

            # Limit VRAM allocation (required for multi-worker on same GPU)
            if gpu_memory_utilization is not None:
                command.extend(
                    ["--gpu-memory-utilization", str(gpu_memory_utilization)]
                )

202
203
204
205
            # Add optional max_model_len if specified
            if max_model_len is not None:
                command.extend(["--max-model-len", str(max_model_len)])

206
207
208
209
210
211
            # Cap block count for predictable KV cache behavior
            if num_gpu_blocks_override is not None:
                command.extend(
                    ["--num-gpu-blocks-override", str(num_gpu_blocks_override)]
                )

212
213
214
215
216
217
218
219
220
            if data_parallel_size is not None:
                # Add DP configuration for external load balancing
                # See: https://docs.vllm.ai/en/v0.10.0/serving/data_parallel_deployment.html#external-load-balancing
                command.extend(
                    [
                        "--data-parallel-size",
                        str(data_parallel_size),
                        # "--data-parallel-address", "127.0.0.1",  # Required for DP coordination
                        # "--data-parallel-rpc-port", "13345",  # RPC port for DP coordination
221
                        # "--kv-transfer-config", '{"kv_connector":"NixlConnector","kv_role":"kv_both"}',  # Required for KV transfer between DP ranks
222
223
224
                    ]
                )

225
226
227
228
            # Use --durable-kv-events to enable JetStream mode (local indexer disabled)
            if durable_kv_events:
                command.append("--durable-kv-events")

229
230
231
232
            # Ports are dynamically allocated for xdist-safe parallel execution.
            system_port = self._system_ports[worker_idx]
            kv_event_port = self._kv_event_ports[worker_idx]
            nixl_port = self._nixl_ports[worker_idx]
233
234
235
236
237
            replay_port = (
                self._replay_ports[worker_idx]
                if worker_idx < len(self._replay_ports)
                else None
            )
238

239
            # Pass KV events config explicitly via CLI
240
241
242
243
244
245
246
247
248
            kv_events_cfg: Dict[str, Any] = {
                "publisher": "zmq",
                "topic": "kv-events",
                "endpoint": f"tcp://*:{kv_event_port}",
                "enable_kv_cache_events": True,
            }
            if replay_port is not None:
                kv_events_cfg["replay_endpoint"] = f"tcp://*:{replay_port}"
            command.extend(["--kv-events-config", json.dumps(kv_events_cfg)])
249

250
            env = os.environ.copy()  # Copy parent environment
251
252
253
254
            env_vars = {
                "CUDA_VISIBLE_DEVICES": gpu_device,
                "DYN_NAMESPACE": self.namespace,
                "DYN_REQUEST_PLANE": request_plane,
255
256
                "DYN_SYSTEM_PORT": str(system_port),
                "VLLM_NIXL_SIDE_CHANNEL_PORT": str(nixl_port),
257
258
259
260
261
262
263
264
                "PYTHONHASHSEED": "0",  # for deterministic event id's
            }

            # Add DYN_FILE_KV if using file storage backend
            if self.store_backend == "file" and "DYN_FILE_KV" in os.environ:
                env_vars["DYN_FILE_KV"] = os.environ["DYN_FILE_KV"]

            env.update(env_vars)
265
266
267
268
269
270
271
272
273
274

            # Create managed process for the worker
            process = ManagedProcess(
                command=command,
                env=env,
                timeout=120,  # Allow time for model loading
                display_output=True,
                health_check_ports=[],
                health_check_urls=[],
                log_dir=request.node.name,
275
                terminate_all_matching_process_names=False,
276
277
278
279
280
            )
            self.worker_processes.append(process)
            if data_parallel_size is not None:
                logger.info(
                    f"Created {data_parallel_size} DP ranks per worker on GPU(s) {gpu_device} "
281
                    f"(gpu_mem={gpu_memory_utilization}, system_port={system_port}) "
282
283
284
285
286
                    f"with endpoint: {self.endpoint}"
                )
            else:
                logger.info(
                    f"Created vLLM worker {worker_idx} on GPU {gpu_device} "
287
                    f"(gpu_mem={gpu_memory_utilization}, system_port={system_port}) "
288
289
290
                    f"with endpoint: {self.endpoint}"
                )

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
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
    @property
    def standalone_indexer_url(self) -> Optional[str]:
        if self._standalone_indexer_port is not None:
            return f"http://localhost:{self._standalone_indexer_port}"
        return None

    @property
    def standalone_indexer_b_url(self) -> Optional[str]:
        if self._standalone_indexer_b_port is not None:
            return f"http://localhost:{self._standalone_indexer_b_port}"
        return None

    def __enter__(self):
        if not self._standalone_indexer:
            return super().__enter__()

        indexer_cmd = [
            *get_kv_indexer_command(),
            "--block-size",
            str(self.block_size),
            "--port",
            str(self._standalone_indexer_port),
        ]
        self._indexer_process = ManagedProcess(
            command=indexer_cmd,
            timeout=120,
            display_output=True,
            health_check_ports=[self._standalone_indexer_port],
            health_check_urls=[],
            log_dir=self._request.node.name,
            terminate_all_matching_process_names=False,
            display_name="dynamo-kv-indexer",
        )
        logger.info(
            "Starting standalone indexer on port %s", self._standalone_indexer_port
        )
        self._indexer_process.__enter__()
        return self

    def __exit__(self, exc_type, exc_val, exc_tb):
        if self._standalone_indexer:
            for process in self.worker_processes:
                process.__exit__(exc_type, exc_val, exc_tb)
            if self._indexer_b_process is not None:
                self._indexer_b_process.__exit__(exc_type, exc_val, exc_tb)
                self._indexer_b_process = None
            if self._indexer_process is not None:
                self._indexer_process.__exit__(exc_type, exc_val, exc_tb)
                self._indexer_process = None
            return

        super().__exit__(exc_type, exc_val, exc_tb)

    async def launch_workers_with_indexer(self, endpoint):
        if not self._standalone_indexer:
            raise RuntimeError(
                "launch_workers_with_indexer requires standalone_indexer=True"
            )

        client = await endpoint.client()
        known_ids: set[int] = set()
        register_url = f"{self.standalone_indexer_url}/register"

        async with aiohttp.ClientSession() as session:
            for worker_idx, process in enumerate(self.worker_processes):
                process.__enter__()

                new_worker_id = None
                for _ in range(120):
                    ids = set(client.instance_ids())
                    new = ids - known_ids
                    if new:
                        new_worker_id = new.pop()
                        known_ids.add(new_worker_id)
                        break
                    await asyncio.sleep(0.5)

                if new_worker_id is None:
                    raise RuntimeError(
                        f"Timed out waiting for vLLM worker {worker_idx} to register "
                        f"(known_ids={known_ids})"
                    )

                zmq_endpoint = f"tcp://127.0.0.1:{self._kv_event_ports[worker_idx]}"
                replay_endpoint = (
                    f"tcp://127.0.0.1:{self._replay_ports[worker_idx]}"
                    if worker_idx < len(self._replay_ports)
                    else None
                )

                payload = {
                    "instance_id": new_worker_id,
                    "endpoint": zmq_endpoint,
                    "dp_rank": 0,
                    "model_name": self.model_name,
                    "block_size": self.block_size,
                }
                if replay_endpoint is not None:
                    payload["replay_endpoint"] = replay_endpoint

                async with session.post(register_url, json=payload) as resp:
                    if resp.status != 201:
                        body = await resp.text()
                        raise RuntimeError(
                            f"Failed to register vLLM instance {new_worker_id}: "
                            f"{resp.status} {body}"
                        )

                self.worker_id_to_zmq_ports[new_worker_id] = {0: zmq_endpoint}
                if replay_endpoint is not None:
                    self._worker_id_to_replay_ports[new_worker_id] = {
                        0: replay_endpoint
                    }

                logger.info(
                    "vLLM worker %s: worker_id=%s, zmq_endpoint=%s, replay_endpoint=%s",
                    worker_idx,
                    new_worker_id,
                    zmq_endpoint,
                    replay_endpoint,
                )

        await wait_for_indexer_workers_active(
            self.standalone_indexer_url, self.worker_id_to_zmq_ports
        )
        logger.info(
            "All %s vLLM workers launched and registered with indexer",
            self.num_workers,
        )

    def launch_indexer(self):
        if not self._standalone_indexer or self._standalone_indexer_b_port is None:
            raise RuntimeError("launch_indexer requires standalone_indexer=True")
        if not self.worker_id_to_zmq_ports:
            raise RuntimeError("launch_indexer requires workers to be registered first")

        worker_entries = []
        for worker_id, zmq_addresses in self.worker_id_to_zmq_ports.items():
            for dp_rank, zmq_endpoint in zmq_addresses.items():
                worker_entries.append(f"{worker_id}:{dp_rank}={zmq_endpoint}")
        workers_arg = ",".join(worker_entries)

        indexer_b_cmd = [
            *get_kv_indexer_command(),
            "--block-size",
            str(self.block_size),
            "--port",
            str(self._standalone_indexer_b_port),
            "--peers",
            f"http://localhost:{self._standalone_indexer_port}",
            "--workers",
            workers_arg,
            "--model-name",
            self.model_name,
        ]
        self._indexer_b_process = ManagedProcess(
            command=indexer_b_cmd,
            timeout=120,
            display_output=True,
            health_check_ports=[self._standalone_indexer_b_port],
            health_check_urls=[],
            log_dir=self._request.node.name,
            terminate_all_matching_process_names=False,
            display_name="dynamo-kv-indexer-b",
        )
        logger.info(
            "Starting standalone indexer B on port %s with peer http://localhost:%s",
            self._standalone_indexer_b_port,
            self._standalone_indexer_port,
        )
        self._indexer_b_process.__enter__()

463
464
465
    process_name = "vLLM worker"
    cleanup_name = "vLLM worker resources"
    init_delay_reason = "initialize NIXL before starting next worker"
466
467


468
@pytest.mark.pre_merge
469
@pytest.mark.gpu_1
470
@pytest.mark.timeout(150)  # ~3x average (~43s/test), rounded up
471
@pytest.mark.parametrize("request_plane", ["tcp"], indirect=True)
472
def test_vllm_kv_router_basic(
473
474
475
476
477
    request,
    runtime_services_dynamic_ports,
    predownload_models,
    set_ucx_tls_no_mm,
    request_plane,
478
):
479
480
481
482
483
484
485
    run_basic_router_test(
        engine_process_cls=VLLMProcess,
        engine_args_name="vllm_args",
        engine_args=VLLM_ARGS,
        num_workers=2,
        single_gpu=True,
        request=request,
486
        request_plane=request_plane,
487
        block_size=BLOCK_SIZE,
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
        model_name=MODEL_NAME,
    )


@pytest.mark.pre_merge
@pytest.mark.gpu_1
@pytest.mark.timeout(150)  # ~3x average (~43s/test), rounded up
@pytest.mark.parametrize("request_plane", ["tcp"], indirect=True)
def test_vllm_kv_router_without_block_size_specified_in_vllm_args(
    request,
    runtime_services_dynamic_ports,
    predownload_models,
    set_ucx_tls_no_mm,
    request_plane,
):
    run_basic_router_test(
        engine_process_cls=VLLMProcess,
        engine_args_name="vllm_args",
        engine_args=VLLM_ARGS_NO_BLOCK_SIZE,
        num_workers=2,
        single_gpu=True,
        request=request,
        request_plane=request_plane,
        block_size=BLOCK_SIZE,
512
513
        model_name=MODEL_NAME,
    )
514
515


516
@pytest.mark.pre_merge
517
@pytest.mark.gpu_1
518
@pytest.mark.timeout(150)  # ~3x average (~43s/test), rounded up
519
@pytest.mark.parametrize("request_plane", ["tcp"], indirect=True)
520
def test_router_decisions_vllm_multiple_workers(
521
522
523
524
525
    request,
    runtime_services_dynamic_ports,
    predownload_models,
    set_ucx_tls_no_mm,
    request_plane,
526
):
527
528
529
530
531
    run_router_decisions_test(
        engine_process_cls=VLLMProcess,
        engine_args_name="vllm_args",
        engine_args=VLLM_ARGS,
        request=request,
532
        request_plane=request_plane,
533
534
535
536
537
538
539
        model_name=MODEL_NAME,
        block_size=BLOCK_SIZE,
        component_name="backend",
        num_workers=2,
        single_gpu=True,
        test_dp_rank=False,
    )
540
541
542


@pytest.mark.gpu_2
543
@pytest.mark.nightly
544
@pytest.mark.parametrize("request_plane", ["tcp"], indirect=True)
545
@pytest.mark.timeout(600)  # 10 min max (multi-GPU + DP startup variance)
546
def test_router_decisions_vllm_dp(
547
548
549
550
551
    request,
    runtime_services_dynamic_ports,
    predownload_models,
    set_ucx_tls_no_mm,
    request_plane,
552
):
553
554
555
556
557
558
559
    """Validate KV cache prefix reuse with vLLM by sending progressive requests with overlapping prefixes.
    Same flow as test_router_decisions_vllm_multiple_workers; force first request to (worker_id, dp_rank=1).
    Dump events from router and verify:
        * All but one (worker_id, dp_rank) should have no events (due to prefix reuse)
        * The (worker_id, dp_rank) with events should have exactly 4 events (one per request)
        * All events should be on the forced (worker_id, dp_rank=1) (verifying forced routing and prefix reuse)
    """
560
561
562
563
564
    run_router_decisions_test(
        engine_process_cls=VLLMProcess,
        engine_args_name="vllm_args",
        engine_args=VLLM_ARGS,
        request=request,
565
        request_plane=request_plane,
566
567
568
569
570
571
572
573
        model_name=MODEL_NAME,
        block_size=BLOCK_SIZE,
        component_name="backend",
        num_workers=1,
        single_gpu=False,
        test_dp_rank=True,
        extra_process_kwargs={"data_parallel_size": 2},
    )
574

575
576
577

@pytest.mark.pre_merge
@pytest.mark.gpu_1
578
@pytest.mark.timeout(150)  # ~3x average (~43s/test), rounded up
579
@pytest.mark.parametrize(
580
    "store_backend,durable_kv_events,request_plane",
581
    [
582
        ("etcd", False, "tcp"),
583
    ],
584
585
    ids=["nats_core"],
    indirect=["durable_kv_events", "request_plane"],
586
)
587
def test_vllm_indexers_sync(
588
589
590
591
592
593
    request,
    runtime_services_dynamic_ports,
    predownload_models,
    file_storage_backend,
    set_ucx_tls_no_mm,
    store_backend,
594
    durable_kv_events,
595
    request_plane,
596
):
597
598
599
600
601
602
    run_indexers_sync_test(
        engine_process_cls=VLLMProcess,
        engine_args_name="vllm_args",
        engine_args=VLLM_ARGS,
        request=request,
        runtime_services_dynamic_ports=runtime_services_dynamic_ports,
603
604
        store_backend=store_backend,
        durable_kv_events=durable_kv_events,
605
606
607
608
        request_plane=request_plane,
        block_size=BLOCK_SIZE,
        model_name=MODEL_NAME,
        num_workers=2,
609
        extra_process_kwargs={"standalone_indexer": True, "zmq_replay": True},
610
    )