From a62a07270540df3d68c88e00389e8673273aebef Mon Sep 17 00:00:00 2001 From: krishung5 Date: Thu, 17 Sep 2026 12:09:57 -0700 Subject: [PATCH 1/6] feat(mm-routing): add Qwen video routing for SGLang Signed-off-by: krishung5 --- components/src/dynamo/sglang/register.py | 12 + .../request_handlers/llm/decode_handler.py | 31 +-- .../request_handlers/llm/mm_disagg_utils.py | 61 +++++ .../tests/test_sglang_frontend_decoding.py | 28 ++ .../tests/test_sglang_multimodal_utils.py | 36 +++ .../sglang/tests/test_sglang_video_routing.py | 138 ++++++++++ components/src/dynamo/sglang/video_routing.py | 157 +++++++++++ container/context.yaml | 6 +- container/templates/sglang_runtime.Dockerfile | 4 + container/templates/wheel_builder.Dockerfile | 4 +- .../additional-media-decoders.md | 7 +- .../parallel-media-decoding.md | 16 +- .../video-decode-gpu-requirements.md | 44 +-- .../sglang/launch/agg_multimodal_router.sh | 3 +- lib/llm/src/local_model/runtime_config.rs | 7 + lib/llm/src/model_card.rs | 24 +- lib/llm/src/preprocessor.rs | 128 +++++++-- lib/llm/src/preprocessor/mm_routing/mod.rs | 33 ++- .../src/preprocessor/mm_routing/nemotron.rs | 1 + lib/llm/src/preprocessor/mm_routing/qwen3.rs | 253 +++++++++++++++++- tests/serve/test_sglang.py | 64 +++-- 21 files changed, 942 insertions(+), 115 deletions(-) create mode 100644 components/src/dynamo/sglang/tests/test_sglang_video_routing.py create mode 100644 components/src/dynamo/sglang/video_routing.py diff --git a/components/src/dynamo/sglang/register.py b/components/src/dynamo/sglang/register.py index 6fb2f1afb4ee..d74250c7d954 100644 --- a/components/src/dynamo/sglang/register.py +++ b/components/src/dynamo/sglang/register.py @@ -41,6 +41,9 @@ runtime_capacity, ) from dynamo.sglang.engine_generate import SGLANG_GENERATE_CAPABILITY +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" @@ -405,6 +408,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 f92ee7493b39..a7c61fcbb684 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 @@ -40,6 +40,7 @@ IMAGE_URL_KEY, VIDEO_URL_KEY, build_disagg_mm_kwargs, + extract_mm_hashes, extract_media_urls, raise_if_unextracted_multimodal, ) @@ -305,32 +306,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: @@ -610,7 +585,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..a7c00d30d3bb 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,9 @@ _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") def _multi_modal_data(request: Dict[str, Any]) -> Dict[str, Any]: @@ -124,6 +127,64 @@ 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 + + flattened: list[str] = [] + for modality in _SGLANG_MM_ITEM_MODALITY_ORDER: + hashes = grouped.get(modality) + if hashes is None: + continue + 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 + 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 d8f2999e25af..9ab8e0254736 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py +++ b/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py @@ -505,6 +505,34 @@ 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": {}, + "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 2a0e0b92af11..1787b99fd6f0 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py +++ b/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py @@ -8,6 +8,7 @@ from dynamo.llm.exceptions import InvalidArgument from dynamo.sglang.request_handlers.llm.mm_disagg_utils import ( build_disagg_mm_kwargs, + extract_mm_hashes, extract_media_urls, raise_if_unextracted_multimodal, ) @@ -79,6 +80,41 @@ 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 = { + "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"] + + +@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..5f62b2c7c15e --- /dev/null +++ b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py @@ -0,0 +1,138 @@ +# 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_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", "LlavaForConditionalGeneration"), + ("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..79e3fd9f7b7c --- /dev/null +++ b/components/src/dynamo/sglang/video_routing.py @@ -0,0 +1,157 @@ +# 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 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 + +try: + from transformers.models.qwen3_vl.video_processing_qwen3_vl import ( + smart_resize as qwen3_smart_resize, + ) +except ImportError: + qwen3_smart_resize = None + +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) + 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, + # SGLang KV events do not attach MM metadata to delimiter/timestamp + # boundary blocks that contain no video placeholder run. + "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 0478fd01571c..80866d8867c6 100644 --- a/container/context.yaml +++ b/container/context.yaml @@ -127,8 +127,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 @@ -147,7 +147,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..c93444030b19 100644 --- a/container/templates/sglang_runtime.Dockerfile +++ b/container/templates/sglang_runtime.Dockerfile @@ -319,6 +319,10 @@ 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. +RUN python3 -c 'from dynamo.llm import MediaDecoder; assert hasattr(MediaDecoder(), "enable_video")' {% else %} ENV IMAGEIO_FFMPEG_EXE= {% endif %} diff --git a/container/templates/wheel_builder.Dockerfile b/container/templates/wheel_builder.Dockerfile index c81a67297659..5ce956eef95d 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,aic-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,aic-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/use-cases/multimodal-serving/additional-media-decoders.md b/docs/fern/pages/use-cases/multimodal-serving/additional-media-decoders.md index d0f3d370c981..5e1a2bef63c1 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,13 @@ 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. +- **VP8 and VP9 video** decode through the in-tree FFmpeg when frontend + decoding is enabled on a CUDA vLLM or SGLang runtime. - **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. -Installing an additional decoder package covers what remains: +Without frontend decoding, the backend worker owns VP8/VP9 decoding and needs +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. 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..06b932c520b3 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,13 +2,13 @@ # 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 @@ -19,7 +19,7 @@ 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 @@ -48,8 +48,8 @@ work, while the embedding cache can skip vision encoding for repeated images. 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 fa41b29c14f4..55f3cc04c8a4 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,19 @@ 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 in its CUDA runtime images: + +- 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. + +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 +131,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 CUDA vLLM +and SGLang 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 a CUDA vLLM or SGLang runtime. For +> other software-decoded cases, install a decode carrier alongside or transcode +> the input before sending it. ### Installing a software decoder 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/local_model/runtime_config.rs b/lib/llm/src/local_model/runtime_config.rs index 5bf98898ebeb..9ab4da734757 100644 --- a/lib/llm/src/local_model/runtime_config.rs +++ b/lib/llm/src/local_model/runtime_config.rs @@ -102,6 +102,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 6db5c6eaf2b2..ac8e36fdc74c 100644 --- a/lib/llm/src/model_card.rs +++ b/lib/llm/src/model_card.rs @@ -19,7 +19,8 @@ use std::sync::{Arc, OnceLock}; use crate::common::checked_file::CheckedFile; use crate::entrypoint::RouterConfig; use crate::local_model::runtime_config::{ - ModelRuntimeConfig, TokenizerBackend, VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ModelRuntimeConfig, SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, TokenizerBackend, + VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, }; use crate::model_type::{ModelInput, ModelType}; @@ -1266,6 +1267,7 @@ impl ModelDeploymentCard { // worker uses the same model-visible prompt expansion. for key in [ VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, ] { append_runtime_contract_checksum( @@ -3257,6 +3259,7 @@ mod ownership_tests { #[test] fn video_processor_runtime_contracts_isolate_worker_sets() { 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, }; @@ -3291,6 +3294,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}), @@ -3303,6 +3323,8 @@ mod ownership_tests { assert_eq!(missing.mdcsum(), unrelated.mdcsum()); assert_ne!(missing.mdcsum(), qwen.mdcsum()); + assert_ne!(missing.mdcsum(), sglang_qwen.mdcsum()); + assert_ne!(qwen.mdcsum(), sglang_qwen.mdcsum()); assert_eq!(qwen.mdcsum(), same_qwen.mdcsum()); assert_ne!(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 a9fd37f3302c..0e3e9d6e6d6d 100644 --- a/lib/llm/src/preprocessor.rs +++ b/lib/llm/src/preprocessor.rs @@ -53,12 +53,13 @@ use std::{ use tokio_util::sync::CancellationToken; use tracing; -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::{ + 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}; @@ -751,10 +752,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 { @@ -762,6 +762,7 @@ struct TrackedMmRoutingReplacement { target_tokens: Vec, worker_tokens: Vec, routing_tokens: Vec, + runless_boundary_uses_mm_metadata: bool, } /// Modality-aware routing payload accumulated in original message order. @@ -789,6 +790,7 @@ enum MmRoutingEntry { event_video_token_id: Option, target_tokens: Vec, replacement_tokens: Vec, + runless_boundary_uses_mm_metadata: bool, }, } @@ -1018,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 frontend +/// keeps the worker tokens. It adds MM metadata only when the worker's KV-event +/// contract does the same; SGLang's Qwen events are token-only in this case. #[cfg(feature = "mm-routing")] fn apply_tracked_mm_replacements( routing_prepend_bos: Option, @@ -1088,7 +1089,12 @@ fn apply_tracked_mm_replacements( let start = worker_tokens.len(); worker_tokens.extend_from_slice(&replacement.worker_tokens); routing_tokens.extend_from_slice(&replacement.routing_tokens); - spans.push((start, worker_tokens.len(), replacement.mm_hash)); + spans.push(( + start, + worker_tokens.len(), + replacement.mm_hash, + replacement.runless_boundary_uses_mm_metadata, + )); token_index += replacement.target_tokens.len(); replacement_index += 1; continue; @@ -1119,8 +1125,8 @@ fn apply_tracked_mm_replacements( let block_end = block_start + block_size; let mm_hashes: Vec = spans .iter() - .filter(|(start, end, _)| *start < block_end && *end > block_start) - .map(|(_, _, mm_hash)| *mm_hash) + .filter(|(start, end, _, _)| *start < block_end && *end > block_start) + .map(|(_, _, mm_hash, _)| *mm_hash) .collect(); if mm_hashes.is_empty() { continue; @@ -1150,15 +1156,24 @@ fn apply_tracked_mm_replacements( } None => { routing_block.copy_from_slice(worker_block); - block_mm_infos[block_index] = Some(BlockExtraInfo { - mm_objects: mm_hashes - .into_iter() - .map(|mm_hash| BlockMmObjectInfo { - mm_hash, - offsets: Vec::new(), - }) - .collect(), - }); + let metadata_hashes = spans + .iter() + .filter(|(start, end, _, uses_metadata)| { + *uses_metadata && *start < block_end && *end > block_start + }) + .map(|(_, _, mm_hash, _)| *mm_hash) + .collect::>(); + if !metadata_hashes.is_empty() { + block_mm_infos[block_index] = Some(BlockExtraInfo { + mm_objects: metadata_hashes + .into_iter() + .map(|mm_hash| BlockMmObjectInfo { + mm_hash, + offsets: Vec::new(), + }) + .collect(), + }); + } } } } @@ -2534,7 +2549,7 @@ impl OpenAIPreprocessor { #[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] let video_routing_processor = { - let qwen_contract = match runtime_config + let vllm_qwen_contract = match runtime_config .get_engine_specific::( VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, ) { @@ -2549,6 +2564,31 @@ impl OpenAIPreprocessor { None } }; + let sglang_qwen_contract = match runtime_config + .get_engine_specific::( + SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ) { + Ok(target) => target, + Err(error) => { + tracing::warn!( + target: "mm_routing", + %error, + key = SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + "invalid SGLang Qwen video processor runtime metadata; exact video routing disabled" + ); + None + } + }; + let qwen_contract = match (vllm_qwen_contract, sglang_qwen_contract) { + (Some(_), Some(_)) => { + tracing::warn!( + target: "mm_routing", + "multiple Qwen video processor contracts were published; exact video routing disabled" + ); + None + } + (contract, None) | (None, contract) => contract, + }; let nemotron_contract = match runtime_config .get_engine_specific::( VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, @@ -3551,6 +3591,8 @@ impl OpenAIPreprocessor { event_video_token_id: routing.event_video_token_id, target_tokens: routing.target_tokens, replacement_tokens: routing.replacement_tokens, + runless_boundary_uses_mm_metadata: routing + .runless_boundary_uses_mm_metadata, }) })(); match video_entry { @@ -4073,6 +4115,7 @@ impl OpenAIPreprocessor { target_tokens: vec![image_token_id], worker_tokens, routing_tokens, + runless_boundary_uses_mm_metadata: true, } } MmRoutingEntry::Video { @@ -4081,6 +4124,7 @@ impl OpenAIPreprocessor { event_video_token_id: _, target_tokens, replacement_tokens, + runless_boundary_uses_mm_metadata, } => { let fill_token = dynamo_kv_router::protocols::pad_value_for_mm_hash(*mm_hash); @@ -4098,6 +4142,7 @@ impl OpenAIPreprocessor { } }) .collect(), + runless_boundary_uses_mm_metadata: *runless_boundary_uses_mm_metadata, } } }; @@ -11951,6 +11996,7 @@ mod tests { 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], + runless_boundary_uses_mm_metadata: true, }; let (tokens, prompt_len, infos) = apply_tracked_mm_replacements( @@ -11970,6 +12016,38 @@ 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], + runless_boundary_uses_mm_metadata: false, + }; + + let (tokens, prompt_len, infos) = apply_tracked_mm_replacements( + None, + &[replacement], + &[1, 9, 2], + 4, + Some(99), + Some(video_token_id), + ) + .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_mixed_boundary_preserves_worker_hash_fallback() { @@ -11989,6 +12067,7 @@ mod tests { pad_value_for_mm_hash(image_hash), 7, ], + runless_boundary_uses_mm_metadata: true, }, TrackedMmRoutingReplacement { mm_hash: video_hash, @@ -12000,6 +12079,7 @@ mod tests { pad_value_for_mm_hash(video_hash), 9, ], + runless_boundary_uses_mm_metadata: true, }, ]; @@ -12052,6 +12132,7 @@ mod tests { pad_value_for_mm_hash(mm_hash), pad_value_for_mm_hash(mm_hash), ], + runless_boundary_uses_mm_metadata: true, }; let replacements = [ replacement(41, image_token_id), @@ -12096,6 +12177,7 @@ mod tests { target_tokens: vec![target], worker_tokens: vec![target], routing_tokens: vec![target], + runless_boundary_uses_mm_metadata: true, }; let replacements = [replacement(41, 10), replacement(42, 20)]; @@ -12127,6 +12209,7 @@ mod tests { target_tokens: vec![7], worker_tokens: vec![100, 19, 18, 18, 20, 101, 19, 18, 20], routing_tokens: vec![100, 19, pad, pad, 20, 101, 19, pad, 20], + runless_boundary_uses_mm_metadata: true, }; let (tokens, prompt_len, block_infos) = @@ -12166,6 +12249,7 @@ mod tests { event_video_token_id: Some(3), target_tokens: vec![3], replacement_tokens: vec![3], + runless_boundary_uses_mm_metadata: true, }; let image = MmRoutingEntry::Image { mm_hash: 2, diff --git a/lib/llm/src/preprocessor/mm_routing/mod.rs b/lib/llm/src/preprocessor/mm_routing/mod.rs index 941cdb242a12..cde1241dd17c 100644 --- a/lib/llm/src/preprocessor/mm_routing/mod.rs +++ b/lib/llm/src/preprocessor/mm_routing/mod.rs @@ -36,11 +36,39 @@ pub(crate) enum QwenVideoResizeMode { RoundTiesEven, } +/// 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, Default, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub(crate) enum QwenVideoRunlessBoundaryHash { + #[default] + MmMetadata, + TokensOnly, +} + /// Worker-reported Qwen video prompt-expansion behavior. -#[derive(Debug, Clone, Copy, Deserialize, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, Deserialize, PartialEq)] pub(crate) struct QwenVideoProcessorContract { pub placeholder_target: QwenVideoPlaceholderTarget, pub resize_mode: QwenVideoResizeMode, + #[serde(default)] + pub runless_boundary_hash: QwenVideoRunlessBoundaryHash, + #[serde(default)] + pub sglang_preprocess: Option, +} + +/// 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. @@ -73,6 +101,9 @@ pub(crate) struct VideoRoutingReplacement { /// Exact chat-template token sequence replaced by the model processor. pub target_tokens: Vec, pub replacement_tokens: Vec, + /// Whether runless boundary blocks carry the media hash separately from + /// their token sequence in the worker's KV event. + pub runless_boundary_uses_mm_metadata: bool, } enum SupportedVideoModel { diff --git a/lib/llm/src/preprocessor/mm_routing/nemotron.rs b/lib/llm/src/preprocessor/mm_routing/nemotron.rs index fbee0e1f55ad..4db70adf37e1 100644 --- a/lib/llm/src/preprocessor/mm_routing/nemotron.rs +++ b/lib/llm/src/preprocessor/mm_routing/nemotron.rs @@ -478,6 +478,7 @@ impl NemotronVideoRoutingSpec { event_video_token_id: None, target_tokens: self.video_target_tokens.clone(), replacement_tokens, + runless_boundary_uses_mm_metadata: true, }) } diff --git a/lib/llm/src/preprocessor/mm_routing/qwen3.rs b/lib/llm/src/preprocessor/mm_routing/qwen3.rs index 50d317185a21..853c042569b9 100644 --- a/lib/llm/src/preprocessor/mm_routing/qwen3.rs +++ b/lib/llm/src/preprocessor/mm_routing/qwen3.rs @@ -7,7 +7,8 @@ use anyhow::{Context, Result}; use serde_json::Value; use super::{ - QwenVideoPlaceholderTarget, QwenVideoProcessorContract, QwenVideoResizeMode, VideoRoutingInput, + QwenVideoPlaceholderTarget, QwenVideoProcessorContract, QwenVideoResizeMode, + QwenVideoRunlessBoundaryHash, SglangQwenVideoPreprocessContract, VideoRoutingInput, VideoRoutingReplacement, config::{read_json, read_model_config, required_token_id, required_usize}, }; @@ -40,9 +41,37 @@ 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 { + frame_count: usize, + width: u32, + height: u32, + source_fps: f64, + sampled_timestamps: Vec, +} + +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 from_model_dir( model_id: &str, @@ -128,6 +157,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 +168,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, + }; + 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 +193,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)) @@ -185,6 +224,105 @@ impl Qwen3VideoRoutingSpec { event_video_token_id: Some(self.video_token_id), target_tokens, replacement_tokens, + runless_boundary_uses_mm_metadata: matches!( + self.runless_boundary_hash, + QwenVideoRunlessBoundaryHash::MmMetadata + ), + }) + } + + fn prepare_input(&self, input: &VideoRoutingInput<'_>) -> 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: input.sampled_timestamps.to_vec(), + }); + }; + + contract.validate()?; + 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, + )?; + anyhow::ensure!( + max_frames >= min_frames, + "mm-routing: SGLang Qwen frontend supplied too few frames" + ); + let requested_frames = input.frame_count as f64 / effective_fps * contract.fps; + 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 { + index + .checked_mul(input.frame_count - 1) + .context("mm-routing: SGLang Qwen frame index overflow")? + / (frame_count - 1) + }; + Ok(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, }) } @@ -351,6 +489,64 @@ fn ensure_matching_value(config: &Value, field: &str, expected: usize) -> Result Ok(()) } +fn ceil_to_factor(value: usize, factor: usize) -> Result { + 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 +595,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 +613,51 @@ 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), } } + #[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 builds_timestamped_qwen_replacement_without_preprocessing_pixels() { let timestamps = [0.0, 1.0, 2.0, 3.0]; @@ -535,6 +774,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 +987,8 @@ mod tests { QwenVideoProcessorContract { placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, }, ) .unwrap(); @@ -797,6 +1040,8 @@ mod tests { QwenVideoProcessorContract { placeholder_target: QwenVideoPlaceholderTarget::BareVideoToken, resize_mode: QwenVideoResizeMode::LegacyCeil, + runless_boundary_hash: QwenVideoRunlessBoundaryHash::MmMetadata, + sglang_preprocess: None, }, ) .is_err() @@ -844,6 +1089,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 4b83914cd469..ddc9936cbb5a 100644 --- a/tests/serve/test_sglang.py +++ b/tests/serve/test_sglang.py @@ -57,6 +57,7 @@ router_selection_chat_payload_default, ) from tests.utils.payloads import ( + CachedTokensChatPayload, ChatPayload, HttpErrorPayload, ImageGenerationPayload, @@ -690,44 +691,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", + # Keep enough video tokens to make the routing sequence materially + # larger than the shared text prefix (the fixture decodes 10 frames). + "DYN_MM_VIDEO_NUM_FRAMES": "32", }, - 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, ) ], ), From 28b6eb649dc7186571861126908371374b6929c9 Mon Sep 17 00:00:00 2001 From: krishung5 Date: Thu, 17 Sep 2026 13:47:37 -0700 Subject: [PATCH 2/6] fix(mm-routing): address SGLang video routing CI failures Signed-off-by: krishung5 --- components/src/dynamo/sglang/register.py | 4 +- .../request_handlers/llm/decode_handler.py | 2 +- .../request_handlers/llm/mm_disagg_utils.py | 39 ++++++++++++++++++- .../tests/test_sglang_frontend_decoding.py | 5 ++- .../tests/test_sglang_multimodal_utils.py | 33 ++++++++++++---- .../sglang/tests/test_sglang_video_routing.py | 4 +- components/src/dynamo/sglang/video_routing.py | 6 ++- container/templates/sglang_runtime.Dockerfile | 2 + lib/llm/src/discovery/watcher.rs | 14 +++---- tests/serve/test_sglang.py | 6 +-- 10 files changed, 87 insertions(+), 28 deletions(-) diff --git a/components/src/dynamo/sglang/register.py b/components/src/dynamo/sglang/register.py index d74250c7d954..3dff723948d7 100644 --- a/components/src/dynamo/sglang/register.py +++ b/components/src/dynamo/sglang/register.py @@ -41,9 +41,7 @@ runtime_capacity, ) from dynamo.sglang.engine_generate import SGLANG_GENERATE_CAPABILITY -from dynamo.sglang.video_routing import ( - publish_sglang_qwen_video_processor_contract, -) +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" 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 84b41a09b9fc..8cecc9ca7cfd 100644 --- a/components/src/dynamo/sglang/request_handlers/llm/decode_handler.py +++ b/components/src/dynamo/sglang/request_handlers/llm/decode_handler.py @@ -40,8 +40,8 @@ IMAGE_URL_KEY, VIDEO_URL_KEY, build_disagg_mm_kwargs, - extract_mm_hashes, extract_media_urls, + extract_mm_hashes, raise_if_unextracted_multimodal, ) 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 a7c00d30d3bb..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 @@ -23,6 +23,11 @@ # 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]: @@ -156,11 +161,11 @@ def extract_mm_hashes(request: Dict[str, Any]) -> list[str] | None: ) return None - flattened: list[str] = [] + hashes_by_modality: dict[str, list[str]] = {} for modality in _SGLANG_MM_ITEM_MODALITY_ORDER: hashes = grouped.get(modality) if hashes is None: - continue + hashes = [] if not isinstance(hashes, list) or not all( isinstance(value, str) for value in hashes ): @@ -170,6 +175,36 @@ def extract_mm_hashes(request: Dict[str, Any]) -> list[str] | None: 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 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 9ab8e0254736..6b8290dc2519 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py +++ b/components/src/dynamo/sglang/tests/test_sglang_frontend_decoding.py @@ -518,7 +518,10 @@ async def fake_async_generate(**kwargs): handler.engine = SimpleNamespace(async_generate=fake_async_generate) request = { "token_ids": [1, 2, 3], - "multi_modal_data": {}, + "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"], 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 8030e25db45b..74d1bfbe59b8 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py +++ b/components/src/dynamo/sglang/tests/test_sglang_multimodal_utils.py @@ -8,8 +8,8 @@ from dynamo.llm.exceptions import InvalidArgument from dynamo.sglang.request_handlers.llm.mm_disagg_utils import ( build_disagg_mm_kwargs, - extract_mm_hashes, extract_media_urls, + extract_mm_hashes, raise_if_unextracted_multimodal, ) from dynamo.sglang.request_handlers.multimodal.worker_handler import StreamProcessor @@ -88,17 +88,39 @@ def test_extract_mm_hashes_preserves_legacy_image_protocol(): 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", [ @@ -109,10 +131,7 @@ def test_extract_mm_hashes_flattens_in_sglang_item_order(): ], ) def test_extract_mm_hashes_rejects_malformed_grouped_protocol(grouped): - assert ( - extract_mm_hashes({"extra_args": {"mm_hashes_by_modality": grouped}}) - is None - ) + assert extract_mm_hashes({"extra_args": {"mm_hashes_by_modality": grouped}}) is None class TestMultimodalGuard: diff --git a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py index 5f62b2c7c15e..c0b468e9a9de 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py +++ b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py @@ -45,9 +45,7 @@ class QwenProcessor(video_routing.transformers.ProcessorMixin): video_config=video_config or {}, _processor=QwenProcessor.__new__(QwenProcessor), ) - return SimpleNamespace( - tokenizer_manager=SimpleNamespace(mm_processor=mm_processor) - ) + return SimpleNamespace(tokenizer_manager=SimpleNamespace(mm_processor=mm_processor)) @pytest.fixture diff --git a/components/src/dynamo/sglang/video_routing.py b/components/src/dynamo/sglang/video_routing.py index 79e3fd9f7b7c..34740c367f78 100644 --- a/components/src/dynamo/sglang/video_routing.py +++ b/components/src/dynamo/sglang/video_routing.py @@ -5,6 +5,7 @@ import json import logging +from collections.abc import Callable from typing import Any, Optional import transformers @@ -16,12 +17,15 @@ 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, + smart_resize as _qwen3_smart_resize, ) except ImportError: qwen3_smart_resize = None +else: + qwen3_smart_resize = _qwen3_smart_resize logger = logging.getLogger(__name__) diff --git a/container/templates/sglang_runtime.Dockerfile b/container/templates/sglang_runtime.Dockerfile index c93444030b19..e33897aed0b5 100644 --- a/container/templates/sglang_runtime.Dockerfile +++ b/container/templates/sglang_runtime.Dockerfile @@ -322,7 +322,9 @@ RUN set -eu; \ # 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/lib/llm/src/discovery/watcher.rs b/lib/llm/src/discovery/watcher.rs index 1bdec6a7ed3c..2cb2a53a44d7 100644 --- a/lib/llm/src/discovery/watcher.rs +++ b/lib/llm/src/discovery/watcher.rs @@ -416,15 +416,15 @@ impl ModelWatcher { validate_policy_worker_role(card, &self.selection_policy)?; // Prepare without exact video routing unless the cohort agreed on a contract. - let removed_video_contract = spec.video_contract.is_none() - && [ + 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, - ] - .into_iter() - .fold(false, |removed, key| { - card.runtime_config.runtime_data.remove(key).is_some() || removed - }); + ] { + removed_video_contract |= card.runtime_config.runtime_data.remove(key).is_some(); + } + } if removed_video_contract { tracing::warn!( target: "mm_routing", diff --git a/tests/serve/test_sglang.py b/tests/serve/test_sglang.py index ddc9936cbb5a..5c07aa2191bd 100644 --- a/tests/serve/test_sglang.py +++ b/tests/serve/test_sglang.py @@ -710,9 +710,9 @@ class SGLangConfig(EngineConfig): ], env={ "DYN_MM_ALLOW_INTERNAL": "1", - # Keep enough video tokens to make the routing sequence materially - # larger than the shared text prefix (the fixture decodes 10 frames). - "DYN_MM_VIDEO_NUM_FRAMES": "32", + # 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=450, frontend_port=DefaultPort.FRONTEND.value, From 6dbaca4487c94295c9bd6a35913e3014fb7ab1ae Mon Sep 17 00:00:00 2001 From: krishung5 Date: Fri, 18 Sep 2026 10:51:35 -0700 Subject: [PATCH 3/6] fix(mm-routing): align SGLang video routing contracts Signed-off-by: krishung5 --- .../sglang/tests/test_sglang_video_routing.py | 1 + lib/llm/src/preprocessor.rs | 56 ++++++++++++ lib/llm/src/preprocessor/mm_routing/qwen3.rs | 89 +++++++++++++++++-- 3 files changed, 140 insertions(+), 6 deletions(-) diff --git a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py index c0b468e9a9de..beedc44e8b93 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py +++ b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py @@ -122,6 +122,7 @@ def test_inherited_video_replacement_publishes_wrapped_target(qwen_preprocessor) ("model_type", "architecture"), [ ("llava", "LlavaForConditionalGeneration"), + ("llava", "Qwen3VLForConditionalGeneration"), ("qwen3_vl", "Qwen3VLForCausalLM"), ], ) diff --git a/lib/llm/src/preprocessor.rs b/lib/llm/src/preprocessor.rs index dffdab90db47..411cd4c75e78 100644 --- a/lib/llm/src/preprocessor.rs +++ b/lib/llm/src/preprocessor.rs @@ -1121,6 +1121,17 @@ fn apply_tracked_mm_replacements( routing_tokens.resize(padded_len, 0); let mut block_mm_infos = vec![None; padded_len / block_size]; + // SGLang publishes token-only KV-event hashes for runless video boundary + // blocks. Its worker has already replaced every image/video placeholder + // with the canonical media pad before hashing, so keep the canonical + // request-side sequence intact and attach no per-block MM metadata. + if replacements + .iter() + .any(|replacement| !replacement.runless_boundary_uses_mm_metadata) + { + return Ok((routing_tokens, expanded_prompt_len, block_mm_infos)); + } + for (block_index, block_start) in (0..padded_len).step_by(block_size).enumerate() { let block_end = block_start + block_size; let mm_hashes: Vec = spans @@ -12240,6 +12251,51 @@ 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], + runless_boundary_uses_mm_metadata: true, + }, + 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], + runless_boundary_uses_mm_metadata: false, + }, + ]; + + 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), + ) + .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() { diff --git a/lib/llm/src/preprocessor/mm_routing/qwen3.rs b/lib/llm/src/preprocessor/mm_routing/qwen3.rs index 853c042569b9..b3dee543c37f 100644 --- a/lib/llm/src/preprocessor/mm_routing/qwen3.rs +++ b/lib/llm/src/preprocessor/mm_routing/qwen3.rs @@ -291,15 +291,15 @@ impl Qwen3VideoRoutingSpec { .map(|index| { let source_index = if frame_count == 1 { 0 + } else if index == frame_count - 1 { + input.frame_count - 1 } else { - index - .checked_mul(input.frame_count - 1) - .context("mm-routing: SGLang Qwen frame index overflow")? - / (frame_count - 1) + let step = (input.frame_count - 1) as f64 / (frame_count - 1) as f64; + (index as f64 * step).floor() as usize }; - Ok(source_index as f64 / effective_fps) + source_index as f64 / effective_fps }) - .collect::>>()?; + .collect(); let min_pixels = contract.video_min_pixels as f64; let max_pixels = (contract.video_max_pixels as f64) @@ -658,6 +658,83 @@ mod tests { 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]; From 4b6af2b241ff0a5e842e5cab7b8d5af9395386be Mon Sep 17 00:00:00 2001 From: krishung5 Date: Fri, 18 Sep 2026 13:30:18 -0700 Subject: [PATCH 4/6] fix(mm-routing): harden SGLang video contracts Signed-off-by: krishung5 --- .../sglang/tests/test_sglang_video_routing.py | 7 + components/src/dynamo/sglang/video_routing.py | 10 +- lib/llm/src/preprocessor.rs | 246 +++++++++++------- lib/llm/src/preprocessor/mm_routing/mod.rs | 183 ++++++++++++- .../src/preprocessor/mm_routing/nemotron.rs | 1 - lib/llm/src/preprocessor/mm_routing/qwen3.rs | 134 ++++++++-- 6 files changed, 461 insertions(+), 120 deletions(-) diff --git a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py index beedc44e8b93..a86c191754a3 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py +++ b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py @@ -109,6 +109,13 @@ def test_processor_override_disables_exact_video_contract(qwen_preprocessor): ) +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) diff --git a/components/src/dynamo/sglang/video_routing.py b/components/src/dynamo/sglang/video_routing.py index 34740c367f78..a0fcbebd0cac 100644 --- a/components/src/dynamo/sglang/video_routing.py +++ b/components/src/dynamo/sglang/video_routing.py @@ -120,6 +120,12 @@ def _resolve_qwen_video_processor_contract(engine: Any) -> Optional[dict[str, An 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 @@ -133,8 +139,8 @@ def _resolve_qwen_video_processor_contract(engine: Any) -> Optional[dict[str, An return { "placeholder_target": placeholder_target, "resize_mode": resize_mode, - # SGLang KV events do not attach MM metadata to delimiter/timestamp - # boundary blocks that contain no video placeholder run. + # 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), diff --git a/lib/llm/src/preprocessor.rs b/lib/llm/src/preprocessor.rs index 411cd4c75e78..5f8403bf22fb 100644 --- a/lib/llm/src/preprocessor.rs +++ b/lib/llm/src/preprocessor.rs @@ -55,7 +55,7 @@ use tracing; #[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] use crate::local_model::runtime_config::{ - SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, + ModelRuntimeConfig, SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, VLLM_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, }; @@ -762,7 +762,6 @@ struct TrackedMmRoutingReplacement { target_tokens: Vec, worker_tokens: Vec, routing_tokens: Vec, - runless_boundary_uses_mm_metadata: bool, } /// Modality-aware routing payload accumulated in original message order. @@ -790,7 +789,6 @@ enum MmRoutingEntry { event_video_token_id: Option, target_tokens: Vec, replacement_tokens: Vec, - runless_boundary_uses_mm_metadata: bool, }, } @@ -1020,9 +1018,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, the frontend -/// keeps the worker tokens. It adds MM metadata only when the worker's KV-event -/// contract does the same; SGLang's Qwen events are token-only in this case. +/// 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, @@ -1031,6 +1029,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, @@ -1089,12 +1088,7 @@ fn apply_tracked_mm_replacements( let start = worker_tokens.len(); worker_tokens.extend_from_slice(&replacement.worker_tokens); routing_tokens.extend_from_slice(&replacement.routing_tokens); - spans.push(( - start, - worker_tokens.len(), - replacement.mm_hash, - replacement.runless_boundary_uses_mm_metadata, - )); + spans.push((start, worker_tokens.len(), replacement.mm_hash)); token_index += replacement.target_tokens.len(); replacement_index += 1; continue; @@ -1121,23 +1115,12 @@ fn apply_tracked_mm_replacements( routing_tokens.resize(padded_len, 0); let mut block_mm_infos = vec![None; padded_len / block_size]; - // SGLang publishes token-only KV-event hashes for runless video boundary - // blocks. Its worker has already replaced every image/video placeholder - // with the canonical media pad before hashing, so keep the canonical - // request-side sequence intact and attach no per-block MM metadata. - if replacements - .iter() - .any(|replacement| !replacement.runless_boundary_uses_mm_metadata) - { - return Ok((routing_tokens, expanded_prompt_len, block_mm_infos)); - } - for (block_index, block_start) in (0..padded_len).step_by(block_size).enumerate() { let block_end = block_start + block_size; let mm_hashes: Vec = spans .iter() - .filter(|(start, end, _, _)| *start < block_end && *end > block_start) - .map(|(_, _, mm_hash, _)| *mm_hash) + .filter(|(start, end, _)| *start < block_end && *end > block_start) + .map(|(_, _, mm_hash)| *mm_hash) .collect(); if mm_hashes.is_empty() { continue; @@ -1165,26 +1148,22 @@ 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); - let metadata_hashes = spans - .iter() - .filter(|(start, end, _, uses_metadata)| { - *uses_metadata && *start < block_end && *end > block_start - }) - .map(|(_, _, mm_hash, _)| *mm_hash) - .collect::>(); - if !metadata_hashes.is_empty() { - block_mm_infos[block_index] = Some(BlockExtraInfo { - mm_objects: metadata_hashes - .into_iter() - .map(|mm_hash| BlockMmObjectInfo { - mm_hash, - offsets: Vec::new(), - }) - .collect(), - }); - } + block_mm_infos[block_index] = Some(BlockExtraInfo { + mm_objects: mm_hashes + .into_iter() + .map(|mm_hash| BlockMmObjectInfo { + mm_hash, + offsets: Vec::new(), + }) + .collect(), + }); } } } @@ -1713,6 +1692,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, @@ -2560,46 +2573,17 @@ impl OpenAIPreprocessor { #[cfg(all(feature = "mm-routing", feature = "media-ffmpeg"))] let video_routing_processor = { - let vllm_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 } }; - let sglang_qwen_contract = match runtime_config - .get_engine_specific::( - SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, - ) { - Ok(target) => target, - Err(error) => { - tracing::warn!( - target: "mm_routing", - %error, - key = SGLANG_QWEN_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, - "invalid SGLang Qwen video processor runtime metadata; exact video routing disabled" - ); - None - } - }; - let qwen_contract = match (vllm_qwen_contract, sglang_qwen_contract) { - (Some(_), Some(_)) => { - tracing::warn!( - target: "mm_routing", - "multiple Qwen video processor contracts were published; exact video routing disabled" - ); - None - } - (contract, None) | (None, contract) => contract, - }; let nemotron_contract = match runtime_config .get_engine_specific::( VLLM_NEMOTRON_VIDEO_PROCESSOR_CONTRACT_RUNTIME_KEY, @@ -3602,8 +3586,6 @@ impl OpenAIPreprocessor { event_video_token_id: routing.event_video_token_id, target_tokens: routing.target_tokens, replacement_tokens: routing.replacement_tokens, - runless_boundary_uses_mm_metadata: routing - .runless_boundary_uses_mm_metadata, }) })(); match video_entry { @@ -4126,7 +4108,6 @@ impl OpenAIPreprocessor { target_tokens: vec![image_token_id], worker_tokens, routing_tokens, - runless_boundary_uses_mm_metadata: true, } } MmRoutingEntry::Video { @@ -4135,7 +4116,6 @@ impl OpenAIPreprocessor { event_video_token_id: _, target_tokens, replacement_tokens, - runless_boundary_uses_mm_metadata, } => { let fill_token = dynamo_kv_router::protocols::pad_value_for_mm_hash(*mm_hash); @@ -4153,7 +4133,6 @@ impl OpenAIPreprocessor { } }) .collect(), - runless_boundary_uses_mm_metadata: *runless_boundary_uses_mm_metadata, } } }; @@ -4164,6 +4143,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, @@ -4171,6 +4159,7 @@ impl OpenAIPreprocessor { block_size, image_token_id, video_token_id, + kv_event_mm_identity, ) { Ok(expanded) => expanded, Err(error) => { @@ -12019,6 +12008,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() { @@ -12130,7 +12161,6 @@ mod tests { 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], - runless_boundary_uses_mm_metadata: true, }; let (tokens, prompt_len, infos) = apply_tracked_mm_replacements( @@ -12140,6 +12170,7 @@ mod tests { 4, Some(99), Some(video_token_id), + mm_routing::KvEventMmIdentity::MmMetadata, ) .unwrap(); @@ -12163,7 +12194,6 @@ mod tests { 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], - runless_boundary_uses_mm_metadata: false, }; let (tokens, prompt_len, infos) = apply_tracked_mm_replacements( @@ -12173,6 +12203,7 @@ mod tests { 4, Some(99), Some(video_token_id), + mm_routing::KvEventMmIdentity::PadValueTokens, ) .unwrap(); @@ -12182,6 +12213,35 @@ mod tests { 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() { @@ -12201,7 +12261,6 @@ mod tests { pad_value_for_mm_hash(image_hash), 7, ], - runless_boundary_uses_mm_metadata: true, }, TrackedMmRoutingReplacement { mm_hash: video_hash, @@ -12213,7 +12272,6 @@ mod tests { pad_value_for_mm_hash(video_hash), 9, ], - runless_boundary_uses_mm_metadata: true, }, ]; @@ -12224,6 +12282,7 @@ mod tests { 4, Some(image_token_id), Some(video_token_id), + mm_routing::KvEventMmIdentity::MmMetadata, ) .unwrap(); @@ -12268,14 +12327,12 @@ mod tests { target_tokens: vec![image_token_id], worker_tokens: vec![image_token_id; 14], routing_tokens: vec![image_pad; 14], - runless_boundary_uses_mm_metadata: true, }, 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], - runless_boundary_uses_mm_metadata: false, }, ]; @@ -12286,6 +12343,7 @@ mod tests { 16, Some(image_token_id), Some(video_token_id), + mm_routing::KvEventMmIdentity::PadValueTokens, ) .unwrap(); @@ -12311,7 +12369,6 @@ mod tests { pad_value_for_mm_hash(mm_hash), pad_value_for_mm_hash(mm_hash), ], - runless_boundary_uses_mm_metadata: true, }; let replacements = [ replacement(41, image_token_id), @@ -12326,6 +12383,7 @@ mod tests { 16, Some(image_token_id), Some(video_token_id), + mm_routing::KvEventMmIdentity::MmMetadata, ) .unwrap(); @@ -12356,7 +12414,6 @@ mod tests { target_tokens: vec![target], worker_tokens: vec![target], routing_tokens: vec![target], - runless_boundary_uses_mm_metadata: true, }; let replacements = [replacement(41, 10), replacement(42, 20)]; @@ -12369,6 +12426,7 @@ mod tests { 4, Some(10), Some(20), + mm_routing::KvEventMmIdentity::MmMetadata, ) .is_err(), "invalid target sequence {token_ids:?} must fail closed" @@ -12388,11 +12446,18 @@ mod tests { target_tokens: vec![7], worker_tokens: vec![100, 19, 18, 18, 20, 101, 19, 18, 20], routing_tokens: vec![100, 19, pad, pad, 20, 101, 19, pad, 20], - runless_boundary_uses_mm_metadata: true, }; - 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]); @@ -12428,7 +12493,6 @@ mod tests { event_video_token_id: Some(3), target_tokens: vec![3], replacement_tokens: vec![3], - runless_boundary_uses_mm_metadata: true, }; let image = MmRoutingEntry::Image { mm_hash: 2, diff --git a/lib/llm/src/preprocessor/mm_routing/mod.rs b/lib/llm/src/preprocessor/mm_routing/mod.rs index cde1241dd17c..10d2ad098cc1 100644 --- a/lib/llm/src/preprocessor/mm_routing/mod.rs +++ b/lib/llm/src/preprocessor/mm_routing/mod.rs @@ -39,25 +39,83 @@ pub(crate) enum QwenVideoResizeMode { /// 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, Default, Deserialize, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub(crate) enum QwenVideoRunlessBoundaryHash { - #[default] MmMetadata, TokensOnly, } -/// Worker-reported Qwen video prompt-expansion behavior. -#[derive(Debug, Clone, Copy, Deserialize, PartialEq)] +/// 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, - #[serde(default)] pub runless_boundary_hash: QwenVideoRunlessBoundaryHash, - #[serde(default)] 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 { @@ -83,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, @@ -101,9 +170,6 @@ pub(crate) struct VideoRoutingReplacement { /// Exact chat-template token sequence replaced by the model processor. pub target_tokens: Vec, pub replacement_tokens: Vec, - /// Whether runless boundary blocks carry the media hash separately from - /// their token sequence in the worker's KV event. - pub runless_boundary_uses_mm_metadata: bool, } enum SupportedVideoModel { @@ -170,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/nemotron.rs b/lib/llm/src/preprocessor/mm_routing/nemotron.rs index 4db70adf37e1..fbee0e1f55ad 100644 --- a/lib/llm/src/preprocessor/mm_routing/nemotron.rs +++ b/lib/llm/src/preprocessor/mm_routing/nemotron.rs @@ -478,7 +478,6 @@ impl NemotronVideoRoutingSpec { event_video_token_id: None, target_tokens: self.video_target_tokens.clone(), replacement_tokens, - runless_boundary_uses_mm_metadata: true, }) } diff --git a/lib/llm/src/preprocessor/mm_routing/qwen3.rs b/lib/llm/src/preprocessor/mm_routing/qwen3.rs index b3dee543c37f..d9d2873b31ad 100644 --- a/lib/llm/src/preprocessor/mm_routing/qwen3.rs +++ b/lib/llm/src/preprocessor/mm_routing/qwen3.rs @@ -1,13 +1,13 @@ // 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, + KvEventMmIdentity, QwenVideoPlaceholderTarget, QwenVideoProcessorContract, QwenVideoResizeMode, QwenVideoRunlessBoundaryHash, SglangQwenVideoPreprocessContract, VideoRoutingInput, VideoRoutingReplacement, config::{read_json, read_model_config, required_token_id, required_usize}, @@ -46,12 +46,12 @@ pub(super) struct Qwen3VideoRoutingSpec { tokenizer: Arc, } -struct PreparedVideoInput { +struct PreparedVideoInput<'a> { frame_count: usize, width: u32, height: u32, source_fps: f64, - sampled_timestamps: Vec, + sampled_timestamps: Cow<'a, [f64]>, } impl SglangQwenVideoPreprocessContract { @@ -73,6 +73,13 @@ impl SglangQwenVideoPreprocessContract { } 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, @@ -80,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( @@ -174,7 +184,7 @@ impl Qwen3VideoRoutingSpec { width: prepared.width, height: prepared.height, source_fps: prepared.source_fps, - sampled_timestamps: &prepared.sampled_timestamps, + sampled_timestamps: prepared.sampled_timestamps.as_ref(), }; let (grid_t, grid_h, grid_w) = self.video_grid(&prepared_input)?; let merge_area = self @@ -224,25 +234,20 @@ impl Qwen3VideoRoutingSpec { event_video_token_id: Some(self.video_token_id), target_tokens, replacement_tokens, - runless_boundary_uses_mm_metadata: matches!( - self.runless_boundary_hash, - QwenVideoRunlessBoundaryHash::MmMetadata - ), }) } - fn prepare_input(&self, input: &VideoRoutingInput<'_>) -> Result { + 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: input.sampled_timestamps.to_vec(), + sampled_timestamps: Cow::Borrowed(input.sampled_timestamps), }); }; - contract.validate()?; anyhow::ensure!( input.frame_count >= contract.frame_factor, "mm-routing: SGLang Qwen video frame count is below frame_factor" @@ -271,11 +276,10 @@ impl Qwen3VideoRoutingSpec { contract.max_frames.min(input.frame_count), contract.frame_factor, )?; - anyhow::ensure!( - max_frames >= min_frames, - "mm-routing: SGLang Qwen frontend supplied too few frames" - ); 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) @@ -322,7 +326,7 @@ impl Qwen3VideoRoutingSpec { height: u32::try_from(height) .context("mm-routing: resized video height exceeds u32")?, source_fps: effective_fps, - sampled_timestamps, + sampled_timestamps: Cow::Owned(sampled_timestamps), }) } @@ -490,6 +494,7 @@ fn ensure_matching_value(config: &Value, field: &str, expected: usize) -> Result } 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) @@ -619,6 +624,101 @@ mod tests { } } + 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(); From dbf8cc83baea788f5230a39b085671365404cf87 Mon Sep 17 00:00:00 2001 From: krishung5 Date: Fri, 18 Sep 2026 14:00:14 -0700 Subject: [PATCH 5/6] test(mm-routing): keep independent model gate cases Signed-off-by: krishung5 --- components/src/dynamo/sglang/tests/test_sglang_video_routing.py | 1 - 1 file changed, 1 deletion(-) diff --git a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py index a86c191754a3..4f01bf3738e1 100644 --- a/components/src/dynamo/sglang/tests/test_sglang_video_routing.py +++ b/components/src/dynamo/sglang/tests/test_sglang_video_routing.py @@ -128,7 +128,6 @@ def test_inherited_video_replacement_publishes_wrapped_target(qwen_preprocessor) @pytest.mark.parametrize( ("model_type", "architecture"), [ - ("llava", "LlavaForConditionalGeneration"), ("llava", "Qwen3VLForConditionalGeneration"), ("qwen3_vl", "Qwen3VLForCausalLM"), ], From 3e9b6672a3b54019797e408ecaa28410b8021a87 Mon Sep 17 00:00:00 2001 From: furionw Date: Tue, 22 Sep 2026 21:49:24 -0700 Subject: [PATCH 6/6] docs(sglang): document video-aware KV routing Signed-off-by: furionw --- .../backends/sglang/multimodal.md | 54 +++++++++++-------- .../backends/sglang/reference-guide.md | 1 + .../backends/sglang-configuration.mdx | 4 +- .../additional-media-decoders.md | 13 ++--- .../multimodal-kv-routing.md | 23 ++++---- .../parallel-media-decoding.md | 16 +++--- .../video-decode-gpu-requirements.md | 13 ++--- 7 files changed, 71 insertions(+), 53 deletions(-) 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 5e1a2bef63c1..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 @@ -10,12 +10,13 @@ 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 when frontend - decoding is enabled on a CUDA vLLM or SGLang runtime. -- **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. + 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. -Without frontend decoding, the backend worker owns VP8/VP9 decoding and needs -the corresponding package from the table below. Installing an additional -decoder package also covers: +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. @@ -101,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 06b932c520b3..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 @@ -11,7 +11,7 @@ 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 @@ -25,23 +25,27 @@ vision encoding. `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 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 569750e5eae9..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,13 +4,14 @@ title: Video Decode GPU Requirements --- -Dynamo provides two video-input decode paths in its CUDA runtime images: +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. + 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 @@ -131,8 +132,8 @@ between the two. ## Behavior when NVDEC is unavailable Hardware decode is additive and never blocks a request on its own: routing falls through -to a software decode path where one exists. VP8 and VP9 frontend decoding on CUDA vLLM -and SGLang does not depend on NVDEC. +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] > The shipped images do not include a software H.264, H.265, or AV1 decoder. @@ -142,7 +143,7 @@ and SGLang does not depend on NVDEC. > 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 VP8 and VP9, use frontend decoding on a CUDA vLLM or SGLang runtime. For +> 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. @@ -205,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. |