vllm_omni.engine.output_processor ¶
MultimodalCompletionOutput dataclass ¶
Bases: CompletionOutput
CompletionOutput with multimodal support.
Inherits all CompletionOutput fields and adds multimodal_output. As a CompletionOutput subclass, compatible with all existing vLLM consumers.
MultimodalOutputProcessor ¶
Bases: OutputProcessor
Handles multimodal output processing.
Captures multimodal outputs from OmniEngineCoreOutput and accumulates them as MultimodalPayload in OmniRequestState, before delegating to the base vLLM OutputProcessor for text handling.
The data flow is: 1. For each EngineCoreOutput with multimodal_output: - Capture into OmniRequestState.add_multimodal_tensor() 2. Base vLLM OutputProcessor handles text detokenization 3. On finish, _consolidate_multimodal_tensors() concatenates accumulated tensors using strategy-based dispatch 4. _new_completion_output() returns MultimodalCompletionOutput
output_modality instance-attribute ¶
output_modality = OutputModality.from_string(
engine_core_output_type
)
abort_requests_collecting_outputs ¶
abort_requests_collecting_outputs(
request_ids: Iterable[str],
*,
internal: bool,
commit_state: bool = True,
) -> tuple[
list[str], list[RequestOutput | PoolingRequestOutput]
]
Abort requests and return terminal abort outputs with partial tokens.
Mirrors upstream OutputProcessor.abort_requests, but also returns abort RequestOutput objects when queue is None (Omni stage pools register OP state without a collector queue).
Terminal abort outputs use RequestOutputKind.CUMULATIVE so DELTA streaming requests still surface the full prefix generated so far. new_token_ids is the detokenizer prefix: Omni's make_request_output writes that list onto CompletionOutput.token_ids.
When commit_state is False, processor mappings stay in place so a failed physical EngineCore abort can retry with the same prefix.
add_request ¶
add_request(
request: EngineCoreRequest,
prompt: str | None,
parent_req: ParentRequest | None = None,
request_index: int = 0,
queue: RequestOutputCollector | None = None,
) -> None
Add a new request to be processed.
Creates an OmniRequestState for the request and registers it for output processing.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
request | EngineCoreRequest | Engine core request to add | required |
prompt | str | None | Optional prompt string for the request | required |
parent_req | ParentRequest | None | Optional parent request for parallel sampling | None |
request_index | int | Index of the request in the batch | 0 |
queue | RequestOutputCollector | None | Optional queue for collecting outputs | None |
Raises:
| Type | Description |
|---|---|
ValueError | If the request ID is already registered |
commit_aborted_request_state ¶
Drop processor mappings after a successful physical EngineCore abort.
MultimodalPayload dataclass ¶
Bases: Mapping
Structured multimodal output payload.
Implements collections.abc.Mapping so that isinstance(payload, dict) style checks in downstream code can be replaced with duck-typing, and payload.get(key), payload[key], key in payload, len(payload) all work seamlessly for both tensors and metadata.
Attributes:
| Name | Type | Description |
|---|---|---|
tensors | dict[str, Tensor] | Dictionary mapping modality/key names to their tensors. |
metadata | dict[str, Any] | Optional dictionary for non-tensor metadata (e.g., sample rate for audio, image dimensions). |
metadata class-attribute instance-attribute ¶
primary_tensor property ¶
Return the first tensor in the payload, or None if empty.
tensors class-attribute instance-attribute ¶
consolidate_metadata ¶
Resolve deferred tensor lists in metadata by keeping the latest value.
Metadata values are per-step snapshots (e.g. sample rate), not content deltas, so the latest value supersedes earlier ones. Nested dicts (unflattened payloads) are resolved one level down.
consolidate_tensors ¶
consolidate_tensors(modality: OutputModality) -> None
Concatenate deferred tensor lists into single tensors.
Tensors are generated content accumulated as chunks, so each key's list is concatenated according to the strategy get_accumulation_strategy(modality, key) resolves for it (e.g. audio waveform chunks along the time dimension, latent frames along the batch dimension). Most keys share modality's default, but a key can be registered for a different strategy -- e.g. a codec-frame matrix that grows along dim 0 rather than the waveform-tuned default for its modality -- via register_key_accumulation_strategy.
from_dict classmethod ¶
from_dict(
data: dict[str, Any] | None,
) -> MultimodalPayload | None
Create a MultimodalPayload from a raw dictionary.
Separates torch.Tensor values into tensors and everything else into metadata.
from_raw classmethod ¶
from_raw(
payload: Any, modality_key: str
) -> MultimodalPayload | None
Create a MultimodalPayload from a raw producer payload.
Accepts a MultimodalPayload (returned as-is), a dict, or a bare tensor (stored under modality_key). Tensors are moved to CPU. Producer-specific dict keys are remapped to the semantic modality key (e.g. "audio", "latent"): AR runners produce {"hidden": ...} and generation runners produce {"model_outputs": ...}.
get ¶
Get a value by key, searching tensors first then metadata.
merged_with ¶
merged_with(
incoming: MultimodalPayload,
) -> MultimodalPayload
Merge incoming onto this payload and return the result.
Content tensors accumulate into lists for deferred concatenation; known sample-rate keys are snapshots replaced immediately, including in DELTA streams that do not consolidate each emission. Missing keys retain their previous value. Other values keep the existing merge behavior. When this payload is empty, incoming is returned as-is, so callers should use the return value: accumulated = accumulated.merged_with(incoming).
OmniRequestOutput dataclass ¶
Bases: RequestOutput
Unified request output for both pipeline stages and diffusion models.
Extends vLLM's RequestOutput so that omni outputs can flow directly through vLLM serving codepaths (which expect prompt_token_ids, outputs, etc. as real attributes). The inherited fields store the LLM generation content; omni-specific fields store pipeline/diffusion extras.
Note: RequestOutput is a plain class (not a dataclass), so all of its attributes are redeclared below as dataclass fields with defaults — the dataclass-generated __init__ replaces RequestOutput.__init__ and must set them itself.
This class handles outputs from: 1. Multi-stage LLM pipelines (with stage_id, final_output_type, and the inherited RequestOutput fields carrying the stage's generation content) 2. Diffusion models (with images, prompt, metrics)
Attributes:
| Name | Type | Description |
|---|---|---|
request_id | str | Unique identifier for this request |
finished | bool | Whether generation is complete |
stage_id | int | None | Identifier of the stage that produced this output (pipeline mode) |
replica_id | int | None | Identifier of the stage replica that produced this output |
final_output_type | str | Type of output ("text", "image", "audio", "latents") |
images | list[Image] | List of generated PIL images (diffusion mode) |
prompt | OmniPromptType | None | The prompt used for generation |
latents | Tensor | None | Optional tensor of latent representations (diffusion mode) |
metrics | Any | Generation metrics. A plain dict for omni outputs; may carry vLLM's request stats object when copied from a raw RequestOutput. |
custom_output property writable ¶
Return custom output data from diffusion pipelines.
ec_transfer_params class-attribute instance-attribute ¶
encoder_prompt_token_ids class-attribute instance-attribute ¶
kv_transfer_params class-attribute instance-attribute ¶
multimodal_output property ¶
multimodal_output: Any
Return the multimodal output payload.
Checks completion outputs first (where multimodal_output is attached by AR stages), then the local _multimodal_output field.
Returns either a MultimodalPayload (Phase 3+) or a plain dict (legacy).
num_cache_creation_tokens class-attribute instance-attribute ¶
num_cache_creation_tokens: int | None = None
outputs class-attribute instance-attribute ¶
stage_durations class-attribute instance-attribute ¶
trajectory_log_probs class-attribute instance-attribute ¶
trajectory_timesteps class-attribute instance-attribute ¶
from_diffusion classmethod ¶
from_diffusion(
request_id: str,
images: list[Image],
prompt: OmniPromptType | None = None,
metrics: dict[str, Any] | None = None,
latents: Tensor | None = None,
trajectory_latents: Tensor | None = None,
trajectory_timesteps: Tensor | None = None,
trajectory_log_probs: Tensor | None = None,
trajectory_decoded: list | None = None,
multimodal_output: dict[str, Any] | None = None,
custom_output: dict[str, Any] | None = None,
final_output_type: str = "image",
stage_durations: dict[str, float] | None = None,
peak_memory_mb: float = 0.0,
finished: bool = True,
) -> OmniRequestOutput
Create output from diffusion model.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
request_id | str | Request identifier | required |
images | list[Image] | Generated images | required |
prompt | OmniPromptType | None | The prompt used | None |
metrics | dict[str, Any] | None | Generation metrics | None |
latents | Tensor | None | Optional latent tensors | None |
trajectory_latents | Tensor | None | Optional stacked trajectory latent tensors | None |
trajectory_timesteps | Tensor | None | Optional stacked trajectory timestep tensors | None |
trajectory_log_probs | Tensor | None | Optional stacked trajectory log-probability tensors | None |
trajectory_decoded | list | None | Optional list of decoded trajectory images | None |
multimodal_output | dict[str, Any] | None | Optional multimodal output dict | None |
custom_output | dict[str, Any] | None | Optional custom output dict (e.g. prompt embeds) | None |
stage_durations | dict[str, float] | None | Optional stage durations (execution time of each stage) dict | None |
peak_memory_mb | float | Peak memory usage in MB | 0.0 |
Returns:
| Type | Description |
|---|---|
OmniRequestOutput | OmniRequestOutput configured for diffusion mode |
from_error classmethod ¶
from_error(
request_id: str,
error_message: str,
*,
status_code: int | None = None,
error_type: str | None = None,
) -> OmniRequestOutput
Create a terminal error output.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
request_id | str | Request identifier | required |
error_message | str | Human-readable error description | required |
Returns:
| Type | Description |
|---|---|
OmniRequestOutput | OmniRequestOutput with |
from_stage_output classmethod ¶
from_stage_output(
source: RequestOutput, **kwargs: Any
) -> OmniRequestOutput
Create an OmniRequestOutput from a stage's raw output.
Copies generation content (outputs, prompt, prompt_token_ids, finished, images, latents, etc.) from source onto the returned object. source may be a vLLM RequestOutput, another OmniRequestOutput (which inherits from RequestOutput).
This is the preferred way to construct an OmniRequestOutput that wraps a stage result.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source | RequestOutput | The stage output whose content is copied onto the new object. | required |
**kwargs | Any | Passed through to the dataclass constructor ( | {} |
Returns:
| Type | Description |
|---|---|
OmniRequestOutput | A new |
OmniRequestState ¶
Bases: RequestState
Request state for omni models, tracking multimodal outputs.
Extends the base RequestState with support for accumulating multimodal tensor outputs (e.g., images, audio, latents) that are produced incrementally during generation.
native_text_stats instance-attribute ¶
native_text_stats = RequestStateStats(
arrival_time=float(arrival_time or 0.0)
)
add_multimodal_tensor ¶
Accumulate a multimodal tensor payload into the request state.
Normalizes incoming payloads (dict or raw tensor) into a MultimodalPayload and merges with any previously accumulated data. Uses list-based deferred concatenation to avoid O(n²) repeated torch.cat calls.
make_request_output ¶
make_request_output(
new_token_ids: list[int],
pooling_output: Tensor | None,
finish_reason: FinishReason | None,
stop_reason: int | str | None,
kv_transfer_params: dict[str, Any] | None = None,
ec_transfer_params: dict[str, Any] | None = None,
*,
routed_experts: Any = None,
) -> OmniRequestOutput | PoolingRequestOutput | None
Create a request output from generation results.
Creates a RequestOutput or PoolingRequestOutput from the generated tokens and accumulated multimodal outputs. Attaches multimodal tensors to the completion output if available.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
new_token_ids | list[int] | List of newly generated token IDs | required |
pooling_output | Tensor | None | Optional pooling output tensor | required |
finish_reason | FinishReason | None | Optional finish reason indicating why generation stopped | required |
stop_reason | int | str | None | Optional stop reason (token ID or stop string) | required |
kv_transfer_params | dict[str, Any] | None | Optional KV cache transfer parameters | None |
ec_transfer_params | dict[str, Any] | None | Optional encoder-cache transfer parameters (6th positional, matching upstream RequestState so that super().process_outputs() calls do not misroute it). | None |
routed_experts | Any | Optional MoE routed-expert ids for this step, attached to the completion output for generation stages (omni-specific keyword; upstream moved this accumulation into RequestState.routed_experts_chunks). | None |
Returns:
| Type | Description |
|---|---|
OmniRequestOutput | PoolingRequestOutput | None | OmniRequestOutput or PoolingRequestOutput if output should be |
OmniRequestOutput | PoolingRequestOutput | None | emitted (based on finish status and output kind), None otherwise |
OutputModality ¶
Bases: Flag
Bit-flag enum for output modalities.
Compose freely with | — no need to enumerate every combination.
Single: OutputModality.TEXT, OutputModality.IMAGE, ... Compound: OutputModality.TEXT | OutputModality.IMAGE (text+image)
Note: POOLING is intentionally excluded. Pooling/embedding is vLLM's native path (pooling_output → PoolingRequestOutput), handled entirely by the base OutputProcessor. vLLM-Omni's layer does not participate.
from_string classmethod ¶
from_string(s: str | None) -> OutputModality
Parse a free-text modality string into an OutputModality flag.
Handles common aliases and compound strings separated by + or ,.
Examples::
OutputModality.from_string("text+image")
# → OutputModality.TEXT | OutputModality.IMAGE
drain_delta_payload ¶
drain_delta_payload(payload: MultimodalPayload) -> None
Remove client-facing delta data while retaining request-level state.
is_non_final_delta_audio_chunk ¶
is_non_final_delta_audio_chunk(
payload: MultimodalPayload, mm_type: str | None
) -> bool
Return whether an audio delta explicitly declares more chunks.
replace_snapshot_keys ¶
replace_snapshot_keys(
accumulated: MultimodalPayload,
incoming: MultimodalPayload,
) -> None
Replace per-chunk metadata with the latest values.