Skip to content

vllm_omni.engine.output_processor

logger module-attribute

logger = init_logger(__name__)

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.

multimodal_output instance-attribute

multimodal_output = multimodal_output

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

engine_core_output_type instance-attribute

engine_core_output_type = engine_core_output_type

output_modality instance-attribute

output_modality = OutputModality.from_string(
    engine_core_output_type
)

abort_requests

abort_requests(request_ids, internal: bool) -> list[str]

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

commit_aborted_request_state(
    request_ids: Iterable[str], *, internal: bool
) -> None

Drop processor mappings after a successful physical EngineCore abort.

pop_native_text_metrics

pop_native_text_metrics(request_id: str) -> dict[str, Any]

process_outputs

process_outputs(
    engine_core_outputs: list[EngineCoreOutput],
    engine_core_timestamp: float | None = None,
    iteration_stats: IterationStats | None = None,
) -> OutputProcessorOutput

remove_request

remove_request(request_id: str) -> None

Rollback one previously registered request if it was never submitted.

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).

is_empty property

is_empty: bool

Return True if the payload has no tensors and no metadata.

metadata class-attribute instance-attribute

metadata: dict[str, Any] = field(default_factory=dict)

primary_tensor property

primary_tensor: Tensor | None

Return the first tensor in the payload, or None if empty.

tensors class-attribute instance-attribute

tensors: dict[str, Tensor] = field(default_factory=dict)

consolidate_metadata

consolidate_metadata() -> None

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(key: str, default: Any = None) -> Any

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).

to_dict

to_dict() -> dict[str, Any]

Convert back to a plain dict (tensors + metadata merged).

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

custom_output: dict[str, Any]

Return custom output data from diffusion pipelines.

ec_transfer_params class-attribute instance-attribute

ec_transfer_params: dict[str, Any] | None = None

encoder_prompt class-attribute instance-attribute

encoder_prompt: str | None = None

encoder_prompt_token_ids class-attribute instance-attribute

encoder_prompt_token_ids: list[int] | None = None

error class-attribute instance-attribute

error: str | None = None

error_status_code class-attribute instance-attribute

error_status_code: int | None = None

error_type class-attribute instance-attribute

error_type: str | None = None

final_output_type class-attribute instance-attribute

final_output_type: str = 'text'

finished class-attribute instance-attribute

finished: bool = True

images class-attribute instance-attribute

images: list[Image] = field(default_factory=list)

is_diffusion_output property

is_diffusion_output: bool

Check if this is a diffusion model output.

is_pipeline_output property

is_pipeline_output: bool

Check if this is a pipeline stage output.

kv_transfer_params class-attribute instance-attribute

kv_transfer_params: dict[str, Any] | None = None

latents class-attribute instance-attribute

latents: Tensor | None = None

lora_request class-attribute instance-attribute

lora_request: Any = None

metrics class-attribute instance-attribute

metrics: Any = field(default_factory=dict)

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

num_cached_tokens class-attribute instance-attribute

num_cached_tokens: int | None = None

num_images property

num_images: int

Return the number of generated images.

outputs class-attribute instance-attribute

outputs: list[CompletionOutput] = field(
    default_factory=list
)

peak_memory_mb class-attribute instance-attribute

peak_memory_mb: float = 0.0

prompt class-attribute instance-attribute

prompt: OmniPromptType | None = None

prompt_logprobs class-attribute instance-attribute

prompt_logprobs: Any = None

prompt_token_ids class-attribute instance-attribute

prompt_token_ids: list[int] | None = None

replica_id class-attribute instance-attribute

replica_id: int | None = None

request_id class-attribute instance-attribute

request_id: str = ''

stage_durations class-attribute instance-attribute

stage_durations: dict[str, float] = field(
    default_factory=dict
)

stage_id class-attribute instance-attribute

stage_id: int | None = None

trajectory_decoded class-attribute instance-attribute

trajectory_decoded: list | None = None

trajectory_latents class-attribute instance-attribute

trajectory_latents: Tensor | None = None

trajectory_log_probs class-attribute instance-attribute

trajectory_log_probs: Tensor | None = None

trajectory_timesteps class-attribute instance-attribute

trajectory_timesteps: Tensor | None = None

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 finished=True and the error field set.

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 (request_id, stage_id, final_output_type, metrics, stage_durations, peak_memory_mb, finished, etc.). Typed as Any because the exact set of valid keys is the dataclass field list, which is validated by cls(**kwargs) at call time.

{}

Returns:

Type Description
OmniRequestOutput

A new OmniRequestOutput with the stage's content flattened onto it.

to_dict

to_dict() -> dict[str, Any]

Convert to dictionary for JSON serialization.

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.

mm_accumulated instance-attribute

mm_type instance-attribute

mm_type: str | None = None

native_text_stats instance-attribute

native_text_stats = RequestStateStats(
    arrival_time=float(arrival_time or 0.0)
)

add_multimodal_tensor

add_multimodal_tensor(
    payload: Any | None, mm_type: str | None
) -> None

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.

apply_streaming_update

apply_streaming_update(update) -> None

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.

AUDIO class-attribute instance-attribute

AUDIO = auto()

IMAGE class-attribute instance-attribute

IMAGE = auto()

LATENT class-attribute instance-attribute

LATENT = auto()

TEXT class-attribute instance-attribute

TEXT = auto()

has_multimodal property

has_multimodal: bool

has_text property

has_text: bool

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.

unflatten_payload

unflatten_payload(flat: dict[str, Any]) -> dict[str, Any]

Unflatten dotted keys back to nested dicts.

Reverse of :func:flatten_payload. hidden_states.layer_N keys are collected into hidden_states.layers.