Skip to content

vllm_omni.diffusion.utils.chunked_video

Turn committed VAE chunks into MP4 bytes without materializing the video.

The producer side is the model's business: a Wan VAE walks a causal frame-by-frame loop, a MiniMax-H3 VAE walks overlapping clips, and Wan S2V decodes one clip per autoregressive iteration. What happens to a chunk once it is final is not: every one of them quantizes to uint8, moves to the host once, and lands in a bounded encoder. That consumer lives here so a model only has to publish finished chunks -- :class:SupportsChunkedVAEDecode for producers that can be driven, :class:ChunkedVideoMP4Session for producers that drive themselves.

The move to the host runs as its own scheduled stage. A chunk is quantized on the accelerator, copied without blocking into one of a fixed number of reusable pinned host slots on a dedicated stream, and only read once its copy-completion event has fired. The producer therefore returns to decoding while the copy is still in flight, and the slot count bounds how many chunks can be in flight at once.

logger module-attribute

logger = init_logger(__name__)

ChunkLease

One chunk in flight: a host slot, its copy event, and its readers.

The slot returns to the ring only after the last reader has finished with it. Being handed to an encoder is not enough -- the encoder's queue holds a view of this buffer until its muxing thread reads it -- so release is driven by the encoder rather than by dequeuing.

discard

discard() -> None

Drop an unread chunk during teardown, keeping the copy's memory alive.

The copy may still be writing into the slot and reading from the source tensor, so this waits for it rather than letting either be reused underneath an in-flight transfer.

ready

ready() -> bool

Whether the copy has landed, without blocking on it.

release

release() -> None

Report one reader done; the last one returns the slot to the ring.

wait

wait() -> ndarray

Block until the copy has landed, then hand over the host frames.

ChunkTransferStats dataclass

What the host transfer cost, for checking that it actually overlapped.

copy_wait_seconds staying near zero while transfers climbs is the signal that copies completed behind the producer; it rising with them means the producer is outrunning the copy stream. slot_wait_seconds rising instead means the encoder is the bottleneck and backpressure has reached the producer.

copy_wait_seconds class-attribute instance-attribute

copy_wait_seconds: float = 0.0

peak_pending_bytes class-attribute instance-attribute

peak_pending_bytes: int = 0

peak_slots_in_use class-attribute instance-attribute

peak_slots_in_use: int = 0

slot_wait_seconds class-attribute instance-attribute

slot_wait_seconds: float = 0.0

transfers class-attribute instance-attribute

transfers: int = 0

ChunkedVideoMP4Session

Encode committed video chunks into one progressive MP4 per batch entry.

Push finished BCTHW chunks as the producer commits them; each batch entry gets its own bounded encoder, so host transfer and H.264 encoding overlap whatever the producer is still decoding.

audio_waveforms holds one waveform per batch entry (None for a silent entry), so a caller with several outputs per prompt repeats a request's waveform across its entries. batch_frames coalesces transfers for producers that publish finer than a transfer is worth; crop trims the decoder's padding to the requested output size. transfer_slots sets how many chunks may be in flight to the host at once.

ring property

ring: PinnedChunkRing | None

The transfer ring this session settled on, once a chunk has flushed.

abort

abort() -> None

finish

finish() -> list[bytes]

Flush what is pending and return one MP4 per batch entry.

push

push(chunk: Tensor) -> None

Queue one committed BCTHW chunk.

transfer_stats

transfer_stats() -> ChunkTransferStats | None

Per-request transfer accounting, or None if no ring was used.

PinnedChunkRing

Fixed-depth pool of reusable host slots for chunk transfers.

A slot walks free -> copying -> ready -> encoding -> free and holds exactly one chunk for that whole trip, so pending transfer bytes are bounded by the depth times the largest chunk of the request. Waiting for a free slot is where encoder backpressure reaches the producer.

depth instance-attribute

depth = depth

slots_in_use property

slots_in_use: int

Slots currently holding a chunk, whether copying, ready or encoding.

stats instance-attribute

abort

abort(error: BaseException | None = None) -> None

Fail every present and future wait for a slot.

reclaim

reclaim(slot: _HostSlot) -> None

Return a slot whose last reader is done.

transfer

transfer(frames: Tensor, *, readers: int) -> ChunkLease

Start one non-blocking copy of frames into a free slot.

chunk_to_uint8_frames

chunk_to_uint8_frames(
    chunk: Tensor, value_range: tuple[float, float]
) -> ndarray

Quantize a BCTHW chunk and block until it reaches the host.

The path a session takes when no pinned ring is available, and the one a caller outside a session gets. A session on an accelerator uses :class:PinnedChunkRing instead so the copy does not stall the producer.

decode_to_mp4

decode_to_mp4(
    vae: Any, z: Tensor, **session_kwargs: Any
) -> list[bytes]

Decode z straight into one progressive MP4 per batch entry.

Drives a VAE that declares :class:SupportsChunkedVAEDecode, so host transfer and encoding overlap the remaining decode and the full video is never materialized. Ranks that own no decode output receive no chunks and get an empty list, matching the empty tensor the full-decode path returns there.

quantize_chunk

quantize_chunk(
    chunk: Tensor, value_range: tuple[float, float]
) -> Tensor

Quantize a BCTHW chunk to contiguous BTHWC uint8 on its own device.

value_range is the interval the producer publishes, which differs per checkpoint (see SupportsChunkedVAEDecode.chunk_value_range). Quantizing on the accelerator means the transfer that follows moves the final bytes rather than float frames, and making the result contiguous there keeps the copy a single flat move instead of a strided host-side gather.