Skip to content

vllm_omni.diffusion.diffusion_kv.kv_connector

Assemble native vLLM KV connectors for Scheduler-owned diffusion pages.

logger module-attribute

logger = init_logger(__name__)

KVReceiveProgress dataclass

Rank-local state survives consumption of Mooncake completion events.

deadlines class-attribute instance-attribute

deadlines: dict[str, float] = field(default_factory=dict)

received class-attribute instance-attribute

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

sent class-attribute instance-attribute

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

submitted class-attribute instance-attribute

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

prepare

prepare(
    active_connector: ActiveKVConnector,
    output: DiffusionSchedulerOutput,
    timeout: float,
) -> KVConnectorOutput

KVTransferRegistrationError

Bases: ValueError

A failed handoff whose destination pages have not been dispatched.

bootstrap_addr_from_kv_transfer_config

bootstrap_addr_from_kv_transfer_config(
    kv_transfer_config: KVTransferConfig | None,
) -> str | None

Read an optional connector bootstrap endpoint from native config.

build_source_kv_transfer_params

build_source_kv_transfer_params(
    *,
    transfer_id: str,
    remote_engine_id: str | None,
    remote_bootstrap_addr: str | None,
) -> dict[str, Any]

Build the opaque producer-side metadata bag for an AR request.

build_target_kv_transfer_params

build_target_kv_transfer_params(
    *,
    source_params: Mapping[str, Any],
    remote_engine_id: str | None,
    remote_bootstrap_addr: str | None,
) -> dict[str, Any]

Build the consumer-side metadata bag without interpreting pages.

commit_kv_load

commit_kv_load(
    connector: KVConnectorBase_V1,
    manager: KVCacheManager,
    requests: tuple[DiffusionKVRequest, ...],
    matched_tokens: list[int],
) -> set[str]

create_scheduler_kv_connector

create_scheduler_kv_connector(
    od_config: OmniDiffusionConfig,
    kv_cache_config: KVCacheConfig | None = None,
    vllm_config: VllmConfig | None = None,
) -> KVConnectorBase_V1 | None

Create a Scheduler-role connector when native config is present.

init_worker_kv_connector

init_worker_kv_connector(
    vllm_config: VllmConfig, kv_cache_config: KVCacheConfig
) -> None

Initialize the Worker-role connector with its rank-local cache plan.

install_mooncake_cfg_fanout

install_mooncake_cfg_fanout(
    connector: KVConnectorBase_V1,
) -> None

Account for multiple CFG destinations on one native Mooncake ticket.

vLLM 0.29 initializes need_send to the paired consumer rank count, but increments sent per consumer request. Adapt only this producer instance, after KV initialization and before serving requests. Keep the upstream address planning, transport, timeout and completion code intact.

mint_transfer_id

mint_transfer_id(request_id: str) -> str

Return the stable ticket shared by one AR -> Diffusion handoff.

native_prefetch_enabled

native_prefetch_enabled(
    od_config: OmniDiffusionConfig,
) -> bool

Validate the opt-in GPU consumer implementation before allocating pages.

parse_kv_transfer_config

parse_kv_transfer_config(
    value: object | None,
) -> KVTransferConfig | None

Normalize a YAML mapping or an existing native config.

prepare_kv_requests

prepare_kv_requests(
    requests: tuple[DiffusionKVRequest, ...],
    params: Mapping[str, Any],
) -> None

shutdown_kv_connector

shutdown_kv_connector(
    *, scheduler_connector: KVConnectorBase_V1 | None = None
) -> None

Shutdown Worker and Scheduler connector objects idempotently.

validate_kv_transfer_boundaries

validate_kv_transfer_boundaries(
    requests: tuple[DiffusionKVRequest, ...],
    matched_tokens: list[int],
) -> None

Validate every CFG row before reserving or registering destination pages.

wait_for_kv_load

wait_for_kv_load(
    active_connector: ActiveKVConnector,
    scheduler_output: DiffusionSchedulerOutput,
    timeout: float,
) -> KVConnectorOutput