Skip to content

vllm_omni.diffusion.sched

Modules:

Name Description
base_scheduler
interface
request_scheduler
sigma_schedule
step_scheduler

BASE_SCHEDULE_KEY module-attribute

BASE_SCHEDULE_KEY = 'base_schedule'

Scheduler module-attribute

Scheduler = RequestScheduler

BaseScheduler

Bases: ABC

Shared queue/state bookkeeping for diffusion schedulers.

kv_connector property

kv_connector

Upstream vLLM Scheduler-role connector, when configured.

max_num_running_reqs instance-attribute

max_num_running_reqs: int = 1

od_config instance-attribute

od_config: OmniDiffusionConfig | None = None

add_request

add_request(request: OmniDiffusionRequest) -> str

close

close() -> None

completed_kv_drains

completed_kv_drains() -> set[str]

fail_incomplete_kv_loads

fail_incomplete_kv_loads(
    transfer_ids: set[str],
) -> set[str]

finish_requests

finish_requests(
    request_ids: str | list[str],
    status: DiffusionRequestStatus,
) -> None

get_admission_wait_decision

get_admission_wait_decision(
    *, now: float, dp_concurrent: bool = False
) -> _AdmissionWaitDecision

Return the admission-delay policy for the next scheduling wave.

get_diffusion_kv_cleanup_targets

get_diffusion_kv_cleanup_targets(
    request_ids: list[str],
) -> list[str | tuple[str, int]]

get_request_state

get_request_state(
    request_id: str,
) -> SchedulerRequestState | None

has_requests

has_requests() -> bool

initialize

initialize(
    od_config: OmniDiffusionConfig,
    *,
    kv_cache_config: KVCacheConfig | None = None,
    scheduler_block_size: int | None = None,
    hash_block_size: int | None = None,
    kv_vllm_config: VllmConfig | None = None,
) -> None

native_kv_poll_output

native_kv_poll_output(
    *, drain_request_ids: list[str] | None = None
) -> DiffusionSchedulerOutput | None

Poll after compute; cancellation/close may wait for selected loads.

num_running_requests

num_running_requests() -> int

num_waiting_requests

num_waiting_requests() -> int

pending_finished_request_ids

pending_finished_request_ids() -> set[str]

Finished requests whose state the engine has not consumed yet.

pop_request_state

pop_request_state(
    request_id: str,
) -> SchedulerRequestState | None

preempt_request

preempt_request(request_id: str) -> bool

release_kv_drains

release_kv_drains(request_ids: set[str]) -> None

Called only after all ranks completed and Worker row cleanup succeeded.

schedule

should_end_admission_wait

should_end_admission_wait(
    decision: _AdmissionWaitDecision,
    *,
    now: float,
    stable_since: float,
) -> bool

Return whether an active admission delay should end.

update_from_output abstractmethod

update_from_output(
    sched_output: DiffusionSchedulerOutput,
    output: BaseRunnerOutput,
) -> set[str]

update_kv_connector_output

update_kv_connector_output(
    output: KVConnectorOutput | None,
) -> None

CachedRequestData dataclass

Cached diffusion requests that only need their request ids resent.

request_ids instance-attribute

request_ids: list[str]

make_empty classmethod

make_empty() -> CachedRequestData

DMD2SigmaSchedule dataclass

Continuous rectified-flow positions pinned by a distilled checkpoint.

A DMD2 student only ever sees the few noise levels it was trained on, so a distilled release ships the exact positions instead of letting the server derive a uniform schedule from num_inference_steps.

This is deliberately distinct from vllm_omni.diffusion.models.dmd2.DMD2Config.denoising_timesteps, which carries integer scheduler timesteps for scheduler-backed pipelines. Here the entries are continuous positions in [0, 1] that still need a per-modality time shift applied, which is what lets one schedule drive several coupled modalities at different shift scales.

base_schedule instance-attribute

base_schedule: tuple[float, ...]

num_inference_steps property

num_inference_steps: int

Denoising steps, i.e. one per interval between sigma boundaries.

from_metadata classmethod

from_metadata(
    metadata: Mapping[str, Any],
    *,
    key: str = BASE_SCHEDULE_KEY,
) -> DMD2SigmaSchedule | None

Read a schedule from checkpoint metadata.

An absent key means the release is not distilled and keeps the legacy uniform schedule. An explicitly empty value is a malformed contract and is rejected rather than silently falling back.

from_positions classmethod

from_positions(
    base_schedule: Sequence[float],
) -> DMD2SigmaSchedule

shifted_sigmas

shifted_sigmas(shift_scale: float) -> list[float]

Apply the rectified-flow time shift for one modality.

DiffusionRequestStatus

Bases: IntEnum

Request status tracked by diffusion scheduler.

FINISHED_ABORTED class-attribute instance-attribute

FINISHED_ABORTED = enum.auto()

FINISHED_COMPLETED class-attribute instance-attribute

FINISHED_COMPLETED = enum.auto()

FINISHED_ERROR class-attribute instance-attribute

FINISHED_ERROR = enum.auto()

PREEMPTED class-attribute instance-attribute

PREEMPTED = enum.auto()

RUNNING class-attribute instance-attribute

RUNNING = enum.auto()

WAITING class-attribute instance-attribute

WAITING = enum.auto()

is_finished staticmethod

is_finished(status: DiffusionRequestStatus) -> bool

DiffusionSchedulerOutput dataclass

Output of a single scheduling cycle.

finished_req_ids instance-attribute

finished_req_ids: set[str]

has_sync_kv_loads property

has_sync_kv_loads: bool

is_empty property

is_empty: bool

kv_connector_metadata class-attribute instance-attribute

kv_connector_metadata: KVConnectorMetadata | None = None

kv_finished_request_ids class-attribute instance-attribute

kv_finished_request_ids: set[str] = field(
    default_factory=set
)

kv_poll_only class-attribute instance-attribute

kv_poll_only: bool = False

kv_prefetch_connector_metadata class-attribute instance-attribute

kv_prefetch_connector_metadata: (
    KVConnectorMetadata | None
) = None

kv_prefetch_job class-attribute instance-attribute

kv_prefetch_job: KVPrefetchJob | None = None

kv_prefetch_request_ids class-attribute instance-attribute

kv_prefetch_request_ids: set[str] = field(
    default_factory=set
)

kv_required_request_ids class-attribute instance-attribute

kv_required_request_ids: set[str] | None = None

kv_transfer_request_ids class-attribute instance-attribute

kv_transfer_request_ids: set[str] = field(
    default_factory=set
)

num_running_reqs instance-attribute

num_running_reqs: int

num_scheduled_reqs property

num_scheduled_reqs: int

num_waiting_reqs instance-attribute

num_waiting_reqs: int

scheduled_cached_reqs instance-attribute

scheduled_cached_reqs: CachedRequestData

scheduled_new_reqs instance-attribute

scheduled_new_reqs: list[NewRequestData]

scheduled_request_ids cached property

scheduled_request_ids: list[str]

All scheduled request ids in this cycle, including both new and cached ones.

step_id instance-attribute

step_id: int

KVPrefetchJob

Bases: TypedDict

Descriptor for prefetching the next request's received KV cache.

kv_sender_info instance-attribute

kv_sender_info: dict[str, Any]

request_id instance-attribute

request_id: str

NewRequestData dataclass

Payload for a newly scheduled diffusion request.

Carries the already-initialized request object so executors and workers do not re-run OmniDiffusionRequest.__post_init__ and mutate sentinel-based fields like guidance_scale_provided.

diffusion_kv_metadata class-attribute instance-attribute

diffusion_kv_metadata: DiffusionKVMetadata | None = None

req instance-attribute

request_id instance-attribute

request_id: str

from_state classmethod

from_state(
    state: SchedulerRequestState,
    *,
    diffusion_kv_metadata: DiffusionKVMetadata
    | None = None,
) -> NewRequestData

RequestScheduler

Bases: BaseScheduler

Scheduler for static request waves, including admission coalescing.

get_admission_wait_decision

get_admission_wait_decision(
    *, now: float, dp_concurrent: bool = False
) -> _AdmissionWaitDecision

should_end_admission_wait

should_end_admission_wait(
    decision: _AdmissionWaitDecision,
    *,
    now: float,
    stable_since: float,
) -> bool

update_from_output

update_from_output(
    sched_output: DiffusionSchedulerOutput,
    output: BaseRunnerOutput,
) -> set[str]

SchedulerInterface

Bases: BaseScheduler

Deprecated compatibility base for custom scheduler injection.

Prefer subclassing :class:BaseScheduler directly. Subclassing this name still works but emits a :class:DeprecationWarning.

SchedulerRequestState dataclass

Scheduler-owned state for one queued OmniDiffusionRequest.

diffusion_kv_requests class-attribute instance-attribute

diffusion_kv_requests: tuple[DiffusionKVRequest, ...] = ()

error class-attribute instance-attribute

error: str | None = None

queued_at class-attribute instance-attribute

queued_at: float = 0.0

req instance-attribute

request_id instance-attribute

request_id: str

sampling_params_key class-attribute instance-attribute

sampling_params_key: (
    StepBatchSamplingParamsKey
    | RequestBatchSamplingParamsKey
    | None
) = None

status class-attribute instance-attribute

is_finished

is_finished() -> bool

StepBatchSamplingParamsKey dataclass

Denoise step level Batch-compatibility key derived from OmniDiffusionSamplingParams.

Only requests with the same key can be batched together. Fields not included here are treated as request-local and do not participate in the current homogeneous batching policy.

boundary_ratio class-attribute instance-attribute

boundary_ratio: float | None = None

cfg_normalize class-attribute instance-attribute

cfg_normalize: bool = False

condition_key class-attribute instance-attribute

condition_key: tuple[Any, ...] | None = None

do_classifier_free_guidance class-attribute instance-attribute

do_classifier_free_guidance: bool = False

fps class-attribute instance-attribute

fps: int | None = None

frame_rate class-attribute instance-attribute

frame_rate: float | None = None

guidance_rescale class-attribute instance-attribute

guidance_rescale: float = 0.0

guidance_scale class-attribute instance-attribute

guidance_scale: float = 0.0

guidance_scale_2 class-attribute instance-attribute

guidance_scale_2: float | None = None

guidance_scale_provided class-attribute instance-attribute

guidance_scale_provided: bool = False

height class-attribute instance-attribute

height: int | None = None

lora_int_id class-attribute instance-attribute

lora_int_id: int | None = None

lora_scale class-attribute instance-attribute

lora_scale: float = 1.0

num_frames class-attribute instance-attribute

num_frames: int = 1

num_outputs_per_prompt class-attribute instance-attribute

num_outputs_per_prompt: int = 1

quality class-attribute instance-attribute

quality: str | None = None

resolution class-attribute instance-attribute

resolution: int | str | None = None

true_cfg_scale class-attribute instance-attribute

true_cfg_scale: float | None = None

use_step_execution class-attribute instance-attribute

use_step_execution: bool = True

width class-attribute instance-attribute

width: int | None = None

StepScheduler

Bases: BaseScheduler

Scheduler that advances each request by one denoise step per update.

add_request

add_request(request: OmniDiffusionRequest) -> str

update_from_output

update_from_output(
    sched_output: DiffusionSchedulerOutput,
    output: BaseRunnerOutput,
) -> set[str]