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.
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_kv_cache_specs ¶
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 ¶
Clear request rows on every Worker after Scheduler retirement.
set_kv_cache_configs ¶
Send rank-local configs and the resolved model length to all Workers.
wait_output_ready ¶
wait_output_ready(
async_output_id: str,
) -> Future[DiffusionOutput]
Resolve deferred output; only asynchronous executors implement this.