Skip to content

vllm.v1.engine.input_processor

Classes:

InputProcessor

Methods:

Source code in vllm/v1/engine/input_processor.py
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
class InputProcessor:
    def __init__(
        self,
        vllm_config: VllmConfig,
        renderer: BaseRenderer | None = None,
        *,
        mm_registry: MultiModalRegistry = MULTIMODAL_REGISTRY,
    ) -> None:
        self.vllm_config = vllm_config
        self.model_config = model_config = vllm_config.model_config
        self.cache_config = vllm_config.cache_config
        self.lora_config = vllm_config.lora_config
        self.scheduler_config = vllm_config.scheduler_config
        self.speculative_config = vllm_config.speculative_config
        self.structured_outputs_config = vllm_config.structured_outputs_config
        self.observability_config = vllm_config.observability_config
        self.diffusion_config = vllm_config.diffusion_config
        # Load the custom logits processor classes once; the returned callable
        # runs their validate_params hooks per request at admission.
        self.validate_logits_processors_params = (
            self._build_logits_processors_params_validator()
        )

        self.generation_config_fields = model_config.try_get_generation_config()

        self.renderer = renderer or renderer_from_config(vllm_config)

        self.supports_mm_inputs = model_config.supports_multimodal_inputs
        self.mm_encoder_cache_size = 0
        self.skip_prompt_length_check = False
        if self.supports_mm_inputs:
            mm_budget = MultiModalBudget(vllm_config, mm_registry)
            self.mm_encoder_cache_size = mm_budget.encoder_cache_size
            self.skip_prompt_length_check = (
                mm_budget.processor.info.skip_prompt_length_check
            )
            mm_budget.reset_cache()  # Not used anymore

        # Raw-prompt preprocessing (tokenization and multimodal processing)
        # is blocking, so async callers should run it on the renderer's
        # thread pool to keep their event loop responsive.
        self.process_inputs_async = make_async(
            self.process_inputs, executor=self.renderer._executor
        )

    @property
    def tokenizer(self) -> TokenizerLike | None:
        return self.renderer.tokenizer

    def get_tokenizer(self) -> TokenizerLike:
        return self.renderer.get_tokenizer()

    def _build_logits_processors_params_validator(
        self,
    ) -> Callable[[SamplingParams], None]:
        """Load the custom logits processor classes and return the per-request
        params validator for the active model runner."""
        custom_logitsprocs = self.model_config.logits_processors
        if self.vllm_config.use_v2_model_runner:
            if self.model_config.runner_type == "pooling" and not custom_logitsprocs:
                return lambda _: None

            from vllm.v1.worker.gpu.sample.logits_processor import (
                build_custom_logits_processors_params_validator,
            )

            return build_custom_logits_processors_params_validator(custom_logitsprocs)

        from vllm.v1.sample.logits_processor import (
            validate_logits_processors_parameters,
        )

        return partial(validate_logits_processors_parameters, custom_logitsprocs)

    def _validate_params(
        self,
        params: SamplingParams | PoolingParams,
        supported_tasks: tuple[SupportedTask, ...],
    ) -> None:
        """Raise `ValueError` if SamplingParams or PoolingParams is not valid."""
        if isinstance(params, SamplingParams):
            supported_generation_tasks = [
                task for task in supported_tasks if task in GENERATION_TASKS
            ]
            if not supported_generation_tasks:
                raise VLLMValidationError("This model does not support generation")

            params.verify(
                self.model_config,
                self.speculative_config,
                self.structured_outputs_config,
                self.tokenizer,
            )
            if params.prompt_logprob_token_ids is not None:
                if not self.vllm_config.use_v2_model_runner:
                    raise VLLMValidationError(
                        "prompt_logprob_token_ids requires the V2 model runner "
                        "(VLLM_USE_V2_MODEL_RUNNER=1).",
                        parameter="prompt_logprob_token_ids",
                    )
                if self.vllm_config.cache_config.kv_sharing_fast_prefill:
                    raise VLLMValidationError(
                        "prompt_logprob_token_ids is incorrect with "
                        "--kv-sharing-fast-prefill; disable it for scoring.",
                        parameter="prompt_logprob_token_ids",
                    )

            self.validate_logits_processors_params(params)

            if self.model_config.is_diffusion:
                # Without --diffusion-config the served canvas is unknown here;
                # the ids and the read-only normalisation are still checked.
                validate_diffusion_sampling_params(
                    params,
                    canvas_length=(
                        self.diffusion_config.canvas_length
                        if self.diffusion_config is not None
                        else None
                    ),
                    vocab_size=self.model_config.get_vocab_size(),
                )

            if self.model_config.return_sampling_mask:
                if params.temperature <= 0:
                    raise ValueError(
                        "sampling distribution replay requires temperature > 0"
                    )
                if params.top_k <= 0:
                    raise ValueError(
                        "sampling distribution replay requires top_k > 0 to "
                        "bound sampling mask size, reduce transfer overhead, "
                        "and avoid potential OOMs"
                    )
            if params.thinking_token_budget is not None and (
                self.vllm_config.reasoning_config is None
                or not self.vllm_config.reasoning_config.enabled
            ):
                raise VLLMValidationError(
                    "thinking_token_budget is set but reasoning_config is "
                    "not configured. Please set --reasoning-parser "
                    "and/or --reasoning-config to use thinking_token_budget."
                )
            if (
                params.trace_decode_token_ids
                and not self.model_config.enable_trace_replay
            ):
                raise VLLMValidationError(
                    "trace_decode_token_ids is set but trace replay is not "
                    "enabled. Start the engine with --enable-trace-replay "
                    "to use it."
                )
        elif isinstance(params, PoolingParams):
            supported_pooling_tasks = [
                task for task in supported_tasks if task in POOLING_TASKS
            ]
            if not supported_pooling_tasks:
                raise VLLMValidationError("This model does not support pooling")

            if params.task is None:
                if "token_embed" in supported_pooling_tasks:
                    params.task = "token_embed"
                elif "token_classify" in supported_pooling_tasks:
                    params.task = "token_classify"
                elif "plugin" in supported_pooling_tasks:
                    params.task = "plugin"

            if params.task not in supported_pooling_tasks:
                raise VLLMValidationError(
                    f"Unsupported task: {params.task!r} "
                    f"Supported tasks: {supported_pooling_tasks}"
                )

            params.verify(self.model_config)
        else:
            raise TypeError(
                f"params must be either SamplingParams or PoolingParams, "
                f"but got {type(params).__name__}"
            )

    def _normalize_trace_replay_params(
        self, sampling_params: SamplingParams, prompt_len: int
    ) -> None:
        """Apply trace replay's generation semantics to request-local params."""
        trace_token_ids = sampling_params.trace_decode_token_ids
        assert trace_token_ids
        assert sampling_params.max_tokens is not None

        max_trace_len = max(self.model_config.max_model_len - prompt_len, 1)
        trace_token_ids = trace_token_ids[:max_trace_len]
        sampling_params.trace_decode_token_ids = trace_token_ids

        # Apply this after the generation config so its EOS token cannot stop
        # replay before the trace is exhausted.
        sampling_params.max_tokens = min(
            len(trace_token_ids), sampling_params.max_tokens
        )
        sampling_params.min_tokens = 0
        sampling_params.ignore_eos = True
        sampling_params._eos_token_id = None
        sampling_params.stop = []
        sampling_params.stop_token_ids = []
        sampling_params._all_stop_token_ids = set()

    def _validate_lora(self, lora_request: LoRARequest | None) -> None:
        if lora_request is None:
            return

        # LoRA request passed in while LoRA is not enabled
        if not self.lora_config:
            raise VLLMValidationError(
                f"Got lora_request {lora_request} but LoRA is not enabled!"
            )

        if self.tokenizer is not None:
            logger.warning_once(
                "vLLM has deprecated support for supporting different "
                "tokenizers for different LoRAs. By default, vLLM uses base "
                "model's tokenizer. If you are using a LoRA "
                "with its own tokenizer, consider specifying `--tokenizer "
                "[lora_path]` to use the LoRA tokenizer."
            )

    def _get_mm_identifier(
        self,
        mm_hash: str,
        lora_request: LoRARequest | None,
    ) -> str:
        """When enable_tower_connector_lora is True, multi-modal embeddings
        vary depending on the LoRA request. Therefore, the mm_hash must be
        generated based on the LoRA request to prevent incorrect cache hits.
        """
        if (
            lora_request is None
            or self.lora_config is None
            or not self.lora_config.enable_tower_connector_lora
        ):
            return mm_hash
        return f"{lora_request.lora_name}:{mm_hash}"

    def inject_into_mm_cache(
        self,
        mm_hashes: dict[str, list[str]],
        mm_kwargs: dict[str, list],
    ) -> None:
        """Inject pre-processed mm_kwargs into the processor cache.

        Call this when mm_kwargs have already been through the HF processor
        externally (e.g. by a frontend that transfers pre-processed tensors
        to the backend).  This ensures MM cache hit rate metrics are reported
        accurately and avoids redundant processing on subsequent requests
        with the same images.

        Uses ``get_and_update_item()`` with an empty prompt_updates list,
        since token expansion has already been handled externally.
        """
        cache = self.renderer.mm_processor_cache
        if cache is None:
            return
        try:
            for modality, hashes in mm_hashes.items():
                items = mm_kwargs.get(modality, [])
                for i, mm_hash in enumerate(hashes):
                    if i < len(items) and items[i] is not None:
                        # Insert into cache via get_and_update_item.
                        # Use the returned item (may be an address for SHM
                        # cache or the original item for LRU cache).
                        items[i], _ = cache.get_and_update_item(
                            (items[i], []),
                            mm_hash,
                        )
            # Update cache stats to reflect the externally processed items
            self.renderer.update_mm_cache_stats()
        except Exception:
            logger.warning(
                "Failed to inject mm_kwargs into processor cache",
                exc_info=True,
            )

    @staticmethod
    def assign_request_id(request: EngineCoreRequest):
        """Replace the externally supplied request ID with an internal request ID
        that adds 8 random characters in order to ensure uniqueness.
        """
        if request.external_req_id is not None:
            raise ValueError(
                "The external_req_id field should not be set on EngineCoreRequests"
                " passed to vLLM; use the request_id field."
            )
        request.external_req_id = request.request_id
        if envs.VLLM_DISABLE_REQUEST_ID_RANDOMIZATION:
            logger.warning_once(
                "VLLM_DISABLE_REQUEST_ID_RANDOMIZATION is set and will be "
                "removed in a future release. Duplicate externally-provided "
                "request IDs may cause failures and/or subtle correctness errors."
            )
        else:
            request.request_id = f"{request.external_req_id}-{random_uuid():.8}"

    def process_inputs(
        self,
        request_id: str,
        prompt: PromptType | EngineInput,
        params: SamplingParams | PoolingParams,
        supported_tasks: tuple[SupportedTask, ...],
        arrival_time: float | None = None,
        lora_request: LoRARequest | None = None,
        tokenization_kwargs: dict[str, Any] | None = None,
        trace_headers: Mapping[str, str] | None = None,
        priority: int = 0,
        data_parallel_rank: int | None = None,
        resumable: bool = False,
        session_id: str | None = None,
        kv_hints: KvHintsEnvelope | None = None,
    ) -> EngineCoreRequest:
        self._validate_params(params, supported_tasks)
        self._validate_lora(lora_request)

        parallel_config = self.vllm_config.parallel_config
        dp_size = parallel_config.data_parallel_size
        dp_local_size = parallel_config.data_parallel_size_local
        num_ranks = dp_local_size if parallel_config.local_engines_only else dp_size
        if data_parallel_rank is not None and not (0 <= data_parallel_rank < num_ranks):
            raise VLLMValidationError(
                f"data_parallel_rank {data_parallel_rank} "
                f"is out of range [0, {num_ranks})."
            )

        if isinstance(prompt, dict) and "type" in prompt:
            if arrival_time is None:
                arrival_time = prompt.get("arrival_time", time.time())  # type: ignore[assignment]

            engine_input: EngineInput = prompt  # type: ignore[assignment]
        else:
            logger.warning_once(
                "Passing raw prompts to InputProcessor is deprecated "
                "and will be removed in the future. You should instead pass "
                "the outputs of Renderer.render_cmpl() or Renderer.render_chat()."
            )

            if arrival_time is None:
                arrival_time = time.time()

            renderer = self.renderer
            model_config = self.model_config

            parsed_prompt = parse_model_prompt(model_config, prompt)
            tok_params = renderer.default_cmpl_tok_params.with_kwargs(
                **(tokenization_kwargs or {})
            )

            (engine_input,) = renderer.render_cmpl(
                [parsed_prompt],
                tok_params,
            )

        current_platform.validate_request(engine_input, params)

        encoder_input, decoder_input = split_enc_dec_input(engine_input)
        self._validate_model_inputs(encoder_input, decoder_input)

        # Mypy can be conservative for TypedDict unions; normalize access.
        if decoder_input["type"] == "embeds":
            prompt_embeds = decoder_input["prompt_embeds"]
            prompt_token_ids = decoder_input.get("prompt_token_ids")
            prompt_is_token_ids = decoder_input.get("is_token_ids")
        else:
            prompt_token_ids = decoder_input["prompt_token_ids"]
            prompt_embeds = None
            prompt_is_token_ids = None

        sampling_params = None
        pooling_params = None
        if isinstance(params, SamplingParams):
            # TODO: can we avoid cloning here in multiproc case?
            sampling_params = params.clone()
            prompt_len = length_from_prompt_token_ids_or_embeds(
                prompt_token_ids, prompt_embeds
            )
            if not 0 <= sampling_params.routed_experts_prompt_start <= prompt_len:
                raise VLLMValidationError(
                    f"routed_experts_prompt_start must be between 0 and "
                    f"the prompt length ({prompt_len}), inclusive.",
                    parameter="routed_experts_prompt_start",
                    value=sampling_params.routed_experts_prompt_start,
                )
            # If unset max tokens, then generate up to the max_model_len.
            if sampling_params.max_tokens is None:
                sampling_params.max_tokens = (
                    self.model_config.max_model_len - prompt_len
                )
                # min_tokens is not checked while max_tokens is unset.
                if sampling_params.min_tokens > sampling_params.max_tokens:
                    raise VLLMValidationError(
                        f"min_tokens must be less than or equal to "
                        f"max_tokens={sampling_params.max_tokens}, got "
                        f"{sampling_params.min_tokens}."
                    )

            sampling_params.update_from_generation_config(
                self.generation_config_fields,
                self.renderer.get_eos_token_id(),
            )
            if self.tokenizer is not None:
                sampling_params.update_from_tokenizer(self.tokenizer)
            if sampling_params.trace_decode_token_ids:
                self._normalize_trace_replay_params(sampling_params, prompt_len)
        else:
            pooling_params = params.clone()

        # Multimodal related.
        mm_features: list[MultiModalFeatureSpec] | None = None

        if decoder_input["type"] == "multimodal":
            decoder_mm_inputs = decoder_input["mm_kwargs"]
            decoder_mm_positions = decoder_input["mm_placeholders"]
            decoder_mm_hashes = decoder_input["mm_hashes"]

            if not all(
                isinstance(leaf, str) for leaf in json_iter_leaves(decoder_mm_hashes)
            ):
                raise ValueError(
                    f"mm_hashes must contain only strings, got: {decoder_mm_hashes}. "
                    "This is likely due to an incorrect custom implementation of "
                    "MultiModalProcessor.apply method."
                )

            # Merge and flatten multimodal placeholders, hashes and inputs
            # from dictionaries to lists, and sort them by each item's position
            # in the input sequence.
            sorted_mm_idxs = argsort_mm_positions(decoder_mm_positions)

            mm_features = []
            for modality, idx in sorted_mm_idxs:
                base_mm_hash = decoder_mm_hashes[modality][idx]
                mm_features.append(
                    MultiModalFeatureSpec(
                        data=decoder_mm_inputs[modality][idx],
                        modality=modality,
                        identifier=self._get_mm_identifier(
                            base_mm_hash,
                            lora_request,
                        ),
                        mm_position=decoder_mm_positions[modality][idx],
                        mm_hash=base_mm_hash,
                    )
                )

        return EngineCoreRequest(
            request_id=request_id,
            prompt_token_ids=prompt_token_ids,
            prompt_embeds=prompt_embeds,
            prompt_is_token_ids=prompt_is_token_ids,
            mm_features=mm_features,
            sampling_params=sampling_params,
            pooling_params=pooling_params,
            arrival_time=arrival_time,
            lora_request=lora_request,
            cache_salt=decoder_input.get("cache_salt"),
            priority=priority,
            data_parallel_rank=data_parallel_rank,
            trace_headers=trace_headers,
            resumable=resumable,
            session_id=session_id,
            kv_hints=kv_hints,
        )

    def _validate_prompt_len(
        self,
        prompt_len: int,
        prompt_type: Literal["encoder", "decoder"],
    ):
        if self.skip_prompt_length_check and prompt_type == "encoder":
            return

        if prompt_len == 0 and prompt_type == "decoder":
            raise VLLMValidationError(f"The {prompt_type} prompt cannot be empty")

        model_config = self.model_config
        max_prompt_len = (
            model_config.max_model_len
            if prompt_type == "decoder"
            else self.mm_encoder_cache_size
        )
        if prompt_len > max_prompt_len:
            if self.supports_mm_inputs:
                suggestion = (
                    "Make sure that `max_model_len` is no smaller than the "
                    "number of text tokens plus multimodal tokens. For image "
                    "inputs, the number of image tokens depends on the number "
                    "of images, and possibly their aspect ratios as well."
                )
            else:
                suggestion = (
                    "Make sure that `max_model_len` is no smaller than the "
                    "number of text tokens."
                )

            raise VLLMValidationError(
                f"The {prompt_type} prompt (length {prompt_len}) is "
                f"longer than the maximum model length of {max_prompt_len}. "
                f"{suggestion}"
            )
        elif prompt_len == max_prompt_len and model_config.runner_type == "generate":
            suggestion = (
                "Make sure that `max_model_len` is no smaller than the "
                "number of text tokens (prompt + requested output tokens)."
            )
            raise VLLMValidationError(
                f"The {prompt_type} prompt (length {prompt_len}) plus the number of "
                f"requested output tokens (at least 1) is longer than the maximum "
                f"model length of {max_prompt_len}. {suggestion}"
            )

    def _validate_model_input(
        self,
        prompt_input: SingletonInput,
        prompt_type: Literal["encoder", "decoder"],
    ) -> None:
        prompt_ids = (
            None
            if prompt_input["type"] == "embeds"
            else prompt_input["prompt_token_ids"]
        )
        prompt_embeds = (
            prompt_input["prompt_embeds"] if prompt_input["type"] == "embeds" else None
        )

        prompt_len = length_from_prompt_token_ids_or_embeds(prompt_ids, prompt_embeds)
        self._validate_prompt_len(prompt_len, prompt_type)

        if prompt_input["type"] == "embeds":
            is_token_ids = prompt_input.get("is_token_ids")
            if is_token_ids is not None and len(is_token_ids) != prompt_len:
                raise VLLMValidationError(
                    "prompt_is_token_ids must have the same length as prompt_embeds "
                    f"(expected {prompt_len}, got {len(is_token_ids)}).",
                    parameter="prompt_is_token_ids",
                )

        if prompt_input["type"] == "multimodal":
            decoder_mm_positions = prompt_input["mm_placeholders"]
            for modality, mm_positions in decoder_mm_positions.items():
                for mm_position in mm_positions:
                    num_embeds = mm_position.get_num_embeds()
                    if num_embeds > self.mm_encoder_cache_size:
                        raise VLLMValidationError(
                            f"The {prompt_type} prompt contains a(n) {modality} item "
                            f"with {num_embeds} embedding tokens, which exceeds the "
                            f"pre-allocated encoder cache size "
                            f"{self.mm_encoder_cache_size}. Please reduce the input "
                            f"size or increase the encoder cache size "
                            f"by setting --limit-mm-per-prompt at startup."
                        )

        # Shared by generate, embedding and pooling requests.
        if prompt_ids:
            self.renderer.validate_token_ids(prompt_ids)

    def _validate_model_inputs(
        self,
        encoder_input: SingletonInput | None,
        decoder_input: SingletonInput,
    ):
        if encoder_input is not None:
            self._validate_model_input(encoder_input, prompt_type="encoder")

        self._validate_model_input(decoder_input, prompt_type="decoder")

_build_logits_processors_params_validator()

Load the custom logits processor classes and return the per-request params validator for the active model runner.

Source code in vllm/v1/engine/input_processor.py
def _build_logits_processors_params_validator(
    self,
) -> Callable[[SamplingParams], None]:
    """Load the custom logits processor classes and return the per-request
    params validator for the active model runner."""
    custom_logitsprocs = self.model_config.logits_processors
    if self.vllm_config.use_v2_model_runner:
        if self.model_config.runner_type == "pooling" and not custom_logitsprocs:
            return lambda _: None

        from vllm.v1.worker.gpu.sample.logits_processor import (
            build_custom_logits_processors_params_validator,
        )

        return build_custom_logits_processors_params_validator(custom_logitsprocs)

    from vllm.v1.sample.logits_processor import (
        validate_logits_processors_parameters,
    )

    return partial(validate_logits_processors_parameters, custom_logitsprocs)

_get_mm_identifier(mm_hash, lora_request)

When enable_tower_connector_lora is True, multi-modal embeddings vary depending on the LoRA request. Therefore, the mm_hash must be generated based on the LoRA request to prevent incorrect cache hits.

Source code in vllm/v1/engine/input_processor.py
def _get_mm_identifier(
    self,
    mm_hash: str,
    lora_request: LoRARequest | None,
) -> str:
    """When enable_tower_connector_lora is True, multi-modal embeddings
    vary depending on the LoRA request. Therefore, the mm_hash must be
    generated based on the LoRA request to prevent incorrect cache hits.
    """
    if (
        lora_request is None
        or self.lora_config is None
        or not self.lora_config.enable_tower_connector_lora
    ):
        return mm_hash
    return f"{lora_request.lora_name}:{mm_hash}"

_normalize_trace_replay_params(sampling_params, prompt_len)

Apply trace replay's generation semantics to request-local params.

Source code in vllm/v1/engine/input_processor.py
def _normalize_trace_replay_params(
    self, sampling_params: SamplingParams, prompt_len: int
) -> None:
    """Apply trace replay's generation semantics to request-local params."""
    trace_token_ids = sampling_params.trace_decode_token_ids
    assert trace_token_ids
    assert sampling_params.max_tokens is not None

    max_trace_len = max(self.model_config.max_model_len - prompt_len, 1)
    trace_token_ids = trace_token_ids[:max_trace_len]
    sampling_params.trace_decode_token_ids = trace_token_ids

    # Apply this after the generation config so its EOS token cannot stop
    # replay before the trace is exhausted.
    sampling_params.max_tokens = min(
        len(trace_token_ids), sampling_params.max_tokens
    )
    sampling_params.min_tokens = 0
    sampling_params.ignore_eos = True
    sampling_params._eos_token_id = None
    sampling_params.stop = []
    sampling_params.stop_token_ids = []
    sampling_params._all_stop_token_ids = set()

_validate_params(params, supported_tasks)

Raise ValueError if SamplingParams or PoolingParams is not valid.

Source code in vllm/v1/engine/input_processor.py
def _validate_params(
    self,
    params: SamplingParams | PoolingParams,
    supported_tasks: tuple[SupportedTask, ...],
) -> None:
    """Raise `ValueError` if SamplingParams or PoolingParams is not valid."""
    if isinstance(params, SamplingParams):
        supported_generation_tasks = [
            task for task in supported_tasks if task in GENERATION_TASKS
        ]
        if not supported_generation_tasks:
            raise VLLMValidationError("This model does not support generation")

        params.verify(
            self.model_config,
            self.speculative_config,
            self.structured_outputs_config,
            self.tokenizer,
        )
        if params.prompt_logprob_token_ids is not None:
            if not self.vllm_config.use_v2_model_runner:
                raise VLLMValidationError(
                    "prompt_logprob_token_ids requires the V2 model runner "
                    "(VLLM_USE_V2_MODEL_RUNNER=1).",
                    parameter="prompt_logprob_token_ids",
                )
            if self.vllm_config.cache_config.kv_sharing_fast_prefill:
                raise VLLMValidationError(
                    "prompt_logprob_token_ids is incorrect with "
                    "--kv-sharing-fast-prefill; disable it for scoring.",
                    parameter="prompt_logprob_token_ids",
                )

        self.validate_logits_processors_params(params)

        if self.model_config.is_diffusion:
            # Without --diffusion-config the served canvas is unknown here;
            # the ids and the read-only normalisation are still checked.
            validate_diffusion_sampling_params(
                params,
                canvas_length=(
                    self.diffusion_config.canvas_length
                    if self.diffusion_config is not None
                    else None
                ),
                vocab_size=self.model_config.get_vocab_size(),
            )

        if self.model_config.return_sampling_mask:
            if params.temperature <= 0:
                raise ValueError(
                    "sampling distribution replay requires temperature > 0"
                )
            if params.top_k <= 0:
                raise ValueError(
                    "sampling distribution replay requires top_k > 0 to "
                    "bound sampling mask size, reduce transfer overhead, "
                    "and avoid potential OOMs"
                )
        if params.thinking_token_budget is not None and (
            self.vllm_config.reasoning_config is None
            or not self.vllm_config.reasoning_config.enabled
        ):
            raise VLLMValidationError(
                "thinking_token_budget is set but reasoning_config is "
                "not configured. Please set --reasoning-parser "
                "and/or --reasoning-config to use thinking_token_budget."
            )
        if (
            params.trace_decode_token_ids
            and not self.model_config.enable_trace_replay
        ):
            raise VLLMValidationError(
                "trace_decode_token_ids is set but trace replay is not "
                "enabled. Start the engine with --enable-trace-replay "
                "to use it."
            )
    elif isinstance(params, PoolingParams):
        supported_pooling_tasks = [
            task for task in supported_tasks if task in POOLING_TASKS
        ]
        if not supported_pooling_tasks:
            raise VLLMValidationError("This model does not support pooling")

        if params.task is None:
            if "token_embed" in supported_pooling_tasks:
                params.task = "token_embed"
            elif "token_classify" in supported_pooling_tasks:
                params.task = "token_classify"
            elif "plugin" in supported_pooling_tasks:
                params.task = "plugin"

        if params.task not in supported_pooling_tasks:
            raise VLLMValidationError(
                f"Unsupported task: {params.task!r} "
                f"Supported tasks: {supported_pooling_tasks}"
            )

        params.verify(self.model_config)
    else:
        raise TypeError(
            f"params must be either SamplingParams or PoolingParams, "
            f"but got {type(params).__name__}"
        )

assign_request_id(request) staticmethod

Replace the externally supplied request ID with an internal request ID that adds 8 random characters in order to ensure uniqueness.

Source code in vllm/v1/engine/input_processor.py
@staticmethod
def assign_request_id(request: EngineCoreRequest):
    """Replace the externally supplied request ID with an internal request ID
    that adds 8 random characters in order to ensure uniqueness.
    """
    if request.external_req_id is not None:
        raise ValueError(
            "The external_req_id field should not be set on EngineCoreRequests"
            " passed to vLLM; use the request_id field."
        )
    request.external_req_id = request.request_id
    if envs.VLLM_DISABLE_REQUEST_ID_RANDOMIZATION:
        logger.warning_once(
            "VLLM_DISABLE_REQUEST_ID_RANDOMIZATION is set and will be "
            "removed in a future release. Duplicate externally-provided "
            "request IDs may cause failures and/or subtle correctness errors."
        )
    else:
        request.request_id = f"{request.external_req_id}-{random_uuid():.8}"

inject_into_mm_cache(mm_hashes, mm_kwargs)

Inject pre-processed mm_kwargs into the processor cache.

Call this when mm_kwargs have already been through the HF processor externally (e.g. by a frontend that transfers pre-processed tensors to the backend). This ensures MM cache hit rate metrics are reported accurately and avoids redundant processing on subsequent requests with the same images.

Uses get_and_update_item() with an empty prompt_updates list, since token expansion has already been handled externally.

Source code in vllm/v1/engine/input_processor.py
def inject_into_mm_cache(
    self,
    mm_hashes: dict[str, list[str]],
    mm_kwargs: dict[str, list],
) -> None:
    """Inject pre-processed mm_kwargs into the processor cache.

    Call this when mm_kwargs have already been through the HF processor
    externally (e.g. by a frontend that transfers pre-processed tensors
    to the backend).  This ensures MM cache hit rate metrics are reported
    accurately and avoids redundant processing on subsequent requests
    with the same images.

    Uses ``get_and_update_item()`` with an empty prompt_updates list,
    since token expansion has already been handled externally.
    """
    cache = self.renderer.mm_processor_cache
    if cache is None:
        return
    try:
        for modality, hashes in mm_hashes.items():
            items = mm_kwargs.get(modality, [])
            for i, mm_hash in enumerate(hashes):
                if i < len(items) and items[i] is not None:
                    # Insert into cache via get_and_update_item.
                    # Use the returned item (may be an address for SHM
                    # cache or the original item for LRU cache).
                    items[i], _ = cache.get_and_update_item(
                        (items[i], []),
                        mm_hash,
                    )
        # Update cache stats to reflect the externally processed items
        self.renderer.update_mm_cache_stats()
    except Exception:
        logger.warning(
            "Failed to inject mm_kwargs into processor cache",
            exc_info=True,
        )