Skip to content

vllm_omni.diffusion.ipc

IPC utilities for transferring large tensors via POSIX shared memory.

Used by Hop1 (GPU worker <-> scheduler) to avoid pickling large video tensors through the MessageQueue. Tensors above _SHM_TENSOR_THRESHOLD are copied into a named shared-memory segment; only a lightweight metadata dict is serialised through the queue.

DIFFUSION_RPC_RESULT_ENVELOPE module-attribute

DIFFUSION_RPC_RESULT_ENVELOPE = 'diffusion_rpc_result'

pack_diffusion_output_shm

pack_diffusion_output_shm(
    output: object, d2h_stream: Stream | None = None
) -> object

Replace large tensors in diffusion worker outputs with SHM handles.

Supports a bare DiffusionOutput, a wrapper object carrying one in .result (for example RunnerOutput), an RPC result envelope carrying the diffusion output in ["result"], a batch wrapper carrying RunnerOutput objects in .runner_outputs, or a DP-tagged dict {"dp_rank": int, "output": DiffusionOutput} used by DP multi-concurrency.

If d2h_stream is provided, D2H copies use that stream (non-blocking on the default stream). The caller must synchronize d2h_stream afterward.

Packing is failure-atomic: every SHM segment created during the call is tracked and unlinked if any part of the payload fails to pack, and field mutations are committed only after the complete payload succeeds.

payload_carries_typed_media

payload_carries_typed_media(output: object) -> bool

True if output carries any typed DiffusionOutput.media payload.

unpack_diffusion_output_shm

unpack_diffusion_output_shm(output: object) -> object

Reconstruct tensors from SHM handles in diffusion worker outputs.