Skip to content

vllm_omni.diffusion.sched.base_scheduler

BatchSamplingParamsKey module-attribute

logger module-attribute

logger = init_logger(__name__)

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

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.