Skip to content

vllm_omni.diffusion.models.lingbot_world.pipeline

Request-scoped LingBot-World v2 causal DMD pipeline.

LINGBOT_DMD_TIMESTEPS module-attribute

LINGBOT_DMD_TIMESTEPS = (1000, 750, 500, 250)

logger module-attribute

logger = init_logger(__name__)

LingBotWorldCausalDMDPipeline

Bases: Module, SupportImageInput, SupportsComponentDiscovery, SupportsStepExecution, InteractionMixin, ProgressBarMixin, DiffusionPipelineProfilerMixin

LingBot-World v2 I2V generation with a request-local causal cache.

device instance-attribute

device = get_local_device()

dummy_run_num_frames class-attribute

dummy_run_num_frames: int = 0

od_config instance-attribute

od_config = od_config

scheduler instance-attribute

scheduler = FlowUniPCMultistepScheduler(**scheduler_kwargs)

supports_chunk_step_grouping class-attribute

supports_chunk_step_grouping: bool = True

supports_step_execution class-attribute

supports_step_execution: bool = True

text_encoder instance-attribute

text_encoder = from_pretrained_with_prefetch(
    UMT5EncoderModel.from_pretrained,
    model,
    subfolder="text_encoder",
    prefetch_list=subfolders,
    local_files_only=local_files_only,
    torch_dtype=dtype,
)

tokenizer instance-attribute

tokenizer = from_pretrained_with_prefetch(
    AutoTokenizer.from_pretrained,
    model,
    subfolder="tokenizer",
    prefetch_list=subfolders,
    local_files_only=local_files_only,
)

transformer instance-attribute

transformer = (
    CausalLingBotWorldTransformer3DModel.from_config(
        transformer_config,
        quant_config=getattr(
            od_config, "quantization_config", None
        ),
        prefix="transformer",
    )
)

vae instance-attribute

vae = from_pretrained_with_prefetch(
    DistributedAutoencoderKLWan.from_pretrained,
    model,
    subfolder="vae",
    prefetch_list=subfolders,
    local_files_only=local_files_only,
    torch_dtype=dtype,
)

vae_scale_factor_spatial instance-attribute

vae_scale_factor_spatial = int(
    getattr(self.vae.config, "scale_factor_spatial", 8)
)

vae_scale_factor_temporal instance-attribute

vae_scale_factor_temporal = int(
    getattr(self.vae.config, "scale_factor_temporal", 4)
)

weights_sources instance-attribute

weights_sources = [
    DiffusersPipelineLoader.ComponentSource(
        model_or_path=model,
        subfolder="transformer",
        revision=None,
        prefix="transformer.",
        fall_back_to_pt=True,
    )
]

ar_diffusion_kv_cache_spec

ar_diffusion_kv_cache_spec() -> ARDiffusionKVCacheSpec

Describe the fixed worker-local cache geometry for realtime ticks.

bind_ar_diffusion_state

bind_ar_diffusion_state(
    session_id: str, state: ARDiffusionKVState
) -> Iterator[None]

close_ar_diffusion_session

close_ar_diffusion_session(session_id: str) -> None

denoise_step

denoise_step(
    input_batch: InputBatch,
    *,
    states: Sequence[StepRequestState] | None = None,
    **kwargs: Any,
) -> Tensor | None

encode_prompt

encode_prompt(
    prompt: str, *, max_sequence_length: int, dtype: dtype
) -> Tensor

Encode prompt with UMT5, reusing the result for a prompt already encoded.

The encode depends only on the whitespace-normalised text, the sequence length and the dtype -- the tokenizer and text encoder are fixed once loaded -- so those three are the key. That covers every caller the same way: a realtime tick repeats its session's prompt on every block, a stepwise request encodes once per request, and a server tends to open many sessions with the same scene prompt.

A hit returns the stored tensor itself. Callers only read it: the cross-attention projection and the DMD transformer both work out of place, so a shared encode is never changed under another session.

forward

load_weights

load_weights(
    weights: Iterable[tuple[str, Tensor]],
) -> set[str]

peek_chunk_media

peek_chunk_media(state: StepRequestState) -> ChunkMediaSpec

Expose this chunk's decoded media extent and latent step count.

post_decode

post_decode(
    state: StepRequestState, **kwargs: Any
) -> DiffusionOutput

prepare_encode

prepare_encode(
    state: StepRequestState, **kwargs: Any
) -> StepRequestState

prepare_next_chunk

prepare_next_chunk(state: StepRequestState) -> None

Prepare the next AR block after chunk-boundary interaction apply.

reset_ar_diffusion_session

reset_ar_diffusion_session(session_id: str) -> None

step_scheduler

step_scheduler(
    state: StepRequestState,
    noise_pred: Tensor,
    **kwargs: Any,
) -> None

get_lingbot_world_post_process_func

get_lingbot_world_post_process_func(
    od_config: OmniDiffusionConfig,
) -> Callable[..., Any]

get_lingbot_world_pre_process_func

get_lingbot_world_pre_process_func(
    od_config: OmniDiffusionConfig,
) -> Callable[[OmniDiffusionRequest], OmniDiffusionRequest]

Materialize request-local files once before dispatching GPU workers.