Skip to content

vllm_omni.diffusion.executor.uniproc_executor

In-process diffusion executor for single-GPU deployments.

Counterpart to :class:MultiprocDiffusionExecutor for num_gpus == 1. The worker is constructed and driven in the engine process, so none of the multiproc machinery is created:

  • no spawned worker subprocess (and no second model load / CUDA context)
  • no MessageQueue shared-memory ring buffers or zmq ipc:// sockets
  • no POSIX /dev/shm segments for output tensors (ipc.py pack/unpack)

Mirrors vLLM's UniProcExecutor (vllm/v1/executor/uniproc_executor.py).

timeout on :meth:collective_rpc is accepted for interface parity but is not enforced: the call runs on the calling thread, so a hung worker blocks it. The multiproc executor enforces deadlines via its result-queue dequeue.

logger module-attribute

logger = init_logger(__name__)

UniProcDiffusionExecutor

Bases: DiffusionExecutor

Runs a single diffusion worker inline, with no IPC of any kind.

is_dead property

is_dead: bool

check_health

check_health() -> None

collective_rpc

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

Invoke method on the single in-process worker.

unique_reply_rank / exec_all_ranks only select which ranks participate, which is degenerate at world size 1. The return shape still matches :class:MultiprocDiffusionExecutor (a bare result when a reply rank is named, a one-element list otherwise).

execute_batch

execute_batch(
    scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput

execute_request

execute_request(
    scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput

execute_step

execute_step(
    scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput

register_failure_callback

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

Register a callback invoked when the inline worker fails fatally.

The multiproc executor fires these from its process monitor. There is no process to monitor here, so they are fired from collective_rpc instead — without them a poisoned device context would leave check_health reporting healthy while every request fails.

shutdown

shutdown() -> None