diff --git a/components/src/dynamo/sglang/register.py b/components/src/dynamo/sglang/register.py index dcaf9e38a4ae..be2afe36e667 100644 --- a/components/src/dynamo/sglang/register.py +++ b/components/src/dynamo/sglang/register.py @@ -53,6 +53,7 @@ effective_gateway_workers, gateway_engine_id, ) +from dynamo.sglang.video_routing import publish_sglang_qwen_video_processor_contract SGLANG_HICACHE_MOONCAKE_RUNTIME_KEY = "sglang_hicache_mooncake" SPEC_DECODE_RUNTIME_KEY = "spec_decode" @@ -433,6 +434,15 @@ async def get_runtime_config( # generation overflow handling to their downstream backend. if engine is not None: publish_token_budget(runtime_config, _get_token_budget(engine, server_args)) + # Hash forwarding currently exists on SGLang's aggregated generation + # path. Do not advertise exact video routing to disaggregated workers, + # whose prefill/decode handlers would otherwise publish incompatible + # KV-event placeholder hashes. + if dynamo_args.frontend_decoding and server_args.disaggregation_mode in ( + None, + "null", + ): + publish_sglang_qwen_video_processor_contract(runtime_config, engine) # set reasoning parser and tool call parser runtime_config.reasoning_parser = dynamo_args.dyn_reasoning_parser runtime_config.tool_call_parser = dynamo_args.dyn_tool_call_parser diff --git a/components/src/dynamo/sglang/request_handlers/llm/decode_handler.py b/components/src/dynamo/sglang/request_handlers/llm/decode_handler.py index 63a3d6871c50..825532de64c8 100644 --- a/components/src/dynamo/sglang/request_handlers/llm/decode_handler.py +++ b/components/src/dynamo/sglang/request_handlers/llm/decode_handler.py @@ -4,7 +4,7 @@ import asyncio import logging import time -from typing import Any, AsyncGenerator, AsyncIterator, Dict, List, Mapping, Optional +from typing import Any, AsyncGenerator, AsyncIterator, Dict, Mapping, Optional import numpy as np import sglang as sgl @@ -47,6 +47,7 @@ VIDEO_URL_KEY, build_disagg_mm_kwargs, extract_media_urls, + extract_mm_hashes, raise_if_unextracted_multimodal, ) from dynamo.sglang.request_utils import request_cache_salt @@ -433,32 +434,6 @@ def _resolve_mm_hashes_supported(engine: Any) -> bool: probe = filter_supported_async_generate_kwargs(engine, {"mm_hashes": None}) return "mm_hashes" in probe - @staticmethod - def _extract_mm_hashes(request: Dict[str, Any]) -> Optional[List[str]]: - """Pull the per-image hashes the Rust frontend forwards via extra_args. - - Returns ``None`` when the field is absent or malformed; SGLang then - recomputes the hash internally via ``hash_feature()``. - """ - extra_args = request.get("extra_args") - if not isinstance(extra_args, dict): - return None - mm_hashes = extra_args.get("mm_hashes") - if not mm_hashes: - return None - if not isinstance(mm_hashes, list): - return None - # Fail closed if a non-string slipped into the list — downstream - # SGLang treats mm_hashes as List[str] and a bad element would - # crash the worker mid-request. Routing falls back to text-prefix. - if not all(isinstance(h, str) for h in mm_hashes): - logging.warning( - "extra_args.mm_hashes contained non-str entries; " - "ignoring routing-side hashes and letting SGLang recompute" - ) - return None - return mm_hashes - def _metadata_uploader_from_request( self, request: Dict[str, Any] ) -> MetadataUploader | None: @@ -797,7 +772,7 @@ async def generate( mm_hashes_kwargs: Dict[str, Any] = {} if self._mm_hashes_supported: - forwarded = self._extract_mm_hashes(request) + forwarded = extract_mm_hashes(request) if forwarded is not None: mm_hashes_kwargs["mm_hashes"] = forwarded diff --git a/components/src/dynamo/sglang/request_handlers/llm/mm_disagg_utils.py b/components/src/dynamo/sglang/request_handlers/llm/mm_disagg_utils.py index e4e9a954d884..0b2d3487dbe5 100644 --- a/components/src/dynamo/sglang/request_handlers/llm/mm_disagg_utils.py +++ b/components/src/dynamo/sglang/request_handlers/llm/mm_disagg_utils.py @@ -20,6 +20,14 @@ _SUPPORTED_MULTIMODAL_CONTENT_TYPES = frozenset( {IMAGE_URL_KEY, AUDIO_URL_KEY, VIDEO_URL_KEY} ) +# BaseMultiModalProcessorOutput.organize_results() builds SGLang's mm_items in +# this order, independent of their order in the original prompt. +_SGLANG_MM_ITEM_MODALITY_ORDER = ("image", "video", "audio") +_MM_DATA_KEY_BY_MODALITY = { + "image": IMAGE_URL_KEY, + "video": VIDEO_URL_KEY, + "audio": AUDIO_URL_KEY, +} def _multi_modal_data(request: Dict[str, Any]) -> Dict[str, Any]: @@ -124,6 +132,94 @@ def extract_media_urls( return urls or None +def extract_mm_hashes(request: Dict[str, Any]) -> list[str] | None: + """Return frontend MM hashes in SGLang's modality-grouped item order.""" + extra_args = request.get("extra_args") + if not isinstance(extra_args, dict): + return None + + grouped = extra_args.get("mm_hashes_by_modality") + if grouped is not None: + if not isinstance(grouped, dict): + logger.warning( + "extra_args.mm_hashes_by_modality is not an object; " + "ignoring routing-side hashes and letting SGLang recompute" + ) + return None + + unknown_modalities = { + str(modality) + for modality, hashes in grouped.items() + if modality not in _SGLANG_MM_ITEM_MODALITY_ORDER and hashes + } + if unknown_modalities: + logger.warning( + "extra_args.mm_hashes_by_modality contains unsupported " + "modalities %s; ignoring routing-side hashes and letting " + "SGLang recompute", + sorted(unknown_modalities), + ) + return None + + hashes_by_modality: dict[str, list[str]] = {} + for modality in _SGLANG_MM_ITEM_MODALITY_ORDER: + hashes = grouped.get(modality) + if hashes is None: + hashes = [] + if not isinstance(hashes, list) or not all( + isinstance(value, str) for value in hashes + ): + logger.warning( + "extra_args.mm_hashes_by_modality[%s] is not a string " + "list; ignoring routing-side hashes and letting SGLang recompute", + modality, + ) + return None + hashes_by_modality[modality] = hashes + + mm_data = request.get("multi_modal_data") + if not isinstance(mm_data, dict): + logger.warning( + "extra_args.mm_hashes_by_modality has no matching " + "multi_modal_data object; ignoring routing-side hashes and " + "letting SGLang recompute" + ) + return None + + flattened: list[str] = [] + for modality in _SGLANG_MM_ITEM_MODALITY_ORDER: + hashes = hashes_by_modality[modality] + media_items = mm_data.get(_MM_DATA_KEY_BY_MODALITY[modality]) + if media_items is None: + media_items = [] + if not isinstance(media_items, list) or len(hashes) != len(media_items): + media_count = ( + len(media_items) if isinstance(media_items, list) else None + ) + logger.warning( + "extra_args.mm_hashes_by_modality[%s] count (%d) does not " + "match multi_modal_data count (%s); ignoring routing-side " + "hashes and letting SGLang recompute", + modality, + len(hashes), + media_count, + ) + return None + flattened.extend(hashes) + return flattened or None + + mm_hashes = extra_args.get("mm_hashes") + if not mm_hashes or not isinstance(mm_hashes, list): + return None + if not all(isinstance(value, str) for value in mm_hashes): + logger.warning( + "extra_args.mm_hashes contained non-str entries; ignoring " + "routing-side hashes and letting SGLang recompute" + ) + return None + return mm_hashes + + def build_disagg_mm_kwargs(request: Dict[str, Any]) -> Dict[str, Any]: """Build media kwargs for a disaggregated worker's ``async_generate`` call. diff --git a/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py b/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py index c116e8a8719d..1dc7c72676fb 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py +++ b/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py @@ -513,6 +513,37 @@ async def fake_async_generate(**kwargs): assert captured["audio_data"] == ["https://example.com/a.wav"] +@pytest.mark.asyncio +async def test_aggregated_forwards_grouped_mm_hashes_in_sglang_item_order(): + handler = _new_decode_handler(enable_frontend_decoding=False) + handler._mm_hashes_supported = True + captured: Dict[str, Any] = {} + + async def fake_async_generate(**kwargs): + captured.update(kwargs) + return _empty_stream() + + handler.engine = SimpleNamespace(async_generate=fake_async_generate) + request = { + "token_ids": [1, 2, 3], + "multi_modal_data": { + "image_url": ["https://example.com/a.jpg"], + "video_url": ["https://example.com/a.mp4"], + }, + "extra_args": { + "mm_hashes_by_modality": { + "video": ["video-a"], + "image": ["image-a"], + } + }, + } + + async for _ in handler.generate(request, _Context()): + pass + + assert captured["mm_hashes"] == ["image-a", "video-a"] + + @pytest.mark.asyncio async def test_aggregated_fd_on_loads_decoded_variants_to_pil(): """With --frontend-decoding, Decoded items are loaded via ImageLoader and diff --git a/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py b/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py index 290dce6e7bc1..74d1bfbe59b8 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py +++ b/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py @@ -9,6 +9,7 @@ from dynamo.sglang.request_handlers.llm.mm_disagg_utils import ( build_disagg_mm_kwargs, extract_media_urls, + extract_mm_hashes, raise_if_unextracted_multimodal, ) from dynamo.sglang.request_handlers.multimodal.worker_handler import StreamProcessor @@ -79,6 +80,60 @@ def test_extract_media_urls_rejects_malformed_payloads(): extract_media_urls({"image_url": ""}, "image_url") +def test_extract_mm_hashes_preserves_legacy_image_protocol(): + request = {"extra_args": {"mm_hashes": ["image-a", "image-b"]}} + + assert extract_mm_hashes(request) == ["image-a", "image-b"] + + +def test_extract_mm_hashes_flattens_in_sglang_item_order(): + request = { + "multi_modal_data": { + "image_url": ["image-a", "image-b"], + "video_url": ["video-a"], + }, + "extra_args": { + "mm_hashes_by_modality": { + "video": ["video-a"], + "image": ["image-a", "image-b"], + }, + }, + } + + assert extract_mm_hashes(request) == ["image-a", "image-b", "video-a"] + + +def test_extract_mm_hashes_rejects_per_modality_count_mismatch(): + request = { + "multi_modal_data": { + "image_url": ["image-a", "image-b"], + "video_url": ["video-a"], + }, + "extra_args": { + "mm_hashes_by_modality": { + # The total count still matches, but the modality association does not. + "image": ["image-a"], + "video": ["video-a", "video-b"], + } + }, + } + + assert extract_mm_hashes(request) is None + + +@pytest.mark.parametrize( + "grouped", + [ + ["not-an-object"], + {"video": "not-a-list"}, + {"video": ["ok", 1]}, + {"future_modality": ["hash"]}, + ], +) +def test_extract_mm_hashes_rejects_malformed_grouped_protocol(grouped): + assert extract_mm_hashes({"extra_args": {"mm_hashes_by_modality": grouped}}) is None + + class TestMultimodalGuard: """Tests for multimodal guard when frontend extraction is missing.""" diff --git a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py new file mode 100644 index 000000000000..4f01bf3738e1 --- /dev/null +++ b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py @@ -0,0 +1,143 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +import json +from types import SimpleNamespace + +import pytest + +from dynamo.llm import ModelRuntimeConfig +from dynamo.sglang import video_routing + +pytestmark = [ + pytest.mark.unit, + pytest.mark.sglang, + pytest.mark.multimodal, + pytest.mark.gpu_0, + pytest.mark.profiled_vram_gib(0), + pytest.mark.pre_merge, +] + + +def _engine( + *, + video_config=None, + model_type="qwen3_vl", + architecture="Qwen3VLForConditionalGeneration", + overrides_video_replacement=True, +): + if overrides_video_replacement: + + class QwenProcessor(video_routing.transformers.ProcessorMixin): + def replace_video_token(self): + return None + + else: + + class QwenProcessor(video_routing.transformers.ProcessorMixin): + pass + + mm_processor = SimpleNamespace( + hf_config=SimpleNamespace( + model_type=model_type, + architectures=[architecture], + ), + video_config=video_config or {}, + _processor=QwenProcessor.__new__(QwenProcessor), + ) + return SimpleNamespace(tokenizer_manager=SimpleNamespace(mm_processor=mm_processor)) + + +@pytest.fixture +def qwen_preprocessor(monkeypatch): + monkeypatch.setattr( + video_routing, + "sglang_qwen_vl", + SimpleNamespace( + IMAGE_FACTOR=28, + VIDEO_MIN_PIXELS=100352, + VIDEO_MAX_PIXELS=602112, + VIDEO_TOTAL_PIXELS=90316800, + FRAME_FACTOR=2, + FPS=2.0, + FPS_MIN_FRAMES=4, + FPS_MAX_FRAMES=768, + ), + ) + monkeypatch.setattr( + video_routing, + "qwen3_smart_resize", + lambda **_: (1216, 4096), + ) + + +def test_publishes_sglang_qwen_video_contract(qwen_preprocessor): + runtime_config = ModelRuntimeConfig() + + video_routing.publish_sglang_qwen_video_processor_contract( + runtime_config, _engine() + ) + + contract = json.loads( + runtime_config.runtime_data[ + video_routing.SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY + ] + ) + assert contract == { + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil", + "runless_boundary_hash": "tokens_only", + "sglang_preprocess": { + "image_factor": 28, + "video_min_pixels": 100352, + "video_max_pixels": 602112, + "video_total_pixels": 90316800, + "frame_factor": 2, + "fps": 2.0, + "min_frames": 4, + "max_frames": 768, + }, + } + + +def test_processor_override_disables_exact_video_contract(qwen_preprocessor): + assert ( + video_routing._resolve_qwen_video_processor_contract( + _engine(video_config={"fps": 1.0}) + ) + is None + ) + + +def test_missing_processor_disables_exact_video_contract(qwen_preprocessor): + engine = _engine() + engine.tokenizer_manager.mm_processor._processor = None + + assert video_routing._resolve_qwen_video_processor_contract(engine) is None + + +def test_inherited_video_replacement_publishes_wrapped_target(qwen_preprocessor): + contract = video_routing._resolve_qwen_video_processor_contract( + _engine(overrides_video_replacement=False) + ) + + assert contract is not None + assert contract["placeholder_target"] == "vision_wrapped_video_token" + + +@pytest.mark.parametrize( + ("model_type", "architecture"), + [ + ("llava", "Qwen3VLForConditionalGeneration"), + ("qwen3_vl", "Qwen3VLForCausalLM"), + ], +) +def test_non_qwen_video_model_does_not_publish_contract( + qwen_preprocessor, model_type, architecture +): + assert ( + video_routing._resolve_qwen_video_processor_contract( + _engine(model_type=model_type, architecture=architecture) + ) + is None + ) diff --git a/components/src/dynamo/sglang/video_routing.py b/components/src/dynamo/sglang/video_routing.py new file mode 100644 index 000000000000..a0fcbebd0cac --- /dev/null +++ b/components/src/dynamo/sglang/video_routing.py @@ -0,0 +1,167 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Publish SGLang's effective Qwen video preprocessing contract.""" + +import json +import logging +from collections.abc import Callable +from typing import Any, Optional + +import transformers + +from dynamo.llm import ModelRuntimeConfig + +try: + from sglang.srt.multimodal.processors import qwen_vl as sglang_qwen_vl +except ImportError: + sglang_qwen_vl = None + +qwen3_smart_resize: Optional[Callable[..., Any]] +try: + from transformers.models.qwen3_vl.video_processing_qwen3_vl import ( + smart_resize as _qwen3_smart_resize, + ) +except ImportError: + qwen3_smart_resize = None +else: + qwen3_smart_resize = _qwen3_smart_resize + +logger = logging.getLogger(__name__) + +SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY = ( + "sglang_qwen_video_processor_contract" +) +QWEN_VIDEO_TARGET_BARE = "bare_video_token" +QWEN_VIDEO_TARGET_WRAPPED = "vision_wrapped_video_token" +QWEN_VIDEO_RESIZE_LEGACY_CEIL = "legacy_ceil" +QWEN_VIDEO_RESIZE_ROUND_TIES_EVEN = "round_ties_even" +QWEN_VIDEO_RUNLESS_BOUNDARY_TOKENS_ONLY = "tokens_only" +QWEN_VIDEO_MODEL_TYPES = {"qwen3_vl", "qwen3_vl_moe", "qwen3_5", "qwen3_5_moe"} +QWEN_VIDEO_ARCHITECTURES = { + "Qwen3VLForConditionalGeneration", + "Qwen3VLMoeForConditionalGeneration", + "Qwen3_5ForConditionalGeneration", + "Qwen3_5MoeForConditionalGeneration", +} + + +def _resolve_qwen_video_resize_mode() -> Optional[str]: + """Identify the installed Transformers Qwen video resize rule.""" + if qwen3_smart_resize is None: + logger.warning( + "Exact SGLang video-aware KV routing disabled because the installed " + "Transformers package has no Qwen3 video processor" + ) + return None + try: + result = qwen3_smart_resize( + num_frames=5, + height=1120, + width=3760, + temporal_factor=2, + factor=32, + min_pixels=4096, + max_pixels=25165824, + ) + except (TypeError, ValueError) as error: + logger.warning( + "Exact SGLang video-aware KV routing disabled because the installed " + "Qwen smart_resize API is unsupported: %s", + error, + ) + return None + if result == (1216, 4096): + return QWEN_VIDEO_RESIZE_LEGACY_CEIL + if result == (1120, 3776): + return QWEN_VIDEO_RESIZE_ROUND_TIES_EVEN + logger.warning( + "Exact SGLang video-aware KV routing disabled because the installed " + "Qwen smart_resize behavior is unsupported: %s", + result, + ) + return None + + +def _resolve_qwen_video_processor_contract(engine: Any) -> Optional[dict[str, Any]]: + """Resolve the two-stage video preprocessing used by this SGLang worker.""" + tokenizer_manager = getattr(engine, "tokenizer_manager", None) + mm_processor = getattr(tokenizer_manager, "mm_processor", None) + hf_config = getattr(mm_processor, "hf_config", None) + if hf_config is None: + model_config = getattr(tokenizer_manager, "model_config", None) + hf_config = getattr(model_config, "hf_config", None) + + model_type = getattr(hf_config, "model_type", None) + if model_type not in QWEN_VIDEO_MODEL_TYPES: + return None + architectures = getattr(hf_config, "architectures", None) or [] + if not QWEN_VIDEO_ARCHITECTURES.intersection(architectures): + return None + if mm_processor is None: + logger.warning( + "Exact SGLang video-aware KV routing disabled because the Qwen " + "multimodal processor is unavailable" + ) + return None + + video_config = getattr(mm_processor, "video_config", None) or {} + if video_config: + logger.warning( + "Exact SGLang video-aware KV routing disabled because engine-level " + "mm_process_config.video can change frame sampling or resizing" + ) + return None + if sglang_qwen_vl is None: + logger.warning( + "Exact SGLang video-aware KV routing disabled because the installed " + "SGLang package has no Qwen video preprocessor" + ) + return None + + processor = getattr(mm_processor, "_processor", None) + if processor is None: + logger.warning( + "Exact SGLang video-aware KV routing disabled because the Qwen " + "processor implementation is unavailable" + ) + return None + processor_impl = getattr(type(processor), "replace_video_token", None) + mixin_impl = getattr(transformers.ProcessorMixin, "replace_video_token", None) + placeholder_target = QWEN_VIDEO_TARGET_WRAPPED + if processor_impl is not None and processor_impl is not mixin_impl: + placeholder_target = QWEN_VIDEO_TARGET_BARE + + resize_mode = _resolve_qwen_video_resize_mode() + if resize_mode is None: + return None + + return { + "placeholder_target": placeholder_target, + "resize_mode": resize_mode, + # Compatibility wire field: SGLang publishes canonical pad-valued KV + # token blocks and no separate MM metadata for this worker. + "runless_boundary_hash": QWEN_VIDEO_RUNLESS_BOUNDARY_TOKENS_ONLY, + "sglang_preprocess": { + "image_factor": int(sglang_qwen_vl.IMAGE_FACTOR), + "video_min_pixels": int(sglang_qwen_vl.VIDEO_MIN_PIXELS), + "video_max_pixels": int(sglang_qwen_vl.VIDEO_MAX_PIXELS), + "video_total_pixels": int(sglang_qwen_vl.VIDEO_TOTAL_PIXELS), + "frame_factor": int(sglang_qwen_vl.FRAME_FACTOR), + "fps": float(sglang_qwen_vl.FPS), + "min_frames": int(sglang_qwen_vl.FPS_MIN_FRAMES), + "max_frames": int(sglang_qwen_vl.FPS_MAX_FRAMES), + }, + } + + +def publish_sglang_qwen_video_processor_contract( + runtime_config: ModelRuntimeConfig, engine: Any +) -> None: + """Publish exact Qwen video preprocessing behavior for the frontend.""" + contract = _resolve_qwen_video_processor_contract(engine) + if contract is not None: + runtime_config.set_engine_specific( + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + json.dumps(contract), + ) diff --git a/container/context.yaml b/container/context.yaml index f83db2faa303..78805b2223b0 100644 --- a/container/context.yaml +++ b/container/context.yaml @@ -140,8 +140,8 @@ sglang: # (framework, python stack, system libs) rather than hiding it. Per-arch stem; # the licenses stage appends -${TARGETARCH}.cdx.json. SGLang's codec-bearing # Python wheels are purged in sglang_runtime.Dockerfile; only the in-tree VP9 - # ffmpeg is copied in (CUDA, for the video-generation encode path) and the - # Dynamo Rust wheel is built without media-ffmpeg. + # ffmpeg is copied in. The CUDA Dynamo Rust wheel links against that same + # codec-limited build for frontend video decoding; XPU remains codec-free. baseline_sbom: cuda@0230b7f2 xpu: base_image: intel/deep-learning-essentials @@ -160,7 +160,7 @@ sglang: # NIXL Python stack — its wheel COPY is narrowed to ai_dynamo*.whl so the SDK # build doesn't leak into the runtime image. nixl_ref: v1.4.0 - enable_media_ffmpeg: "false" + enable_media_ffmpeg: "true" enable_gpu_memory_service: "true" enable_kvbm: "false" enable_modelexpress: "true" diff --git a/container/templates/sglang_runtime.Dockerfile b/container/templates/sglang_runtime.Dockerfile index 27e673b60657..e33897aed0b5 100644 --- a/container/templates/sglang_runtime.Dockerfile +++ b/container/templates/sglang_runtime.Dockerfile @@ -319,6 +319,12 @@ RUN set -eu; \ echo "ERROR: shipped ffmpeg ($ff) exposes an H.264/H.265/AAC/NVENC encoder" >&2; \ exit 1; \ fi + +# Frontend video decoding is part of the shipped SGLang CUDA contract. Fail the +# image build if the runtime wheel was accidentally compiled without it. +{% if target not in ("dev", "local-dev") %} +RUN python3 -c 'from dynamo.llm import MediaDecoder; assert hasattr(MediaDecoder(), "enable_video")' +{% endif %} {% else %} ENV IMAGEIO_FFMPEG_EXE= {% endif %} diff --git a/container/templates/wheel_builder.Dockerfile b/container/templates/wheel_builder.Dockerfile index 0ec910f80637..d34f6c574c66 100644 --- a/container/templates/wheel_builder.Dockerfile +++ b/container/templates/wheel_builder.Dockerfile @@ -574,9 +574,7 @@ COPY components/ /opt/dynamo/components/ # Build ai-dynamo (pure Python) and ai-dynamo-runtime (maturin) wheels ARG USE_SCCACHE ARG TARGETARCH -{% if framework != "sglang" %} ARG ENABLE_MEDIA_FFMPEG -{% endif %} RUN --mount=type=secret,id=aws-web-identity-token,target=/run/secrets/aws-token \ --mount=type=secret,id=aws-role-arn,env=AWS_ROLE_ARN \ --mount=type=cache,target=/root/.cargo/registry,sharing=shared \ @@ -593,7 +591,7 @@ RUN --mount=type=secret,id=aws-web-identity-token,target=/run/secrets/aws-token cd /opt/dynamo && \ uv build --wheel --out-dir /opt/dynamo/dist && \ cd /opt/dynamo/lib/bindings/python && \ -{% if framework == "sglang" %} maturin build --release --features "kv-indexer,slot-tracker,select-service,mm-routing,ais-forward-pass,request-trace-s3" --out /opt/dynamo/dist && \ +{% if framework == "sglang" and device == "xpu" %} maturin build --release --features "kv-indexer,slot-tracker,select-service,mm-routing,ais-forward-pass,request-trace-s3" --out /opt/dynamo/dist && \ {% else %} if [ "$ENABLE_MEDIA_FFMPEG" = "true" ]; then \ # Skip maturin's built-in repair: it would graft the in-tree libav* into the # wheel, which the codec gate rejects. Repair with those sonames excluded so diff --git a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md index 8c4cb9819b03..5ed8e62167f2 100644 --- a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md +++ b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md @@ -12,26 +12,24 @@ This document provides a comprehensive guide for multimodal inference using SGLa |----------|--------------|------------|---------------|-------| | **Image** | HTTP/HTTPS URL | Yes | Yes | Vision encoder generates embeddings | | **Image** | Data URL (Base64) | No | No | | -| **Video** | HTTP/HTTPS, `file://`, `data:` | No | Yes, H.264/H.265, Qwen2-family only | Needs the encode worker; decoded on NVDEC, then the vision encoder produces embeddings | +| **Video** | HTTP/HTTPS, `file://`, `data:` | Yes, VP8/VP9 on CUDA for Qwen3-VL and Qwen3.5 | Yes, H.264/H.265 for Qwen2-family only | Aggregated video requires `--frontend-decoding`; disaggregated video requires the encode worker | | **Audio** | HTTP/HTTPS URL | No | No | Not supported in SGLang backend | > [!IMPORTANT] -> **Video input is limited to H.264 and H.265, and requires a separate encode -> worker.** The runtime image ships no software video decoder, so H.264 and H.265 -> video is decoded on the GPU by NVDEC. Video in any other format (VP8, VP9, AV1) -> cannot be decoded at all. +> SGLang has two video-input paths with different codec and model support. > -> "Encode worker" here means Dynamo's `--disaggregation-mode encode` component, -> which runs the model's vision encoder to turn frames into embeddings. It does -> not encode video — nothing in this path produces a video stream. Hardware -> decode is wired into that worker (as used by `multimodal_epd.sh`). In an -> aggregated deployment SGLang resolves and decodes the media URL itself, so -> Dynamo never sees the bytes and cannot route them to NVDEC — video input is -> therefore unavailable in aggregated deployments of this image. +> - Aggregated Qwen3-VL and Qwen3.5 workers on CUDA support VP8 and VP9 with +> `--frontend-decoding`. Dynamo's Rust frontend decodes and transfers the +> sampled frames; SGLang still runs model-specific preprocessing and vision +> encoding. The multimodal router launcher uses this path for exact +> video-aware KV routing. +> - Disaggregated Qwen2-family deployments support H.264 and H.265 through the +> frontend-facing `--disaggregation-mode encode` worker. NVDEC performs the +> decode before the vision encoder produces embeddings. > -> Video is also skipped for model types whose preprocessing cannot accept -> pre-decoded frames (the Qwen3-VL family); those requests fall back to the URL -> path, which has no decoder in this image. Use a Qwen2-family vision model. +> The shipped SGLang image does not provide an AV1 input decoder. SGLang XPU +> images also omit the frontend FFmpeg decoder, so the aggregated VP8/VP9 path +> is CUDA-only. > > NVDEC requires a GPU with a video decode engine and a container granted the > `video` driver capability — see @@ -79,26 +77,33 @@ SGLang supports EPD, EP/D, E/PD, and E/P/D patterns. See [Multimodal Model Servi ### SGLang-Specific Characteristics -- **Vision Encoder in Python**: Encode worker uses SGLang's MMEncoder for model-agnostic vision encoding +- **Vision Encoder in Python**: SGLang runs model-specific vision encoding after the frontend or worker decodes the media - **Token Expansion**: Single `<|image_pad|>` token replaced with N tokens based on embedding shape - **NIXL Transfer**: Embeddings transferred from Encoder → PD Worker using NIXL -- **No Rust Processing**: All tokenization and image handling happens in Python +- **Optional Rust Frontend Decode**: `--frontend-decoding` moves media fetch and decode to the Rust frontend while retaining SGLang preprocessing and vision encoding ## Multimodal KV Routing Multimodal KV routing works with SGLang's aggregated worker topology. It is independent of the E/PD and E/P/D encoder-disaggregation patterns described later in this guide. -SGLang RadixAttention includes a per-image `pad_value` token in its prefix-cache key. Dynamo must use that same token in the routing view: +SGLang RadixAttention includes a per-media `pad_value` token in its prefix-cache key. Dynamo must use that same token in the routing view: -1. The Rust frontend hashes each image and calculates its expanded token count. +1. The Rust frontend hashes each image or sampled video and calculates its expanded token count. 2. It derives `pad_value = MM_PAD_SHIFT_VALUE + (mm_hash % 2^30)`. -3. It substitutes that value for the image placeholder in the routing-only token view. +3. It substitutes that value for the media placeholder in the routing-only token view. 4. It forwards the original hash list through `GenerateReqInput.mm_hashes`. 5. SGLang uses the supplied hash when constructing its own `pad_value`, keeping the router and RadixAttention cache keys aligned. The frontend selects this behavior automatically when the worker's `ModelDeploymentCard` reports `backend_framework="sglang"`. -Step 4 needs an SGLang build that accepts `mm_hashes`. Dynamo's SGLang image carries the upstream patch that adds it; a custom build without it is still supported, but Dynamo detects the missing argument when the worker starts and routes on the text prefix alone, so image identity no longer contributes to cache overlap. +For Qwen3-VL and Qwen3.5 video, `--frontend-decoding` also makes an aggregated +worker publish its effective frame-sampling, resize, timestamp, and placeholder +contract. The frontend enables exact video routing only when every worker in the +discovery cohort publishes the same compatible contract. Unsupported models, +legacy workers, mixed contracts, or a processor override fall back to +text-prefix routing without failing the request. + +Step 4 needs an SGLang build that accepts `mm_hashes`. Dynamo's SGLang image carries the upstream patch that adds it; a custom build without it is still supported, but Dynamo detects the missing argument when the worker starts and routes on the text prefix alone, so media identity no longer contributes to cache overlap. Launch an aggregated deployment with multimodal KV routing: @@ -107,7 +112,8 @@ cd $DYNAMO_HOME bash examples/backends/sglang/launch/agg_multimodal_router.sh ``` -The launcher configures KV events on each worker and sets `--router-mode kv` with a matching frontend and worker block size. +The launcher enables `--frontend-decoding`, configures KV events on each worker, +and sets `--router-mode kv` with a matching frontend and worker block size. | Variable | Default | Purpose | |----------|---------|---------| @@ -115,6 +121,7 @@ The launcher configures KV events on each worker and sets `--router-mode kv` wit | `NUM_WORKERS` | `2` | Number of SGLang workers | | `BLOCK_SIZE` | `16` | Worker page size and frontend KV block size | | `KV_EVENTS_PORT_BASE` | `29090` | Starting port for per-worker KV event publishers | +| `DYN_MM_VIDEO_NUM_FRAMES` | `32` | Maximum sampled video frames; keep this value consistent across the frontend and workers | | `SGLANG_EXTRA_ARGS` | unset | Additional arguments for `python -m dynamo.sglang` | ### Version Requirements @@ -123,13 +130,14 @@ The Dynamo SGLang image includes both routing prerequisites: - Dynamo is built with the `mm-routing` Rust feature. - SGLang 0.5.13 or later includes `GenerateReqInput.mm_hashes` support. Dynamo currently pins 0.5.19. +- The CUDA image includes the codec-limited VP8/VP9 frontend decoder. Custom installations on SGLang 0.5.12 or earlier need the `mm_hashes` change from [sgl-project/sglang#25300](https://github.com/sgl-project/sglang/pull/25300). Without it, requests still complete but fall back to text-prefix-only routing. Prefer upgrading to 0.5.13 or later instead of patching an installed package. -Enable `DYN_LOG=info,mm_routing=debug` to inspect image-token counts, multimodal hashes, and the selected worker's overlap. A repeated request should select the same worker with high block overlap. +Enable `DYN_LOG=info,mm_routing=debug` to inspect media-token counts, multimodal hashes, the video processor contract, and the selected worker's overlap. A repeated request should select the same worker with high block overlap. For the user-facing workflow, see [Multimodal KV Routing](../../../../../use-cases/multimodal-serving/multimodal-kv-routing.md). diff --git a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/reference-guide.md b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/reference-guide.md index 325da5bd4d5b..44c020c90b75 100644 --- a/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/reference-guide.md +++ b/docs/fern/pages/developer-guide/knowledge-base/modular-components/backends/sglang/reference-guide.md @@ -43,6 +43,7 @@ These arguments are added by Dynamo on top of SGLang's native arguments. For the | `--embedding-worker` | `DYN_SGL_EMBEDDING_WORKER` | `false` | Run as embedding worker (also sets SGLang's `--is-embedding`) | | `--enable-multimodal` | `DYN_SGL_ENABLE_MULTIMODAL` | `false` | Allow [multimodal](multimodal.md) inputs on this worker | | `--dedicated-mm-encoder` | `DYN_SGL_DEDICATED_MM_ENCODER` | `false` | Select the internal encode-worker topology for multimodal PD/P/D workers | +| `--frontend-decoding` | `DYN_SGL_FRONTEND_DECODING` | `false` | Decode images and supported videos in the Rust frontend and transfer pixels through NIXL. Use it on aggregated/native workers or the frontend-facing encode worker, not internal multimodal workers. | | `--image-diffusion-worker` | `DYN_SGL_IMAGE_DIFFUSION_WORKER` | `false` | Run as [image diffusion](../../../../../use-cases/diffusion/workflows/text-to-image.md#sglang) worker | | `--video-generation-worker` | `DYN_SGL_VIDEO_GENERATION_WORKER` | `false` | Run as [video generation](../../../../../use-cases/diffusion/workflows/text-to-video.md#sglang) worker | | `--disagg-config` | `DYN_SGL_DISAGG_CONFIG` | `None` | Path to YAML disaggregation config file | diff --git a/docs/fern/pages/reference/backends/sglang-configuration.mdx b/docs/fern/pages/reference/backends/sglang-configuration.mdx index 6f0a3fc196cf..9a6fe49bd15f 100644 --- a/docs/fern/pages/reference/backends/sglang-configuration.mdx +++ b/docs/fern/pages/reference/backends/sglang-configuration.mdx @@ -121,7 +121,7 @@ The SGLang wrapper selects a worker role from these flags and the native `--disa ## Multimodal decoding - Decode multimodal images in the Rust frontend and transfer the pre-decoded pixels to the backend via NIXL RDMA, bypassing in-engine HTTP fetch and decode. Incompatible with the dedicated encoder topology (`--disaggregation-mode=encode` or `--dedicated-mm-encoder`); use it on a native worker instead. + Decode multimodal images and supported videos in the Rust frontend and transfer the pre-decoded pixels to the backend via NIXL RDMA, bypassing in-engine HTTP fetch and decode. Use it on an aggregated or native worker, or on the frontend-facing `--disaggregation-mode=encode` worker. Do not set it on internal workers selected by `--dedicated-mm-encoder`. Environment variable: `DYN_SGL_FRONTEND_DECODING` @@ -151,7 +151,7 @@ This flag is retained for backward compatibility and will be removed in a future - `--disagg-config` and `--disagg-config-key` must both be set, or both unset. - `--dedicated-mm-encoder` requires `--enable-multimodal`. - `--dedicated-mm-encoder` is valid only with `--disaggregation-mode=pd`, `--disaggregation-mode=prefill`, or `--disaggregation-mode=decode`; do not combine it with `--disaggregation-mode=encode`. -- `--frontend-decoding` cannot be combined with the dedicated encoder topology (`--disaggregation-mode=encode` or `--dedicated-mm-encoder`). +- `--frontend-decoding` cannot be combined with internal EPD workers selected by `--dedicated-mm-encoder`; in an EPD deployment, set it only on the frontend-facing `--disaggregation-mode=encode` worker. ## Related pages diff --git a/docs/fern/pages/use-cases/multimodal-serving/additional-media-decoders.md b/docs/fern/pages/use-cases/multimodal-serving/additional-media-decoders.md index d0f3d370c981..286053336b58 100644 --- a/docs/fern/pages/use-cases/multimodal-serving/additional-media-decoders.md +++ b/docs/fern/pages/use-cases/multimodal-serving/additional-media-decoders.md @@ -9,10 +9,14 @@ Dynamo's runtime images ship a deliberately small media stack. The in-tree FFmpe Two input classes are already covered without installing anything: -- **VP8 and VP9 video** decode through the in-tree FFmpeg. -- **H.264 and H.265 video** decode on the GPU through NVDEC, which every backend uses by default. NVDEC needs a GPU with a video decode engine and a container granted the `video` driver capability. See [Video Decode GPU Requirements](video-decode-gpu-requirements.md) for the hardware and capability matrix. +- **VP8 and VP9 video** decode through the in-tree FFmpeg when frontend + decoding is enabled on vLLM, or SGLang on CUDA. +- **H.264 and H.265 video** decode on the GPU through NVDEC on supported backend paths. NVDEC needs a GPU with a video decode engine and a container granted the `video` driver capability. See [Video Decode GPU Requirements](video-decode-gpu-requirements.md) for the hardware and capability matrix. -Installing an additional decoder package covers what remains: +Without frontend decoding, the backend worker owns VP8/VP9 decoding. Follow the +[same-version replacement procedure](#notes-and-limits) for vLLM video; other +backends need the corresponding package from the table below. Installing an +additional decoder package also covers: - **AAC and other compressed audio**, which NVDEC does not decode at all. - **H.264 and H.265 on hosts where NVDEC is unavailable** — no video decode engine on the GPU, or a container without the `video` capability. @@ -98,7 +102,7 @@ The default pip timeout is 600 seconds (`--timeout-s` overrides it; `0` disables - For H.264 and H.265, prefer NVDEC. Granting the container the `video` driver capability decodes those formats on the GPU with no extra package. Install a software decoder when that is not an option, or when the input is audio. - Installing a decoder package brings in that wheel's bundled media libraries. The runtime images are scanned for media components at build time; a package installed afterwards is not covered by that scan. Review what your deployment ships — a baked image layer keeps the change visible and reviewable. - On TensorRT-LLM, the install puts back `opencv-python-headless`, which those images deliberately do not ship. H.264 and H.265 already decode there through NVDEC, so install it only for a host where NVDEC is unavailable. -- The vLLM images ship OpenCV already, rebuilt from source with every video backend disabled. It covers still images — which multimodal Mistral models need, because `mistral_common` resizes every image through `cv2` — and decodes no video at all. Video input on vLLM goes through NVDEC. To decode video in software instead, swap that build for the PyPI wheel of the same version: +- The vLLM images ship OpenCV already, rebuilt from source with every video backend disabled. It covers still images — which multimodal Mistral models need, because `mistral_common` resizes every image through `cv2` — and decodes no video at all. Backend video input on vLLM goes through NVDEC. To decode video in software instead, swap that build for the PyPI wheel of the same version: ```bash VERSION=$(pip show opencv-python-headless | awk '/^Version:/{print $2}') diff --git a/docs/fern/pages/use-cases/multimodal-serving/multimodal-kv-routing.md b/docs/fern/pages/use-cases/multimodal-serving/multimodal-kv-routing.md index 7f5315d196e2..df5e96474d93 100644 --- a/docs/fern/pages/use-cases/multimodal-serving/multimodal-kv-routing.md +++ b/docs/fern/pages/use-cases/multimodal-serving/multimodal-kv-routing.md @@ -45,9 +45,9 @@ URI string, while frontend decoding hashes the decoded bytes. Exact video routing requires frontend decoding. The frontend samples the video and computes an XXH3-64 identity over the sampled RGB frames, decoded shape and data type, and video metadata such as the source frame rate and sampled timestamps. Two URLs that decode to the same sampled content and metadata receive the same identity. -Video prompt expansion is model-specific. The frontend uses a model adapter to reproduce the worker's resize, patch, temporal grouping, delimiter, timestamp, and pruning behavior in a routing-only token sequence. At worker startup, vLLM reports the installed processor behavior that affects this sequence. Dynamo disables exact video routing if the reported contract does not match the frontend adapter. +Video prompt expansion is model-specific. The frontend uses a model adapter to reproduce the worker's resize, patch, temporal grouping, delimiter, timestamp, and pruning behavior in a routing-only token sequence. At worker startup, vLLM and eligible SGLang workers report the installed processor behavior that affects this sequence. Dynamo disables exact video routing if the reported contract does not match the frontend adapter. -The selected worker receives the same per-video hash, and its KV events use that identity. This keeps the frontend's routing blocks aligned with the blocks published by vLLM without running the vision model in the frontend. +The selected worker receives the same per-video hash, and its KV events use that identity. This keeps the frontend's routing blocks aligned with the blocks published by the backend without running the vision model in the frontend. @@ -65,7 +65,9 @@ The selected worker receives the same per-video hash, and its KV events use that The frontend hashes each image and converts that hash into the same `pad_value` token that SGLang RadixAttention uses for its prefix-cache key. The matching token view lets the router measure image overlap before selecting a worker. - Dynamo's SGLang image includes the required hash-forwarding support. Exact video routing is not available on this path; video requests use text-prefix routing. + For aggregated Qwen3-VL and Qwen3.5 workers on CUDA, frontend decoding extends the same mechanism to video. The worker publishes its effective frame-sampling, resize, timestamp, and placeholder contract. The frontend reproduces that contract in its routing-only token view and forwards the sampled video's hash to SGLang. + + Dynamo's SGLang image includes the required hash-forwarding support. Exact video routing fails closed to text-prefix routing when a worker lacks support, the discovery cohort contains legacy or mismatched contracts, or the installed processor behavior is incompatible. The frontend hashes each image, represents that identity in its routing token view, and forwards the hash as `multi_modal_uuids`. TensorRT-LLM workers publish matching KV events so the router can identify cached image blocks. @@ -115,7 +117,7 @@ The selected worker receives the same per-video hash, and its KV events use that bash examples/backends/sglang/launch/agg_multimodal_router.sh ``` - The launcher configures KV events and matching frontend and worker block sizes. + The launcher enables `--frontend-decoding`, configures KV events, and sets matching frontend and worker block sizes. Set `DYN_MM_VIDEO_NUM_FRAMES` before launching to change the maximum sampled frame count from its default of `32`; keep the same value for the frontend and every worker. See [SGLang Multimodal](../../developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md#multimodal-kv-routing) for prerequisites, configuration, fallback behavior, and verification. @@ -155,11 +157,11 @@ curl http://localhost:8000/v1/chat/completions \ JSON ``` -Send the request twice to populate and then reuse the worker's KV cache. Exact routing through the default Rust frontend supports VP8 or VP9 video in MP4, WebM, or MKV containers. The in-tree frontend decoder does not support H.264 or H.265. Re-encode those inputs to VP9, or disable frontend decoding and let the backend decode them through NVDEC or an installed software decoder; with backend-only decoding, video identity does not contribute to exact routing. See [Additional Media Decoders](additional-media-decoders.md) and [Video Decode GPU Requirements](video-decode-gpu-requirements.md). +Send the request twice to populate and then reuse the worker's KV cache. Exact routing through the default Rust frontend supports VP8 or VP9 video in MP4, WebM, or MKV containers. SGLang provides this decoder on CUDA only. The in-tree frontend decoder does not support H.264 or H.265. Re-encode those inputs to VP9, or use a backend-supported NVDEC or software-decoder path; with backend-only decoding, video identity does not contribute to exact routing. See [Additional Media Decoders](additional-media-decoders.md) and [Video Decode GPU Requirements](video-decode-gpu-requirements.md). ## Verify Video Routing -The vLLM multimodal router launcher enables routing debug logs by default. For a custom launch, set `DYN_LOG=info,mm_routing=debug,dynamo_kv_router::scheduling=debug,dynamo_llm::kv_router=debug`, then check for: +The vLLM and SGLang multimodal router launchers enable routing debug logs by default. For a custom launch, set `DYN_LOG=info,mm_routing=debug,dynamo_kv_router::scheduling=debug,dynamo_llm::kv_router=debug`, then check for: - `exact video-aware KV routing enabled` during worker registration. - `video routing metadata resolved` after the frontend decodes an eligible request. @@ -173,10 +175,11 @@ If startup or request validation cannot prove that the frontend and worker produ The default Rust frontend supports exact video routing for: -- Qwen3-VL and Qwen3.5 dense and mixture-of-experts models. Video pruning is not supported on this path. -- Nemotron 3 Nano Omni. The adapter reproduces temporal grouping, frame separators, dynamic resolution, and optional EVS pruning. +- vLLM Qwen3-VL and Qwen3.5 dense and mixture-of-experts models. Video pruning is not supported on this path. +- vLLM Nemotron 3 Nano Omni. The adapter reproduces temporal grouping, frame separators, dynamic resolution, and optional EVS pruning. +- Aggregated SGLang Qwen3-VL and Qwen3.5 dense and mixture-of-experts models on CUDA with `--frontend-decoding`. -Exact video routing falls back to text-prefix routing for unsupported models, opaque client-provided UUIDs, non-empty `mm_processor_kwargs`, audio in the same request, or adjacent video objects. Nemotron uses the same `` placeholder for image and video features, so its exact path currently supports one video-only media object per request; mixed image-and-video or multiple-video Nemotron requests fall back. +Exact video routing falls back to text-prefix routing for unsupported models, opaque client-provided UUIDs, non-empty `mm_processor_kwargs`, audio in the same request, or adjacent video objects. SGLang also falls back when workers publish missing, incompatible, or different processor contracts. Nemotron uses the same `` placeholder for image and video features, so its exact path currently supports one video-only media object per request; mixed image-and-video or multiple-video Nemotron requests fall back. ## Support Matrix @@ -184,7 +187,7 @@ Exact video routing falls back to text-prefix routing for unsupported models, op |---------|--------------|--------|-------| | [vLLM](../../developer-guide/knowledge-base/modular-components/backends/vllm/multimodal.md#multimodal-kv-routing) | Rust frontend (default) | Yes | Exact image routing supports Qwen2-VL, Qwen2.5-VL, Qwen3-VL, LLaVA 1.5, LLaVA-NeXT, Llama 4, Kimi K2.5/K2.6, Qwen3.5, Qwen3.6, and Nemotron 3 Nano Omni. Exact video routing requires frontend decoding and supports Qwen3-VL, Qwen3.5, and Nemotron 3 Nano Omni. Other media layouts use text-prefix routing. | | [vLLM](../../developer-guide/knowledge-base/modular-components/backends/vllm/multimodal.md#multimodal-kv-routing) | Python chat processor | Yes | Uses vLLM's multimodal processor to derive image and video feature hashes and positions for models supported by vLLM. | -| [SGLang](../../developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md#multimodal-kv-routing) | Rust frontend (default) | Images | Hash forwarding is upstream in SGLang 0.5.13+; Dynamo pins 0.5.19. Video requests use text-prefix routing. | +| [SGLang](../../developer-guide/knowledge-base/modular-components/backends/sglang/multimodal.md#multimodal-kv-routing) | Rust frontend (default) | Yes | Exact image routing uses SGLang's hash-forwarding support. Exact video routing supports aggregated Qwen3-VL and Qwen3.5 workers on CUDA with frontend decoding; other video requests use text-prefix routing. | | [TensorRT-LLM](../../developer-guide/knowledge-base/modular-components/backends/tensorrt-llm/multimodal.md#multimodal-kv-routing) | Rust frontend (default) | Images | Exact image routing supports the Qwen2-VL family and Kimi K2.5/K2.6. Video requests and other models use text-prefix routing. | Nemotron 3 Nano Omni exact image and video routing currently requires the vLLM backend. Other backends fall back to text-prefix routing for this model. diff --git a/docs/fern/pages/use-cases/multimodal-serving/parallel-media-decoding.md b/docs/fern/pages/use-cases/multimodal-serving/parallel-media-decoding.md index ddf50b733c0b..71908d4d506c 100644 --- a/docs/fern/pages/use-cases/multimodal-serving/parallel-media-decoding.md +++ b/docs/fern/pages/use-cases/multimodal-serving/parallel-media-decoding.md @@ -2,16 +2,16 @@ # SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. # SPDX-License-Identifier: Apache-2.0 title: Parallel Media Decoding -subtitle: Decode image inputs concurrently in the Rust frontend and transfer pixels to inference backends +subtitle: Decode media inputs concurrently in the Rust frontend and transfer pixels to inference backends --- -Parallel media decoding moves image fetching, base64 decoding, and image -decompression from the inference backend to the NVIDIA Dynamo Rust frontend. -The frontend decodes images concurrently on a CPU worker pool and transfers -the decoded pixel buffers to the backend through NIXL. +Parallel media decoding moves media fetching and decoding from the inference +backend to the NVIDIA Dynamo Rust frontend. The frontend decodes supported +images and videos on a CPU worker pool and transfers the decoded pixel buffers +to the backend through NIXL. The backend still runs its model-specific multimodal processor and vision -encoder. This feature changes where image input is decoded; it does not skip +encoder. This feature changes where media input is decoded; it does not skip vision encoding. ## Support Matrix @@ -19,37 +19,41 @@ vision encoding. | Input modality | vLLM | SGLang | TensorRT-LLM | | --- | --- | --- | --- | | Image | Agg | Agg | Agg | -| Video | Not supported | Not supported | Not supported | +| Video | Agg (VP8/VP9) | Agg (VP8/VP9) | Not supported | | Audio | Not supported | Not supported | Not supported | `Agg` refers to an aggregated worker. The entries in this matrix represent the supported topologies for frontend decoding. +SGLang frontend video decoding is available only on CUDA. SGLang XPU images do +not include the in-tree FFmpeg decoder. Image frontend decoding is not subject +to that codec restriction. + This matrix describes parallel media decoding, not the overall multimodal support of each backend. A backend can support video or audio by decoding it on the worker even when the frontend decoding path does not support that modality. ## When to Use -Use parallel media decoding when image preprocessing consumes a significant +Use parallel media decoding when media preprocessing consumes a significant part of request latency or backend CPU time. It is most useful for workloads with: -- Concurrent requests containing HTTP, HTTPS, or base64-encoded images -- Multiple images in one request -- Backend workers whose request path is constrained by image fetching or +- Concurrent requests containing HTTP, HTTPS, or base64-encoded media +- Multiple images or supported videos in one request +- Backend workers whose request path is constrained by media fetching or decompression Parallel media decoding can also be combined with the [embedding -cache](embedding-cache.md). Frontend decoding reduces image input processing +cache](embedding-cache.md). Frontend decoding reduces media input processing work, while the embedding cache can skip vision encoding for repeated images. ## How It Works For each request, the frontend: -1. Fetches the image URL or decodes the base64 data URL. -2. Decodes images concurrently on a CPU worker pool. +1. Fetches the media URL or decodes the base64 data URL. +2. Decodes supported images or VP8/VP9 videos on a CPU worker pool. 3. Registers the decoded pixel buffer with NIXL. 4. Sends the buffer descriptor to the selected backend worker. diff --git a/docs/fern/pages/use-cases/multimodal-serving/video-decode-gpu-requirements.md b/docs/fern/pages/use-cases/multimodal-serving/video-decode-gpu-requirements.md index d12a713da732..6b6d807768ac 100644 --- a/docs/fern/pages/use-cases/multimodal-serving/video-decode-gpu-requirements.md +++ b/docs/fern/pages/use-cases/multimodal-serving/video-decode-gpu-requirements.md @@ -4,17 +4,20 @@ title: Video Decode GPU Requirements --- -Dynamo decodes H.264 and H.265 (HEVC) video input on the GPU using NVDEC, NVIDIA's -dedicated hardware video decoder, through -[PyNvVideoCodec](https://pypi.org/project/PyNvVideoCodec/). - -Other formats — VP8, VP9 and AV1 — have **no video-input decoder** in the shipped images. -The in-tree VP8/VP9 FFmpeg serves the video *output* (generation) path; it is not wired to -video input, and the Rust `media-ffmpeg` decoder is not built into these images. Video -input decodes through Python carriers (OpenCV, PyAV, decord) that the images deliberately -omit — the vLLM images do ship OpenCV, but built without any video backend, for still-image -work only — so a VP8/VP9/AV1 clip fails with an unsupported-codec error unless a carrier -that decodes video is installed alongside. +Dynamo provides two video-input decode paths. CUDA runtime images can use both; +vLLM CPU and XPU runtimes can use the frontend software path: + +- H.264 and H.265 (HEVC) decode on the GPU using NVDEC, NVIDIA's dedicated + hardware video decoder, through + [PyNvVideoCodec](https://pypi.org/project/PyNvVideoCodec/). +- VP8 and VP9 decode on the CPU through Dynamo's codec-limited, in-tree FFmpeg + when frontend decoding is enabled on vLLM, or SGLang on CUDA. + +The in-tree FFmpeg does not include H.264, H.265, or AV1 decoders. Without +frontend decoding, video input remains owned by the backend and requires its +Python decode carrier (OpenCV, PyAV, or decord). The images deliberately omit +those wider software-decode carriers by default; vLLM does ship OpenCV, but it +is built without a video backend for still-image work only. This page covers which GPUs provide NVDEC, what the container must expose, and how Dynamo behaves when hardware decode is unavailable. @@ -129,20 +132,20 @@ between the two. ## Behavior when NVDEC is unavailable Hardware decode is additive and never blocks a request on its own: routing falls through -to the software decode path where one exists. +to a software decode path where one exists. VP8 and VP9 frontend decoding on vLLM, or +SGLang on CUDA, does not depend on NVDEC. > [!IMPORTANT] -> In the shipped images there is no software decode path for video input, for any format. -> The Python carriers that decode video input (OpenCV, PyAV, decord) are deliberately not -> installed — the vLLM images ship OpenCV built without a video backend, which resizes -> still images and opens no video — and the in-tree VP8/VP9 FFmpeg serves the video -> *output* path rather than input. So if NVDEC is unavailable, H.264 and H.265 fail with an unsupported-codec error -> — and VP8, VP9 and AV1 fail the same way whether NVDEC is available or not, since NVDEC -> does not decode them either. +> The shipped images do not include a software H.264, H.265, or AV1 decoder. +> Therefore, if NVDEC is unavailable, H.264 and H.265 fail with an +> unsupported-codec error unless the backend's wider Python decode carrier is +> installed. AV1 also requires an additional carrier because Dynamo does not +> route it through NVDEC and the in-tree FFmpeg excludes it. > > Grant the container the `video` driver capability so NVDEC can serve H.264 and H.265. -> For the other formats, install a decode carrier alongside, or transcode the input to -> H.264/H.265 before sending it. +> For VP8 and VP9, use frontend decoding on vLLM, or SGLang on CUDA. For +> other software-decoded cases, install a decode carrier alongside or transcode +> the input before sending it. ### Installing a software decoder @@ -203,7 +206,7 @@ NVDEC as above needs no additional change. Encode performance does not depend on | Variable | Default | Purpose | |----------|---------|---------| -| `DYN_DISABLE_NVDEC` | unset | Set to `1` to skip hardware decode. In a shipped image that leaves video input with no decoder at all, so it is a debugging switch rather than a fallback. Read as a boolean: `1`/`true`/`yes` disable, anything else does not. | +| `DYN_DISABLE_NVDEC` | unset | Set to `1` to skip hardware decode. In a shipped image that leaves H.264 and H.265 input with no decoder, so it is a debugging switch rather than a fallback for those codecs. VP8 and VP9 frontend decoding is unaffected. Read as a boolean: `1`/`true`/`yes` disable, anything else does not. | | `DYN_NVDEC_GPU_ID` | `0` | GPU ordinal used for decode. | | `DYN_MM_VIDEO_NUM_FRAMES` | `32` | Frames sampled uniformly from each clip. | | `DYN_MM_MAX_FILE_SIZE_MB` | `64` | Maximum size in MiB for each remote image, audio, or video download. TensorRT-LLM uses its backend-specific `--max-file-size-mb` option (`DYN_TRTLLM_MAX_FILE_SIZE_MB`) instead. | diff --git a/examples/backends/sglang/launch/agg_multimodal_router.sh b/examples/backends/sglang/launch/agg_multimodal_router.sh index 7d3146603cc5..cdc3b34cf4aa 100755 --- a/examples/backends/sglang/launch/agg_multimodal_router.sh +++ b/examples/backends/sglang/launch/agg_multimodal_router.sh @@ -127,6 +127,7 @@ for i in $(seq 1 "${NUM_WORKERS}"); do python -m dynamo.sglang \ --model-path "${MODEL}" \ --served-model-name "${MODEL}" \ + --frontend-decoding \ --page-size "${BLOCK_SIZE}" \ --context-length "${MAX_MODEL_LEN}" \ --tp 1 \ @@ -180,7 +181,7 @@ done echo echo "Architecture: Rust frontend (MM-aware KV router) -> ${NUM_WORKERS}x SGLang workers" echo " - mm_hashes forwarded to SGLang GenerateReqInput.mm_hashes -> matching pad_value" -echo " - Image dims via header-only HTTP fetch (Range: bytes=0-65535)" +echo " - Images and videos decoded once in the frontend and transferred over NIXL" echo " - No PyO3, no GIL, no Python deps in the routing path" echo echo "Press Ctrl+C to stop all services" diff --git a/lib/llm/src/discovery/watcher.rs b/lib/llm/src/discovery/watcher.rs index f5c7bb6adeef..a3355a253db8 100644 --- a/lib/llm/src/discovery/watcher.rs +++ b/lib/llm/src/discovery/watcher.rs @@ -31,8 +31,8 @@ use crate::{ kv_router::plugins::RouterPluginBuilder, kv_router::{EncoderRouter, PrefillRouter, RouterLoadSource, RoutingLoadContext}, local_model::runtime_config::{ - TokenizerBackend, VLLM_INFERENCE_V1_GENERATE_CAPABILITY, - VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, TokenizerBackend, + VLLM_INFERENCE_V1_GENERATE_CAPABILITY, VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, }, model_card::ModelDeploymentCard, model_type::{ModelInput, ModelType}, @@ -415,13 +415,16 @@ impl ModelWatcher { validate_policy_worker_role(card, &self.plugins)?; // Prepare without exact video routing unless the cohort agreed on a contract. - if spec.video_contract.is_none() - && card - .runtime_config - .runtime_data - .remove(VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY) - .is_some() - { + let mut removed_video_contract = false; + if spec.video_contract.is_none() { + for key in [ + VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ] { + removed_video_contract |= card.runtime_config.runtime_data.remove(key).is_some(); + } + } + if removed_video_contract { tracing::warn!( target: "mm_routing", model_name = card.name(), @@ -1389,13 +1392,21 @@ fn lora_projection_fingerprint(card: &ModelDeploymentCard) -> anyhow::Result Option { - let mut contract = card - .runtime_config - .runtime_data - .get(VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY)? - .clone(); - canonicalize_json(&mut contract); - Some(blake3::hash(contract.to_string().as_bytes()).to_string()) + let mut contracts = serde_json::Map::new(); + for key in [ + VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ] { + if let Some(contract) = card.runtime_config.runtime_data.get(key) { + contracts.insert(key.to_string(), contract.clone()); + } + } + if contracts.is_empty() { + return None; + } + let mut contracts = serde_json::Value::Object(contracts); + canonicalize_json(&mut contracts); + Some(blake3::hash(contracts.to_string().as_bytes()).to_string()) } fn canonicalize_json(value: &mut serde_json::Value) { @@ -1447,6 +1458,43 @@ mod tests { use dynamo_runtime::{Runtime, distributed::DistributedConfig}; use futures::StreamExt; + #[test] + fn qwen_video_contract_digest_is_canonical_and_engine_specific() { + fn card_with_contract(key: &str, contract: serde_json::Value) -> ModelDeploymentCard { + let mut card = ModelDeploymentCard::with_name_only("model"); + card.runtime_config + .runtime_data + .insert(key.to_string(), contract); + card + } + + let absent = ModelDeploymentCard::with_name_only("model"); + assert_eq!(qwen_video_contract_digest(&absent), None); + + let vllm = card_with_contract( + VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + serde_json::json!({"placeholder_target": "bare_video_token", "resize_mode": "round_ties_even"}), + ); + let reordered_vllm = card_with_contract( + VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + serde_json::json!({"resize_mode": "round_ties_even", "placeholder_target": "bare_video_token"}), + ); + let sglang = card_with_contract( + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + serde_json::json!({"placeholder_target": "bare_video_token", "resize_mode": "round_ties_even"}), + ); + + assert_eq!( + qwen_video_contract_digest(&vllm), + qwen_video_contract_digest(&reordered_vllm) + ); + assert_ne!( + qwen_video_contract_digest(&vllm), + qwen_video_contract_digest(&sglang), + "engine-specific contracts must not share a cohort fingerprint" + ); + } + #[tokio::test] async fn retired_worker_set_prevents_late_prefill_from_retained_chat_pipeline() { use crate::protocols::common::llm_backend::{BackendOutput, PreprocessedRequest}; diff --git a/lib/llm/src/local_model/runtime_config.rs b/lib/llm/src/local_model/runtime_config.rs index 9a7e8a9f4c42..7fd56f49bd73 100644 --- a/lib/llm/src/local_model/runtime_config.rs +++ b/lib/llm/src/local_model/runtime_config.rs @@ -100,6 +100,13 @@ pub const VLLM_INFERENCE_V1_GENERATE_CAPABILITY: &str = "vllm_inference_v1_gener pub const VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY: &str = "vllm_qwen_video_processor_contract"; +/// Worker-reported Qwen3 video prompt-expansion contract used by SGLang. +/// +/// SGLang performs an additional frame-selection and spatial-resize stage +/// before the Transformers processor, so this cannot share vLLM's contract. +pub const SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY: &str = + "sglang_qwen_video_processor_contract"; + /// Worker-reported Nemotron Nano Omni video prompt-expansion contract used by /// vLLM. Absence disables exact video routing for mixed-version safety. pub const VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY: &str = diff --git a/lib/llm/src/model_card.rs b/lib/llm/src/model_card.rs index 534e95052899..204b3607aa16 100644 --- a/lib/llm/src/model_card.rs +++ b/lib/llm/src/model_card.rs @@ -47,7 +47,7 @@ fn append_runtime_contract_checksum( canonicalize_json_object_keys(&mut value); let value = serde_json::to_vec(&value).expect("serializing serde_json::Value cannot fail"); - // These contracts control model-visible media prompt expansion. Workers + // This contract controls model-visible media prompt expansion. Workers // with different contracts must not share a cohort whose preprocessor is // built from one representative card. bytes.extend_from_slice(b"\0dynamo/model-card/runtime-contract/v1\0"); @@ -3375,6 +3375,7 @@ mod ownership_tests { #[test] fn video_processor_runtime_contract_checksum_boundaries() { use crate::local_model::runtime_config::{ + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, }; @@ -3409,6 +3410,23 @@ mod ownership_tests { "resize_mode": "round_ties_even" }), ); + let sglang_qwen = card_with_contract( + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + serde_json::json!({ + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil", + "sglang_preprocess": { + "image_factor": 28, + "video_min_pixels": 100352, + "video_max_pixels": 602112, + "video_total_pixels": 90316800, + "frame_factor": 2, + "fps": 2.0, + "min_frames": 4, + "max_frames": 768 + } + }), + ); let nemotron = card_with_contract( VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, serde_json::json!({"video_pruning_rate": 0.5}), @@ -3421,6 +3439,8 @@ mod ownership_tests { assert_eq!(missing.mdcsum(), unrelated.mdcsum()); assert_eq!(missing.mdcsum(), qwen.mdcsum()); + assert_eq!(missing.mdcsum(), sglang_qwen.mdcsum()); + assert_eq!(qwen.mdcsum(), sglang_qwen.mdcsum()); assert_eq!(qwen.mdcsum(), same_qwen.mdcsum()); assert_eq!(qwen.mdcsum(), different_qwen.mdcsum()); assert_ne!(missing.mdcsum(), nemotron.mdcsum()); diff --git a/lib/llm/src/preprocessor.rs b/lib/llm/src/preprocessor.rs index 5816b0b5040e..ae8802c581c5 100644 --- a/lib/llm/src/preprocessor.rs +++ b/lib/llm/src/preprocessor.rs @@ -54,12 +54,13 @@ use std::{ use tokio_util::sync::CancellationToken; use tracing::{self, Instrument}; -use crate::local_model::runtime_config::{TOKEN_BUDGET_RUNTIME_KEY, TokenBudget}; #[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] use crate::local_model::runtime_config::{ + ModelRuntimeConfig, SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, }; +use crate::local_model::runtime_config::{TOKEN_BUDGET_RUNTIME_KEY, TokenBudget}; #[cfg(feature = "mm-routing")] use crate::model_card::ModelInfoType; use crate::model_card::{ModelDeploymentCard, ModelInfo, PromptFormatterArtifact}; @@ -753,10 +754,9 @@ pub struct MmImageEntry { } /// One replacement tracked in both the worker-visible and canonical routing -/// token spaces. vLLM includes MM metadata on every block intersecting a -/// feature span, including timestamp/delimiter-only boundary blocks. Those -/// blocks need the worker token form plus `block_mm_infos`; blocks with an -/// exact placeholder/object mapping use the canonical pad-value form. +/// token spaces. Blocks with an exact placeholder/object mapping use the +/// canonical pad-value form. For a timestamp/delimiter-only boundary block, +/// the worker contract determines whether the hash also includes MM metadata. #[cfg(feature = "mm-routing")] #[derive(Debug, Clone, PartialEq, Eq)] struct TrackedMmRoutingReplacement { @@ -1020,10 +1020,9 @@ fn append_mm_routing_replacement_with_fill( /// normalizer block by block. /// /// Most blocks use canonical pad-value tokens. If a feature-span boundary -/// does not contain an exact ordered placeholder/object mapping, vLLM keeps -/// the worker tokens and hashes the block's MM metadata instead. Reproducing -/// that fallback here keeps both sides identical without discarding the media -/// identity carried by an ambiguous boundary block. +/// does not contain an exact ordered placeholder/object mapping, the worker's +/// KV-event identity contract determines whether to keep those canonical +/// tokens or fall back to worker tokens plus MM metadata. #[cfg(feature = "mm-routing")] fn apply_tracked_mm_replacements( routing_prepend_bos: Option, @@ -1032,6 +1031,7 @@ fn apply_tracked_mm_replacements( block_size: usize, image_token_id: Option, video_token_id: Option, + kv_event_mm_identity: mm_routing::KvEventMmIdentity, ) -> Result<( Vec, usize, @@ -1150,6 +1150,11 @@ fn apply_tracked_mm_replacements( "frontend MM replacement differs from KV-event normalization" ); } + None if kv_event_mm_identity == mm_routing::KvEventMmIdentity::PadValueTokens => { + // SGLang has already replaced every placeholder with the + // canonical media pad before publishing this KV-event block. + // The request-side routing block is therefore complete as-is. + } None => { routing_block.copy_from_slice(worker_block); block_mm_infos[block_index] = Some(BlockExtraInfo { @@ -1708,6 +1713,40 @@ pub struct OpenAIPreprocessor { pub(crate) const LORA_NAME_CONTEXT_KEY: &str = "discovery.lora_name"; +#[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] +fn resolve_qwen_video_processor_contract( + runtime_config: &ModelRuntimeConfig, +) -> Result> { + let vllm_contract = runtime_config + .get_engine_specific::( + VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ) + .with_context(|| { + format!( + "invalid Qwen video processor runtime metadata under {VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY}" + ) + })? + .map(mm_routing::QwenVideoProcessorContract::try_from) + .transpose()?; + let sglang_contract = runtime_config + .get_engine_specific::( + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ) + .with_context(|| { + format!( + "invalid Qwen video processor runtime metadata under {SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY}" + ) + })? + .map(mm_routing::QwenVideoProcessorContract::try_from) + .transpose()?; + + anyhow::ensure!( + vllm_contract.is_none() || sglang_contract.is_none(), + "multiple Qwen video processor contracts were published" + ); + Ok(vllm_contract.or(sglang_contract)) +} + impl OpenAIPreprocessor { fn omitted_max_tokens_default( prompt_len: usize, @@ -2559,16 +2598,12 @@ impl OpenAIPreprocessor { #[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] let video_routing_processor = { - let qwen_contract = match runtime_config - .get_engine_specific::( - VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, - ) { - Ok(target) => target, + let qwen_contract = match resolve_qwen_video_processor_contract(&runtime_config) { + Ok(contract) => contract, Err(error) => { tracing::warn!( target: "mm_routing", %error, - key = VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, "invalid Qwen video processor runtime metadata; exact video routing disabled" ); None @@ -4151,6 +4186,15 @@ impl OpenAIPreprocessor { // the frontend-tokenized prompt. let routing_bos = routing_bos_to_prepend(self.routing_prepend_bos, image_counter_required); + #[cfg(feature = "media-ffmpeg")] + let kv_event_mm_identity = self + .video_routing_processor + .as_ref() + .map_or(mm_routing::KvEventMmIdentity::MmMetadata, |processor| { + processor.kv_event_mm_identity() + }); + #[cfg(not(feature = "media-ffmpeg"))] + let kv_event_mm_identity = mm_routing::KvEventMmIdentity::MmMetadata; match apply_tracked_mm_replacements( routing_bos, &replacements, @@ -4158,6 +4202,7 @@ impl OpenAIPreprocessor { block_size, image_token_id, video_token_id, + kv_event_mm_identity, ) { Ok(expanded) => expanded, Err(error) => { @@ -12456,6 +12501,48 @@ mod tests { ); } + #[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] + #[test] + fn dual_qwen_video_contracts_disable_exact_routing() { + let mut runtime_config = ModelRuntimeConfig::default(); + runtime_config + .set_engine_specific( + VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + serde_json::json!({ + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil" + }), + ) + .unwrap(); + runtime_config + .set_engine_specific( + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + serde_json::json!({ + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil", + "runless_boundary_hash": "tokens_only", + "sglang_preprocess": { + "image_factor": 28, + "video_min_pixels": 100352, + "video_max_pixels": 602112, + "video_total_pixels": 90316800, + "frame_factor": 2, + "fps": 2.0, + "min_frames": 4, + "max_frames": 768 + } + }), + ) + .unwrap(); + + let error = resolve_qwen_video_processor_contract(&runtime_config).unwrap_err(); + assert!( + error + .to_string() + .contains("multiple Qwen video processor contracts") + ); + } + #[cfg(feature = "mm-routing")] #[test] fn exif_transposed_dimensions_match_vllm_image_loading() { @@ -12576,6 +12663,7 @@ mod tests { 4, Some(99), Some(video_token_id), + mm_routing::KvEventMmIdentity::MmMetadata, ) .unwrap(); @@ -12586,6 +12674,67 @@ mod tests { assert!(infos[1].is_none()); } + #[cfg(feature = "mm-routing")] + #[test] + fn tracked_video_boundary_matches_token_only_worker_contract() { + use dynamo_kv_router::protocols::pad_value_for_mm_hash; + + let video_token_id = 100; + let mm_hash = 41; + let video_pad = pad_value_for_mm_hash(mm_hash); + let replacement = TrackedMmRoutingReplacement { + mm_hash, + target_tokens: vec![9], + worker_tokens: vec![3, 4, 5, 6, video_token_id, video_token_id], + routing_tokens: vec![3, 4, 5, 6, video_pad, video_pad], + }; + + let (tokens, prompt_len, infos) = apply_tracked_mm_replacements( + None, + &[replacement], + &[1, 9, 2], + 4, + Some(99), + Some(video_token_id), + mm_routing::KvEventMmIdentity::PadValueTokens, + ) + .unwrap(); + + assert_eq!(prompt_len, 8); + assert_eq!(&tokens[..4], &[1, 3, 4, 5]); + assert_eq!(&tokens[4..], &[6, video_pad, video_pad, 2]); + assert!(infos.iter().all(Option::is_none)); + } + + #[cfg(feature = "mm-routing")] + #[test] + fn tracked_token_only_worker_still_validates_normalizable_blocks() { + let video_token_id = 100; + let replacement = TrackedMmRoutingReplacement { + mm_hash: 41, + target_tokens: vec![9], + worker_tokens: vec![video_token_id, video_token_id], + routing_tokens: vec![1, 2], + }; + + let error = apply_tracked_mm_replacements( + None, + &[replacement], + &[9], + 4, + Some(99), + Some(video_token_id), + mm_routing::KvEventMmIdentity::PadValueTokens, + ) + .unwrap_err(); + + assert!( + error + .to_string() + .contains("frontend MM replacement differs from KV-event normalization") + ); + } + #[cfg(feature = "mm-routing")] #[test] fn tracked_mixed_boundary_preserves_worker_hash_fallback() { @@ -12626,6 +12775,7 @@ mod tests { 4, Some(image_token_id), Some(video_token_id), + mm_routing::KvEventMmIdentity::MmMetadata, ) .unwrap(); @@ -12653,6 +12803,50 @@ mod tests { assert!(infos[1].is_none()); } + #[cfg(feature = "mm-routing")] + #[test] + fn tracked_mixed_token_only_boundary_preserves_canonical_pads() { + use dynamo_kv_router::protocols::pad_value_for_mm_hash; + + let image_token_id = 99; + let video_token_id = 100; + let image_hash = 41; + let video_hash = 42; + let image_pad = pad_value_for_mm_hash(image_hash); + let video_pad = pad_value_for_mm_hash(video_hash); + let replacements = [ + TrackedMmRoutingReplacement { + mm_hash: image_hash, + target_tokens: vec![image_token_id], + worker_tokens: vec![image_token_id; 14], + routing_tokens: vec![image_pad; 14], + }, + TrackedMmRoutingReplacement { + mm_hash: video_hash, + target_tokens: vec![video_token_id], + worker_tokens: vec![7, 8, video_token_id, video_token_id], + routing_tokens: vec![7, 8, video_pad, video_pad], + }, + ]; + + let (tokens, prompt_len, infos) = apply_tracked_mm_replacements( + None, + &replacements, + &[image_token_id, video_token_id], + 16, + Some(image_token_id), + Some(video_token_id), + mm_routing::KvEventMmIdentity::PadValueTokens, + ) + .unwrap(); + + assert_eq!(prompt_len, 18); + assert_eq!(&tokens[..14], &[image_pad; 14]); + assert_eq!(&tokens[14..18], &[7, 8, video_pad, video_pad]); + assert!(tokens[18..].iter().all(|token| *token == 0)); + assert!(infos.iter().all(Option::is_none)); + } + #[cfg(feature = "mm-routing")] #[test] fn tracked_replacements_preserve_image_video_image_order() { @@ -12682,6 +12876,7 @@ mod tests { 16, Some(image_token_id), Some(video_token_id), + mm_routing::KvEventMmIdentity::MmMetadata, ) .unwrap(); @@ -12724,6 +12919,7 @@ mod tests { 4, Some(10), Some(20), + mm_routing::KvEventMmIdentity::MmMetadata, ) .is_err(), "invalid target sequence {token_ids:?} must fail closed" @@ -12745,8 +12941,16 @@ mod tests { routing_tokens: vec![100, 19, pad, pad, 20, 101, 19, pad, 20], }; - let (tokens, prompt_len, block_infos) = - apply_tracked_mm_replacements(None, &[replacement], &[7], 4, Some(18), None).unwrap(); + let (tokens, prompt_len, block_infos) = apply_tracked_mm_replacements( + None, + &[replacement], + &[7], + 4, + Some(18), + None, + mm_routing::KvEventMmIdentity::MmMetadata, + ) + .unwrap(); assert_eq!(prompt_len, 9); assert_eq!(tokens, [100, 19, pad, pad, 20, 101, 19, pad, 20, 0, 0, 0]); diff --git a/lib/llm/src/preprocessor/mm_routing/mod.rs b/lib/llm/src/preprocessor/mm_routing/mod.rs index 941cdb242a12..10d2ad098cc1 100644 --- a/lib/llm/src/preprocessor/mm_routing/mod.rs +++ b/lib/llm/src/preprocessor/mm_routing/mod.rs @@ -36,11 +36,97 @@ pub(crate) enum QwenVideoResizeMode { RoundTiesEven, } -/// Worker-reported Qwen video prompt-expansion behavior. +/// How the worker hashes a block that intersects a video expansion but has no +/// video placeholder run. vLLM carries MM metadata in its KV event for these +/// boundary blocks; SGLang emits only the block's token IDs. #[derive(Debug, Clone, Copy, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub(crate) enum QwenVideoRunlessBoundaryHash { + MmMetadata, + TokensOnly, +} + +/// Engine-independent Qwen video prompt-expansion behavior used internally. +#[derive(Debug, Clone, Copy, PartialEq)] pub(crate) struct QwenVideoProcessorContract { pub placeholder_target: QwenVideoPlaceholderTarget, pub resize_mode: QwenVideoResizeMode, + pub runless_boundary_hash: QwenVideoRunlessBoundaryHash, + pub sglang_preprocess: Option, +} + +/// Qwen video contract accepted only under the vLLM runtime key. +#[derive(Debug, Clone, Copy, Deserialize, PartialEq)] +pub(crate) struct VllmQwenVideoProcessorContract { + pub placeholder_target: QwenVideoPlaceholderTarget, + pub resize_mode: QwenVideoResizeMode, + #[serde(default)] + runless_boundary_hash: Option, + #[serde(default)] + sglang_preprocess: Option, +} + +impl TryFrom for QwenVideoProcessorContract { + type Error = anyhow::Error; + + fn try_from(contract: VllmQwenVideoProcessorContract) -> Result { + anyhow::ensure!( + contract.sglang_preprocess.is_none(), + "mm-routing: vLLM Qwen contract must not publish SGLang preprocessing" + ); + anyhow::ensure!( + contract + .runless_boundary_hash + .unwrap_or(QwenVideoRunlessBoundaryHash::MmMetadata) + == QwenVideoRunlessBoundaryHash::MmMetadata, + "mm-routing: vLLM Qwen contract must use MM metadata KV-event identity" + ); + Ok(Self { + placeholder_target: contract.placeholder_target, + resize_mode: contract.resize_mode, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, + }) + } +} + +/// Qwen video contract accepted only under the SGLang runtime key. +#[derive(Debug, Clone, Copy, Deserialize, PartialEq)] +pub(crate) struct SglangQwenVideoProcessorContract { + pub placeholder_target: QwenVideoPlaceholderTarget, + pub resize_mode: QwenVideoResizeMode, + pub runless_boundary_hash: QwenVideoRunlessBoundaryHash, + pub sglang_preprocess: SglangQwenVideoPreprocessContract, +} + +impl TryFrom for QwenVideoProcessorContract { + type Error = anyhow::Error; + + fn try_from(contract: SglangQwenVideoProcessorContract) -> Result { + anyhow::ensure!( + contract.runless_boundary_hash == QwenVideoRunlessBoundaryHash::TokensOnly, + "mm-routing: SGLang Qwen contract must use pad-value token KV-event identity" + ); + Ok(Self { + placeholder_target: contract.placeholder_target, + resize_mode: contract.resize_mode, + runless_boundary_hash: contract.runless_boundary_hash, + sglang_preprocess: Some(contract.sglang_preprocess), + }) + } +} + +/// SGLang's Qwen video preprocessing stage before Transformers runs. +#[derive(Debug, Clone, Copy, Deserialize, PartialEq)] +pub(crate) struct SglangQwenVideoPreprocessContract { + pub image_factor: usize, + pub video_min_pixels: usize, + pub video_max_pixels: usize, + pub video_total_pixels: usize, + pub frame_factor: usize, + pub fps: f64, + pub min_frames: usize, + pub max_frames: usize, } /// Worker-reported Nemotron video prompt-expansion behavior. @@ -55,6 +141,17 @@ pub(crate) struct VideoProcessorContracts { pub nemotron: Option, } +/// How a worker represents multimodal identity in published KV-event blocks. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum KvEventMmIdentity { + /// Ambiguous blocks retain worker token IDs and carry media hashes as + /// separate block metadata. + MmMetadata, + /// Media placeholder runs are replaced with hash-derived pad values before + /// the worker publishes the block; no separate MM metadata is attached. + PadValueTokens, +} + /// Geometry and temporal metadata visible to a model's video processor. pub(crate) struct VideoRoutingInput<'a> { pub frame_count: usize, @@ -139,4 +236,103 @@ impl VideoRoutingProcessor { SupportedVideoModel::TestStub => anyhow::bail!("test video routing processor stub"), } } + + pub(crate) fn kv_event_mm_identity(&self) -> KvEventMmIdentity { + match &self.model { + SupportedVideoModel::Qwen3(spec) => spec.kv_event_mm_identity(), + SupportedVideoModel::Nemotron(_) => KvEventMmIdentity::MmMetadata, + #[cfg(test)] + SupportedVideoModel::TestStub => KvEventMmIdentity::MmMetadata, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const PUBLISHED_SGLANG_QWEN_CONTRACT: &str = r#"{ + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil", + "runless_boundary_hash": "tokens_only", + "sglang_preprocess": { + "image_factor": 28, + "video_min_pixels": 100352, + "video_max_pixels": 602112, + "video_total_pixels": 90316800, + "frame_factor": 2, + "fps": 2.0, + "min_frames": 4, + "max_frames": 768 + } + }"#; + + #[test] + fn parses_published_sglang_qwen_contract() { + let published: SglangQwenVideoProcessorContract = + serde_json::from_str(PUBLISHED_SGLANG_QWEN_CONTRACT).unwrap(); + let contract = QwenVideoProcessorContract::try_from(published).unwrap(); + + assert_eq!( + contract, + QwenVideoProcessorContract { + placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, + resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::TokensOnly, + sglang_preprocess: Some(SglangQwenVideoPreprocessContract { + image_factor: 28, + video_min_pixels: 100_352, + video_max_pixels: 602_112, + video_total_pixels: 90_316_800, + frame_factor: 2, + fps: 2.0, + min_frames: 4, + max_frames: 768, + }), + } + ); + } + + #[test] + fn rejects_incomplete_sglang_qwen_contract() { + let missing_preprocess = r#"{ + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil", + "runless_boundary_hash": "tokens_only" + }"#; + let missing_identity = r#"{ + "placeholder_target": "bare_video_token", + "resize_mode": "legacy_ceil", + "sglang_preprocess": { + "image_factor": 28, + "video_min_pixels": 100352, + "video_max_pixels": 602112, + "video_total_pixels": 90316800, + "frame_factor": 2, + "fps": 2.0, + "min_frames": 4, + "max_frames": 768 + } + }"#; + + assert!( + serde_json::from_str::(missing_preprocess).is_err() + ); + assert!( + serde_json::from_str::(missing_identity).is_err() + ); + } + + #[test] + fn rejects_cross_engine_qwen_contract_semantics() { + let vllm_with_sglang_preprocess: VllmQwenVideoProcessorContract = + serde_json::from_str(PUBLISHED_SGLANG_QWEN_CONTRACT).unwrap(); + assert!(QwenVideoProcessorContract::try_from(vllm_with_sglang_preprocess).is_err()); + + let sglang_with_vllm_identity = + PUBLISHED_SGLANG_QWEN_CONTRACT.replace("\"tokens_only\"", "\"mm_metadata\""); + let sglang_with_vllm_identity: SglangQwenVideoProcessorContract = + serde_json::from_str(&sglang_with_vllm_identity).unwrap(); + assert!(QwenVideoProcessorContract::try_from(sglang_with_vllm_identity).is_err()); + } } diff --git a/lib/llm/src/preprocessor/mm_routing/qwen3.rs b/lib/llm/src/preprocessor/mm_routing/qwen3.rs index 50d317185a21..d9d2873b31ad 100644 --- a/lib/llm/src/preprocessor/mm_routing/qwen3.rs +++ b/lib/llm/src/preprocessor/mm_routing/qwen3.rs @@ -1,13 +1,14 @@ // SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. // SPDX-License-Identifier: Apache-2.0 -use std::{path::Path, sync::Arc}; +use std::{borrow::Cow, path::Path, sync::Arc}; use anyhow::{Context, Result}; use serde_json::Value; use super::{ - QwenVideoPlaceholderTarget, QwenVideoProcessorContract, QwenVideoResizeMode, VideoRoutingInput, + KvEventMmIdentity, QwenVideoPlaceholderTarget, QwenVideoProcessorContract, QwenVideoResizeMode, + QwenVideoRunlessBoundaryHash, SglangQwenVideoPreprocessContract, VideoRoutingInput, VideoRoutingReplacement, config::{read_json, read_model_config, required_token_id, required_usize}, }; @@ -40,10 +41,45 @@ pub(super) struct Qwen3VideoRoutingSpec { vision_end_token_id: TokenIdType, placeholder_target: QwenVideoPlaceholderTarget, resize_mode: QwenVideoResizeMode, + runless_boundary_hash: QwenVideoRunlessBoundaryHash, + sglang_preprocess: Option, tokenizer: Arc, } +struct PreparedVideoInput<'a> { + frame_count: usize, + width: u32, + height: u32, + source_fps: f64, + sampled_timestamps: Cow<'a, [f64]>, +} + +impl SglangQwenVideoPreprocessContract { + fn validate(&self) -> Result<()> { + anyhow::ensure!( + self.image_factor > 0 + && self.video_min_pixels > 0 + && self.video_max_pixels >= self.video_min_pixels + && self.video_total_pixels > 0 + && self.frame_factor > 0 + && self.min_frames > 0 + && self.max_frames >= self.min_frames + && self.fps.is_finite() + && self.fps > 0.0, + "mm-routing: invalid SGLang Qwen video preprocessing contract" + ); + Ok(()) + } +} + impl Qwen3VideoRoutingSpec { + pub(super) fn kv_event_mm_identity(&self) -> KvEventMmIdentity { + match self.runless_boundary_hash { + QwenVideoRunlessBoundaryHash::MmMetadata => KvEventMmIdentity::MmMetadata, + QwenVideoRunlessBoundaryHash::TokensOnly => KvEventMmIdentity::PadValueTokens, + } + } + pub(super) fn from_model_dir( model_id: &str, expected_model_type: &str, @@ -51,6 +87,9 @@ impl Qwen3VideoRoutingSpec { tokenizer: Arc, processor_contract: QwenVideoProcessorContract, ) -> Result { + if let Some(contract) = processor_contract.sglang_preprocess { + contract.validate()?; + } let expected_architecture = expected_architecture(expected_model_type) .context("mm-routing: Qwen video model_type has no registered architecture")?; let model_config = read_model_config( @@ -128,6 +167,8 @@ impl Qwen3VideoRoutingSpec { vision_end_token_id: required_token_id(&model_config, "vision_end_token_id", "Qwen")?, placeholder_target: processor_contract.placeholder_target, resize_mode: processor_contract.resize_mode, + runless_boundary_hash: processor_contract.runless_boundary_hash, + sglang_preprocess: processor_contract.sglang_preprocess, tokenizer, }) } @@ -137,7 +178,15 @@ impl Qwen3VideoRoutingSpec { input: &VideoRoutingInput<'_>, ) -> Result { self.validate_input(input)?; - let (grid_t, grid_h, grid_w) = self.video_grid(input)?; + let prepared = self.prepare_input(input)?; + let prepared_input = VideoRoutingInput { + frame_count: prepared.frame_count, + width: prepared.width, + height: prepared.height, + source_fps: prepared.source_fps, + sampled_timestamps: prepared.sampled_timestamps.as_ref(), + }; + let (grid_t, grid_h, grid_w) = self.video_grid(&prepared_input)?; let merge_area = self .spatial_merge_size .checked_mul(self.spatial_merge_size) @@ -154,7 +203,7 @@ impl Qwen3VideoRoutingSpec { .checked_mul(tokens_per_grid) .context("mm-routing: Qwen video token count overflow")?; - let grid_timestamps = self.grid_timestamps(input, grid_t)?; + let grid_timestamps = self.grid_timestamps(&prepared_input, grid_t)?; let mut replacement_tokens = Vec::with_capacity( base_video_tokens .checked_add(grid_t.saturating_mul(8)) @@ -188,6 +237,99 @@ impl Qwen3VideoRoutingSpec { }) } + fn prepare_input<'a>(&self, input: &VideoRoutingInput<'a>) -> Result> { + let Some(contract) = self.sglang_preprocess else { + return Ok(PreparedVideoInput { + frame_count: input.frame_count, + width: input.width, + height: input.height, + source_fps: input.source_fps, + sampled_timestamps: Cow::Borrowed(input.sampled_timestamps), + }); + }; + + anyhow::ensure!( + input.frame_count >= contract.frame_factor, + "mm-routing: SGLang Qwen video frame count is below frame_factor" + ); + + let first_frame = (input.sampled_timestamps[0] * input.source_fps).round_ties_even(); + let last_frame = + (input.sampled_timestamps[input.frame_count - 1] * input.source_fps).round_ties_even(); + let span_frames = last_frame - first_frame; + anyhow::ensure!( + input.frame_count == 1 || span_frames > 0.0, + "mm-routing: SGLang Qwen sampled video has no positive frame span" + ); + let effective_fps = if input.frame_count > 1 { + (input.frame_count - 1) as f64 * input.source_fps / span_frames + } else { + input.source_fps + }; + anyhow::ensure!( + effective_fps.is_finite() && effective_fps > 0.0, + "mm-routing: SGLang Qwen effective fps is invalid" + ); + + let min_frames = ceil_to_factor(contract.min_frames, contract.frame_factor)?; + let max_frames = floor_to_factor( + contract.max_frames.min(input.frame_count), + contract.frame_factor, + )?; + let requested_frames = input.frame_count as f64 / effective_fps * contract.fps; + // Match SGLang's min(max(requested, min_frames), max_frames) order. + // For a 2-3 frame clip, max_frames can be below the configured + // minimum and SGLang intentionally collapses the result to 2. + let requested_frames = requested_frames + .max(min_frames as f64) + .min(max_frames as f64) + .min(input.frame_count as f64); + let frame_count = floor_to_factor(requested_frames as usize, contract.frame_factor)?; + anyhow::ensure!( + frame_count >= contract.frame_factor && frame_count <= input.frame_count, + "mm-routing: SGLang Qwen selected frame count is invalid" + ); + + // Match np.linspace(0, total_frames - 1, num=frame_count, dtype=int64). + let sampled_timestamps = (0..frame_count) + .map(|index| { + let source_index = if frame_count == 1 { + 0 + } else if index == frame_count - 1 { + input.frame_count - 1 + } else { + let step = (input.frame_count - 1) as f64 / (frame_count - 1) as f64; + (index as f64 * step).floor() as usize + }; + source_index as f64 / effective_fps + }) + .collect(); + + let min_pixels = contract.video_min_pixels as f64; + let max_pixels = (contract.video_max_pixels as f64) + .min( + contract.video_total_pixels as f64 / frame_count as f64 + * contract.frame_factor as f64, + ) + .max((min_pixels * 1.05) as usize as f64); + let (height, width) = sglang_smart_resize( + usize::try_from(input.height).context("mm-routing: video height exceeds usize")?, + usize::try_from(input.width).context("mm-routing: video width exceeds usize")?, + contract.image_factor, + contract.video_min_pixels as f64, + max_pixels, + )?; + + Ok(PreparedVideoInput { + frame_count, + width: u32::try_from(width).context("mm-routing: resized video width exceeds u32")?, + height: u32::try_from(height) + .context("mm-routing: resized video height exceeds u32")?, + source_fps: effective_fps, + sampled_timestamps: Cow::Owned(sampled_timestamps), + }) + } + fn validate_input(&self, input: &VideoRoutingInput<'_>) -> Result<()> { anyhow::ensure!( input.frame_count > 0, @@ -351,6 +493,65 @@ fn ensure_matching_value(config: &Value, field: &str, expected: usize) -> Result Ok(()) } +fn ceil_to_factor(value: usize, factor: usize) -> Result { + anyhow::ensure!(factor > 0, "mm-routing: SGLang Qwen frame_factor is zero"); + value + .checked_add(factor - 1) + .map(|value| value / factor * factor) + .context("mm-routing: SGLang Qwen frame rounding overflow") +} + +fn floor_to_factor(value: usize, factor: usize) -> Result { + anyhow::ensure!(factor > 0, "mm-routing: SGLang Qwen frame_factor is zero"); + Ok(value / factor * factor) +} + +fn sglang_smart_resize( + height: usize, + width: usize, + factor: usize, + min_pixels: f64, + max_pixels: f64, +) -> Result<(usize, usize)> { + anyhow::ensure!( + height > 0 && width > 0, + "mm-routing: video dimensions are zero" + ); + let aspect_ratio = height.max(width) as f64 / height.min(width) as f64; + anyhow::ensure!( + aspect_ratio <= 200.0, + "mm-routing: SGLang Qwen video aspect ratio exceeds 200:1" + ); + + let mut resized_height = + ((height as f64 / factor as f64).round_ties_even() as usize * factor).max(factor); + let mut resized_width = + ((width as f64 / factor as f64).round_ties_even() as usize * factor).max(factor); + let resized_pixels = resized_height + .checked_mul(resized_width) + .context("mm-routing: SGLang Qwen resized area overflow")? as f64; + let source_pixels = height + .checked_mul(width) + .context("mm-routing: SGLang Qwen source area overflow")? as f64; + + if resized_pixels > max_pixels { + let beta = (source_pixels / max_pixels).sqrt(); + resized_height = (height as f64 / beta / factor as f64).floor() as usize * factor; + resized_width = (width as f64 / beta / factor as f64).floor() as usize * factor; + } else if resized_pixels < min_pixels { + let beta = (min_pixels / source_pixels).sqrt(); + resized_height = (height as f64 * beta / factor as f64).ceil() as usize * factor; + resized_width = (width as f64 * beta / factor as f64).ceil() as usize * factor; + } + + anyhow::ensure!( + resized_height > 0 && resized_width > 0, + "mm-routing: SGLang Qwen smart resize produced a zero dimension" + ); + + Ok((resized_height, resized_width)) +} + #[cfg(test)] mod tests { use crate::tokenizers::{Encoding, traits::DecodeResult}; @@ -399,6 +600,8 @@ mod tests { vision_end_token_id: 151653, placeholder_target: QwenVideoPlaceholderTarget::VisionWrappedVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, tokenizer: Arc::new(TimestampTokenizer), } } @@ -415,10 +618,223 @@ mod tests { vision_end_token_id: 151653, placeholder_target: QwenVideoPlaceholderTarget::VisionWrappedVideoToken, resize_mode, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, tokenizer: Arc::new(TimestampTokenizer), } } + fn sglang_preprocess_contract() -> SglangQwenVideoPreprocessContract { + SglangQwenVideoPreprocessContract { + image_factor: 28, + video_min_pixels: 100_352, + video_max_pixels: 602_112, + video_total_pixels: 90_316_800, + frame_factor: 2, + fps: 2.0, + min_frames: 4, + max_frames: 768, + } + } + + #[test] + fn vllm_input_path_borrows_timestamps() { + let timestamps = [0.0, 1.0]; + let input = VideoRoutingInput { + frame_count: 2, + width: 224, + height: 224, + source_fps: 1.0, + sampled_timestamps: ×tamps, + }; + + let prepared = production_geometry_spec(QwenVideoResizeMode::LegacyCeil) + .prepare_input(&input) + .unwrap(); + + assert!(matches!(prepared.sampled_timestamps, Cow::Borrowed(_))); + } + + #[test] + fn sglang_short_clips_clamp_to_available_frame_factor() { + let mut spec = production_geometry_spec(QwenVideoResizeMode::LegacyCeil); + spec.sglang_preprocess = Some(sglang_preprocess_contract()); + + for frame_count in [2, 3] { + let timestamps: Vec<_> = (0..frame_count).map(|index| index as f64 / 30.0).collect(); + let input = VideoRoutingInput { + frame_count, + width: 320, + height: 240, + source_fps: 30.0, + sampled_timestamps: ×tamps, + }; + + assert_eq!(spec.prepare_input(&input).unwrap().frame_count, 2); + } + } + + #[test] + fn sglang_rejects_frames_below_frame_factor() { + let timestamps = [0.0]; + let input = VideoRoutingInput { + frame_count: 1, + width: 320, + height: 240, + source_fps: 30.0, + sampled_timestamps: ×tamps, + }; + let mut spec = production_geometry_spec(QwenVideoResizeMode::LegacyCeil); + spec.sglang_preprocess = Some(sglang_preprocess_contract()); + + assert!(spec.prepare_input(&input).is_err()); + } + + #[test] + fn rejects_invalid_sglang_contract_during_spec_creation() { + let mut invalid = sglang_preprocess_contract(); + invalid.frame_factor = 0; + let contract = QwenVideoProcessorContract { + placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, + resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::TokensOnly, + sglang_preprocess: Some(invalid), + }; + + let error = Qwen3VideoRoutingSpec::from_model_dir( + "Qwen/Qwen3-VL-2B-Instruct", + "qwen3_vl", + tempfile::tempdir().unwrap().path(), + Arc::new(TimestampTokenizer), + contract, + ) + .err() + .expect("invalid SGLang contract must disable exact routing at startup"); + + assert!(error.to_string().contains("invalid SGLang Qwen")); + } + + #[test] + fn ceil_to_factor_rejects_zero_factor() { + assert!(ceil_to_factor(4, 0).is_err()); + } + + #[test] + fn reproduces_sglang_frame_sampling_and_pre_resize() { + let timestamps: Vec<_> = (0..32).map(|index| index as f64 * 10.0 / 31.0).collect(); + let input = VideoRoutingInput { + frame_count: 32, + width: 320, + height: 240, + source_fps: 30.0, + sampled_timestamps: ×tamps, + }; + let mut spec = production_geometry_spec(QwenVideoResizeMode::LegacyCeil); + spec.sglang_preprocess = Some(SglangQwenVideoPreprocessContract { + image_factor: 28, + video_min_pixels: 128 * 28 * 28, + video_max_pixels: 768 * 28 * 28, + video_total_pixels: (128000.0 * 28.0 * 28.0 * 0.9) as usize, + frame_factor: 2, + fps: 2.0, + min_frames: 4, + max_frames: 768, + }); + + let prepared = spec.prepare_input(&input).unwrap(); + + assert_eq!(prepared.frame_count, 20); + assert_eq!((prepared.width, prepared.height), (392, 280)); + assert!((prepared.source_fps - 3.1).abs() < 1e-9); + assert_eq!(prepared.sampled_timestamps.first(), Some(&0.0)); + assert!((prepared.sampled_timestamps.last().unwrap() - 10.0).abs() < 1e-9); + let prepared_input = VideoRoutingInput { + frame_count: prepared.frame_count, + width: prepared.width, + height: prepared.height, + source_fps: prepared.source_fps, + sampled_timestamps: &prepared.sampled_timestamps, + }; + assert_eq!(spec.video_grid(&prepared_input).unwrap(), (10, 18, 24)); + } + + #[test] + fn sglang_frame_sampling_matches_numpy_float_linspace() { + let timestamps: Vec<_> = (0..46) + .map(|index| index as f64 * 510.0 / 45.0 / 30.0) + .collect(); + let input = VideoRoutingInput { + frame_count: 46, + width: 320, + height: 240, + source_fps: 30.0, + sampled_timestamps: ×tamps, + }; + let mut spec = production_geometry_spec(QwenVideoResizeMode::LegacyCeil); + spec.sglang_preprocess = Some(SglangQwenVideoPreprocessContract { + image_factor: 28, + video_min_pixels: 128 * 28 * 28, + video_max_pixels: 768 * 28 * 28, + video_total_pixels: (128000.0 * 28.0 * 28.0 * 0.9) as usize, + frame_factor: 2, + fps: 2.0, + min_frames: 4, + max_frames: 768, + }); + + let prepared = spec.prepare_input(&input).unwrap(); + let selected_indices: Vec<_> = prepared + .sampled_timestamps + .iter() + .map(|timestamp| (timestamp * prepared.source_fps).round() as usize) + .collect(); + + assert_eq!(prepared.frame_count, 34); + assert_eq!( + selected_indices, + [ + 0, 1, 2, 4, 5, 6, 8, 9, 10, 12, 13, 14, 16, 17, 19, 20, 21, 23, 24, 25, 27, 28, 29, + 31, 32, 34, 35, 36, 38, 39, 40, 42, 43, 45, + ] + ); + } + + #[test] + fn sglang_pre_resize_uses_ties_to_even() { + assert_eq!( + sglang_smart_resize(406, 700, 28, 100_352.0, 602_112.0).unwrap(), + (392, 700) + ); + } + + #[test] + fn sglang_long_video_applies_total_pixel_budget() { + let timestamps: Vec<_> = (0..400).map(|index| index as f64 / 2.0).collect(); + let input = VideoRoutingInput { + frame_count: 400, + width: 1920, + height: 1080, + source_fps: 2.0, + sampled_timestamps: ×tamps, + }; + let mut spec = production_geometry_spec(QwenVideoResizeMode::LegacyCeil); + spec.sglang_preprocess = Some(SglangQwenVideoPreprocessContract { + image_factor: 28, + video_min_pixels: 100_352, + video_max_pixels: 602_112, + video_total_pixels: 90_316_800, + frame_factor: 2, + fps: 2.0, + min_frames: 4, + max_frames: 768, + }); + + let prepared = spec.prepare_input(&input).unwrap(); + + assert_eq!(prepared.frame_count, 400); + assert_eq!((prepared.height, prepared.width), (504, 896)); + } + #[test] fn builds_timestamped_qwen_replacement_without_preprocessing_pixels() { let timestamps = [0.0, 1.0, 2.0, 3.0]; @@ -535,6 +951,8 @@ mod tests { vision_end_token_id: 151653, placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, tokenizer: Arc::new(CheckpointTimestampTokenizer { seconds_token_id: 6486, }), @@ -746,6 +1164,8 @@ mod tests { QwenVideoProcessorContract { placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, }, ) .unwrap(); @@ -797,6 +1217,8 @@ mod tests { QwenVideoProcessorContract { placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, }, ) .is_err() @@ -844,6 +1266,8 @@ mod tests { QwenVideoProcessorContract { placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, }, ) .err() diff --git a/tests/serve/test_sglang.py b/tests/serve/test_sglang.py index a4b7f028f80f..829950246229 100644 --- a/tests/serve/test_sglang.py +++ b/tests/serve/test_sglang.py @@ -50,6 +50,7 @@ router_selection_chat_payload_default, ) from tests.utils.payloads import ( + CachedTokensChatPayload, ChatPayload, HttpErrorPayload, ImageGenerationPayload, @@ -683,44 +684,63 @@ class SGLangConfig(EngineConfig): "video_agg_fd_qwen": SGLangConfig( name="video_agg_fd_qwen", directory=sglang_dir, - script_name="agg_vision.sh", + script_name="agg_multimodal_router.sh", marks=[ pytest.mark.multimodal, pytest.mark.gpu_1, - pytest.mark.profiled_vram_gib(10.0), + pytest.mark.profiled_vram_gib(18.7), pytest.mark.requested_sglang_kv_tokens(8736), - pytest.mark.timeout(390), + pytest.mark.timeout(500), pytest.mark.pre_merge, - # TODO: Enable media-ffmpeg in the SGLang container build, then - # remove this skip. Frontend video decoding requires the Dynamo - # binding to be built with media-ffmpeg support. - pytest.mark.skip(reason="SGLang container lacks media-ffmpeg support"), ], model="Qwen/Qwen3-VL-2B-Instruct", script_args=[ - "--model-path", + "--model", "Qwen/Qwen3-VL-2B-Instruct", - "--frontend-decoding", + "--num-workers", + "2", + "--single-gpu", ], env={ "DYN_MM_ALLOW_INTERNAL": "1", - "DYN_MM_VIDEO_NUM_FRAMES": "4", + # Decode all 10 frames in the fixture so the routing sequence is + # materially larger than the shared text prefix. + "DYN_MM_VIDEO_NUM_FRAMES": "10", }, - timeout=360, + timeout=450, frontend_port=DefaultPort.FRONTEND.value, request_payloads=[ - chat_payload( - [ - {"type": "text", "text": "Describe the video in detail"}, - { - "type": "video_url", - "video_url": {"url": MULTIMODAL_VIDEO_URL}, - }, - ], - repeat_count=1, + CachedTokensChatPayload( + body={ + "messages": [ + { + "role": "user", + "content": [ + { + "type": "text", + "text": "Describe the video in detail", + }, + { + "type": "video_url", + "video_url": {"url": MULTIMODAL_VIDEO_URL}, + }, + ], + } + ], + "max_tokens": 100, + "temperature": 0.0, + "stream": False, + }, + repeat_count=3, expected_response=MULTIMODAL_VIDEO_EXPECTED, - temperature=0.0, - max_tokens=100, + min_cached_tokens=128, + require_rust_processor_init=True, + min_routing_total_blocks=10, + # A text-only hit is at most a small fraction of this video + # request. Requiring high router-side overlap proves the media + # hashes matched the worker KV events instead of accepting a + # cached text prefix as a false positive. + min_avg_kv_hit_rate=0.9, ) ], ),