Skip to content

vllm_omni.diffusion.diffusion_engine

logger module-attribute

logger = init_logger(__name__)

DiffusionEngine

The diffusion engine for vLLM-Omni diffusion models.

default_diffusion_model_runner_cls class-attribute instance-attribute

default_diffusion_model_runner_cls: str | None = None

dp_concurrent class-attribute instance-attribute

dp_concurrent: bool = False

execution_mode instance-attribute

execution_mode = self._resolve_execution_mode(od_config)

od_config instance-attribute

od_config = od_config

abort

abort(request_id: str | Iterable[str]) -> None

add_req_and_wait_for_response

add_req_and_wait_for_response(
    request: OmniDiffusionRequest,
) -> DiffusionOutput

add_request

add_request(request: OmniDiffusionRequest) -> str

async_add_req_and_stream_response

async_add_req_and_stream_response(
    request: OmniDiffusionRequest,
) -> AsyncGenerator[DiffusionOutput, None]

async_add_req_and_wait_for_response async

async_add_req_and_wait_for_response(
    request: OmniDiffusionRequest,
) -> DiffusionOutput

Deprecated compatibility wrapper over async_add_req_and_stream_response().

Use async_add_req_and_stream_response() for new callers. This method drains the unified output stream and returns only the final DiffusionOutput, matching the historical non-streaming behavior.

async_collective_rpc async

async_collective_rpc(
    method: str,
    timeout: float | None = None,
    args: tuple = (),
    kwargs: dict | None = None,
    unique_reply_rank: int | None = None,
) -> Any

Async variant of :meth:collective_rpc for event-loop callers.

Enqueue a task keyed by a future and await the result without blocking the loop.

close

close() -> None

collective_rpc

collective_rpc(
    method: str,
    timeout: float | None = None,
    args: tuple = (),
    kwargs: dict | None = None,
    unique_reply_rank: int | None = None,
) -> Any

Call a method on worker processes and get results immediately.

The call is enqueued and executed by the engine's busy loop between scheduler steps, so it is naturally serialized against per-request execute_fn() invocations without any explicit mutual-exclusion lock.

Parameters:

Name Type Description Default
method str

The method name (str) to execute on workers

required
timeout float | None

Optional timeout in seconds

None
args tuple

Positional arguments for the method

()
kwargs dict | None

Keyword arguments for the method

None
unique_reply_rank int | None

If set, only get reply from this rank

None

Returns:

Type Description
Any

Single result if unique_reply_rank is provided, otherwise list of results

get_output_stream async

get_output_stream(
    request_id: str,
) -> AsyncGenerator[DiffusionOutput, None]

make_engine staticmethod

make_engine(
    config: OmniDiffusionConfig,
    scheduler: BaseScheduler | None = None,
) -> DiffusionEngine

Factory method to create the engine selected by config.engine_backend.

Parameters:

Name Type Description Default
config OmniDiffusionConfig

The configuration for the diffusion engine.

required
scheduler BaseScheduler | None

Optional scheduler override. When omitted, the selected engine chooses the scheduler from its execution mode.

None

Returns:

Type Description
DiffusionEngine

An instance of the resolved DiffusionEngine (sub)class.

postprocess_output

postprocess_output(
    request: OmniDiffusionRequest, output: DiffusionOutput
) -> list[OmniRequestOutput]

Convert a DiffusionOutput to a list of OmniRequestOutput.

profile

profile(
    is_start: bool = True, profile_prefix: str | None = None
) -> None

Start or stop profiling on all diffusion workers.

Parameters:

Name Type Description Default
is_start bool

True to start profiling, False to stop.

True
profile_prefix str | None

Optional prefix for trace filename.

None

resolve_engine_class staticmethod

resolve_engine_class(
    config: OmniDiffusionConfig,
) -> type[DiffusionEngine]

Resolve the engine class selected by config.engine_backend.

Mirrors DiffusionExecutor.get_class: accepts "default", a DiffusionEngine subclass, or an import-path string (e.g. a deploy config's engine_backend). Kept separate from :meth:make_engine so the selection is testable without constructing an engine (which runs a dummy forward).

Parameters:

Name Type Description Default
config OmniDiffusionConfig

The configuration for the diffusion engine.

required

Returns:

Type Description
type[DiffusionEngine]

The DiffusionEngine (sub)class to instantiate.

run_startup_warmup

run_startup_warmup() -> None

step async

Deprecated compatibility wrapper over step_streaming().

Use step_streaming() for new callers. This method drains the unified output stream and returns only the final output batch, matching the historical non-streaming step() behavior.

step_streaming async

step_streaming(
    request: OmniDiffusionRequest,
) -> AsyncGenerator[list[OmniRequestOutput], None]

DiffusionExecutionMode

Bases: str, Enum

REQUEST_BATCH class-attribute instance-attribute

REQUEST_BATCH = 'request_batch'

STEP_BATCH class-attribute instance-attribute

STEP_BATCH = 'step_batch'

get_dummy_run_num_frames

get_dummy_run_num_frames(
    model_class_name: str, supports_audio_input: bool
) -> int

Get num_frames for the dummy warmup run. Returns 0 to skip warmup.

image_color_format

image_color_format(model_class_name: str) -> str

supports_audio_output

supports_audio_output(model_class_name: str) -> bool

supports_multimodal_input

supports_multimodal_input(
    od_config: OmniDiffusionConfig,
) -> tuple[bool, bool]

supports_request_batch

supports_request_batch(
    od_config: OmniDiffusionConfig,
) -> bool

supports_request_cancellation

supports_request_cancellation(
    od_config: OmniDiffusionConfig,
) -> bool

Whether the local pipeline checks cooperative cancellation boundaries.