vllm_omni.diffusion.diffusion_kv.kv_connector ¶
Assemble native vLLM KV connectors for Scheduler-owned diffusion pages.
KVReceiveProgress dataclass ¶
Rank-local state survives consumption of Mooncake completion events.
deadlines class-attribute instance-attribute ¶
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 ¶
Initialize the Worker-role connector with its rank-local cache plan.
install_mooncake_cfg_fanout ¶
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 ¶
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 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