diff --git a/examples/backends/vllm/mm_router_worker/mm_router_worker.py b/examples/backends/vllm/mm_router_worker/mm_router_worker.py index 5d02984f4665..683487055fd7 100644 --- a/examples/backends/vllm/mm_router_worker/mm_router_worker.py +++ b/examples/backends/vllm/mm_router_worker/mm_router_worker.py @@ -11,7 +11,7 @@ Usage: python -m examples.backends.vllm.mm_router_worker \ --model Qwen/Qwen3-VL-8B-Instruct \ - --namespace default \ + --namespace dynamo \ --component mm_router \ --endpoint generate \ --downstream-component backend \ @@ -21,6 +21,7 @@ import argparse import asyncio import logging +import os import signal import uvloop @@ -58,8 +59,8 @@ def parse_args() -> argparse.Namespace: parser.add_argument( "--namespace", type=str, - default="default", - help="Dynamo namespace", + default=os.environ.get("DYN_NAMESPACE", "dynamo"), + help="Dynamo namespace (default: DYN_NAMESPACE env or 'dynamo')", ) parser.add_argument( "--component", @@ -88,6 +89,15 @@ def parse_args() -> argparse.Namespace: help="Downstream vLLM workers' endpoint name", ) + # Router configuration + parser.add_argument( + "--no-router-kv-events", + action="store_true", + default=False, + help="Use approximate KV routing (no KV events from workers). " + "Required for hybrid models like Qwen3.5 that cannot emit KV events.", + ) + return parser.parse_args() @@ -135,12 +145,17 @@ def signal_handler(): logger.info(f"Found {len(instance_ids)} workers: {list(instance_ids)}") # Create KvRouter to select workers based on KV overlap + kv_router_config = KvRouterConfig( + use_kv_events=not args.no_router_kv_events, + ) kv_router = KvRouter( endpoint=downstream_endpoint, block_size=args.block_size, - kv_router_config=KvRouterConfig(), + kv_router_config=kv_router_config, + ) + logger.info( + f"KvRouter created successfully (use_kv_events={not args.no_router_kv_events})" ) - logger.info("KvRouter created successfully") # Initialize tokenizer and processor for MM processing logger.info(f"Loading tokenizer from {args.model}...") diff --git a/examples/backends/vllm/qwen35/launch.sh b/examples/backends/vllm/qwen35/launch.sh new file mode 100755 index 000000000000..8c8362af369b --- /dev/null +++ b/examples/backends/vllm/qwen35/launch.sh @@ -0,0 +1,124 @@ +#!/bin/bash +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# Aggregated multimodal serving for Qwen3.5 hybrid models with MM-aware approximate KV routing. +# +# Qwen3.5 is a multimodal hybrid model (GatedDeltaNet + Gated Attention + Vision Encoder). +# It supports images and video but has hybrid architecture constraints: +# +# 1. DO NOT pass --kv-events-config or --enable-kv-cache-events: +# vLLM disables the Hybrid KV Cache Manager when kv_events_config is set, +# but Qwen3.5's mixed KV cache specs (GDN + FullAttention) cannot be unified +# into one type. +# +# 2. Use --mamba-cache-mode align (not "all"): +# Qwen3.5 raises NotImplementedError with mamba_cache_mode="all". +# +# 3. Approximate KV routing (--no-router-kv-events): +# Hybrid models cannot emit KV events to the router. The MM Router Worker +# predicts cache state from its own routing decisions using prefix hashing. +# +# 4. Disaggregated P/D is NOT supported for hybrid models in vLLM: +# HybridKVCacheCoordinator asserts dcp_world_size == 1. +# +# 5. Use TCP transport for multimodal payloads (NATS has 1MB limit). +# +# Architecture: +# Frontend (--router-mode round-robin) +# -> MM Router Worker (approximate KV routing + multimodal hash) +# -> vLLM Worker (--enable-multimodal --mamba-cache-mode align) + +set -e +trap 'echo Cleaning up...; kill 0' EXIT + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "$SCRIPT_DIR/../../../common/gpu_utils.sh" +source "$SCRIPT_DIR/../../../common/launch_utils.sh" + +MODEL="${MODEL:-Qwen/Qwen3.5-0.8B}" + +EXTRA_ARGS=() +while [[ $# -gt 0 ]]; do + case $1 in + --model) MODEL="$2"; shift 2 ;; + *) EXTRA_ARGS+=("$1"); shift ;; + esac +done + +MAX_MODEL_LEN="${MAX_MODEL_LEN:-4096}" +MAX_CONCURRENT_SEQS="${MAX_CONCURRENT_SEQS:-2}" +HTTP_PORT="${DYN_HTTP_PORT:-8000}" +BLOCK_SIZE="${BLOCK_SIZE:-16}" + +# TCP transport: avoids NATS 1MB payload limit for base64-encoded images +export DYN_REQUEST_PLANE=tcp + +print_launch_banner --no-curl "Launching Qwen3.5 Multimodal + MM Router + Approx KV (1 GPU)" "$MODEL" "$HTTP_PORT" \ + "Backend: dynamo.vllm --enable-multimodal --mamba-cache-mode align" \ + "MM Router: MM-aware approximate KV routing (--no-router-kv-events)" \ + "Frontend: round-robin to MM Router" \ + "Transport: TCP (multimodal payloads)" + +print_curl_footer < str: class VLLMWorkerProcess(ManagedProcess): - """vLLM backend worker that emits KV events.""" - - def __init__(self, request, *, system_port: int, kv_event_port: int): - super().__init__( - command=[ - "python3", - "-m", - "dynamo.vllm", - "--model", - VLLM_MM_MODEL, - "--enable-multimodal", - "--gpu-memory-utilization", - "0.85", - "--max-model-len", - "8192", - "--served-model-name", - f"{VLLM_MM_MODEL}__internal", + """vLLM backend worker that optionally emits KV events.""" + + def __init__(self, request, *, system_port: int, kv_event_port: int | None = None): + cmd = [ + "python3", + "-m", + "dynamo.vllm", + "--model", + VLLM_MM_MODEL, + "--enable-multimodal", + "--gpu-memory-utilization", + "0.85", + "--max-model-len", + "8192", + "--served-model-name", + f"{VLLM_MM_MODEL}__internal", + ] + if kv_event_port is not None: + cmd += [ "--kv-events-config", ( f'{{"publisher":"zmq","topic":"kv-events",' f'"endpoint":"tcp://*:{kv_event_port}",' f'"enable_kv_cache_events": true}}' ), - ], + ] + super().__init__( + command=cmd, env=_make_process_env(DYN_SYSTEM_PORT=str(system_port)), health_check_urls=[ (f"http://localhost:{system_port}/health", _check_ready) @@ -135,29 +139,32 @@ def __init__(self, request, *, system_port: int, kv_event_port: int): class VLLMMMRouterWorkerProcess(ManagedProcess): - """vLLM MM router worker.""" - - def __init__(self, request, *, system_port: int): + """vLLM MM router worker (exact or approximate KV routing).""" + + def __init__(self, request, *, system_port: int, approx_routing: bool = False): + cmd = [ + "python3", + "-m", + "examples.backends.vllm.mm_router_worker", + "--model", + VLLM_MM_MODEL, + "--namespace", + NAMESPACE, + "--component", + "mm_router", + "--endpoint", + "generate", + "--downstream-component", + "backend", + "--downstream-endpoint", + "generate", + "--block-size", + str(BLOCK_SIZE), + ] + if approx_routing: + cmd.append("--no-router-kv-events") super().__init__( - command=[ - "python3", - "-m", - "examples.backends.vllm.mm_router_worker", - "--model", - VLLM_MM_MODEL, - "--namespace", - NAMESPACE, - "--component", - "mm_router", - "--endpoint", - "generate", - "--downstream-component", - "backend", - "--downstream-endpoint", - "generate", - "--block-size", - str(BLOCK_SIZE), - ], + command=cmd, env=_make_process_env( DYN_SYSTEM_USE_ENDPOINT_HEALTH_STATUS='["generate"]', DYN_SYSTEM_PORT=str(system_port), @@ -207,17 +214,27 @@ def mm_runtime_services(request): os.environ.pop("ETCD_ENDPOINTS", None) -@pytest.fixture(scope="module") +@pytest.fixture(scope="module", params=[False, True], ids=["exact_kv", "approx_kv"]) def start_vllm_mm_services( request, mm_runtime_services ) -> Generator[tuple[int, ManagedProcess], None, None]: - frontend_port, vllm_port, router_port, kv_event_port = allocate_ports( - count=4, start_port=10000 - ) + approx_routing = request.param + + if approx_routing: + frontend_port, vllm_port, router_port = allocate_ports( + count=3, start_port=10000 + ) + kv_event_port = None + else: + frontend_port, vllm_port, router_port, kv_event_port = allocate_ports( + count=4, start_port=10000 + ) with VLLMWorkerProcess(request, system_port=vllm_port, kv_event_port=kv_event_port): time.sleep(10) - with VLLMMMRouterWorkerProcess(request, system_port=router_port) as router_proc: + with VLLMMMRouterWorkerProcess( + request, system_port=router_port, approx_routing=approx_routing + ) as router_proc: time.sleep(3) with FrontendProcess(request, frontend_port=frontend_port): yield frontend_port, router_proc