Skip to content

vllm_omni.diffusion.executor.abstract

DiffusionExecutor

Bases: ABC

Abstract base class for Diffusion executors.

is_dead abstractmethod property

is_dead: bool

Whether the executor is shut down or has failed fatally.

od_config instance-attribute

od_config = od_config

uses_multiproc class-attribute instance-attribute

uses_multiproc: bool = False

check_health abstractmethod

check_health() -> None

Check if the executor and workers are healthy.

collective_rpc abstractmethod

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

Execute a method on workers.

determine_available_kv_memory

determine_available_kv_memory(
    profile_requests: list[OmniDiffusionRequest],
) -> list[int]

Profile and collect the KV memory budget on every Worker rank.

drop_output

drop_output(async_output_id: str) -> None

Reclaim an async output that will never be waited on (e.g. an aborted request).

Only executors with an async output path (result pump) cache outputs that a consumer must later claim; executors without one have nothing to reclaim and can keep the default no-op implementation.

execute_batch abstractmethod

execute_batch(
    scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput

Execute request-mode work through the request-batch path.

execute_request abstractmethod

execute_request(
    scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput

Execute request-mode work from a scheduler output.

execute_step abstractmethod

execute_step(
    scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput

Execute step-mode work from a scheduler output.

get_class staticmethod

get_class(
    od_config: OmniDiffusionConfig,
) -> type[DiffusionExecutor]

get_kv_cache_specs

get_kv_cache_specs() -> list[dict[str, KVCacheSpec]]

Collect rank-local native specs after every Worker loads its model.

prepare_kv_for_forward

prepare_kv_for_forward(
    scheduler_output: DiffusionSchedulerOutput,
) -> KVConnectorOutput | None

register_failure_callback

register_failure_callback(
    callback: Callable[[], None],
) -> None

Register a callback invoked when the executor fatally fails.

Executors without a background failure monitor can keep the default no-op implementation.

remove_diffusion_kv_requests

remove_diffusion_kv_requests(
    request_ids: list[str | tuple[str, int]],
) -> None

Clear request rows on every Worker after Scheduler retirement.

set_kv_cache_configs

set_kv_cache_configs(
    kv_cache_configs: list[KVCacheConfig],
    resolved_max_model_len: int,
) -> None

Send rank-local configs and the resolved model length to all Workers.

shutdown abstractmethod

shutdown() -> None

Shutdown the executor and release resources.

wait_output_ready

wait_output_ready(
    async_output_id: str,
) -> Future[DiffusionOutput]

Resolve deferred output; only asynchronous executors implement this.