nemo_rl.models.generation.vllm.vllm_sparse_refit#

Remote sparse-refit receiver lifecycle for vLLM generation workers.

Module Contents#

Classes#

_StagedSparsePayload

VllmSparseRefitReceiver

Own the optional transport server, apply queue, and relay resources.

Functions#

Data#

API#

nemo_rl.models.generation.vllm.vllm_sparse_refit.logger#

‘getLogger(…)’

nemo_rl.models.generation.vllm.vllm_sparse_refit._warn_unauthenticated_refit_server(transport: str) None#
class nemo_rl.models.generation.vllm.vllm_sparse_refit._StagedSparsePayload#

Bases: typing.NamedTuple

path: str#

None

started_at: float#

None

finished_at: float#

None

save_s: float#

None

nemo_rl.models.generation.vllm.vllm_sparse_refit._stage_sparse_payload(
serialized: bytes,
staging_dir: str,
) nemo_rl.models.generation.vllm.vllm_sparse_refit._StagedSparsePayload#
class nemo_rl.models.generation.vllm.vllm_sparse_refit.VllmSparseRefitReceiver(worker: Any)#

Own the optional transport server, apply queue, and relay resources.

Initialization

set_worker_hostnames(hostnames: list[str]) None#
start_sync_server() None#
shutdown() None#
_enqueue_sparse_payload_apply(
payload: bytes,
payload_key: tuple[str, int, int],
checksum: str,
verification_candidates: int = 0,
) dict[str, Any]#
_submit_pending_sparse_payloads() None#
_notify_refit_apply_waiters(
_future: concurrent.futures.Future[dict[str, Any]],
) None#
_collect_refit_apply_results(
futures: list[concurrent.futures.Future[dict[str, Any]]],
) dict[str, Any]#
static _refit_collective_response(
worker_results: Any,
) dict[str, Any]#
_refit_collective_rpc(
method: str,
args: tuple[Any, ...],
) Any#
update_weights_from_serialized_sparse_payloads(
serialized_payloads: tuple[bytes, ...],
) dict[str, Any]#

Apply a FIFO batch of sparse deltas through one collective RPC.

update_weights_from_staged_sparse_payloads(
staged_payloads: tuple[concurrent.futures.Future[nemo_rl.models.generation.vllm.vllm_sparse_refit._StagedSparsePayload], ...],
) dict[str, Any]#
_flush_queued_sparse_payloads() dict[str, Any]#
_prepare_sparse_refit_info(
request: dict[str, Any],
) dict[str, Any]#
async _apply_s3_manifest_payload(
manifest: dict[str, Any],
) dict[str, Any]#
_apply_zmq_payload(
compressed: bytes,
metadata: collections.abc.Mapping[str, Any],
) dict[str, Any]#
setup_api_server(app: Any) None#
report_refit_server_base_url() str | None#
start_zmq_sparse_refit_relay() str#
configure_zmq_sparse_refit_relay(relay_addresses: list[str]) None#
stop_zmq_sparse_refit_relay() None#
flush_zmq_sparse_refit_relay(
transfer_id: str,
expected_payloads: int = 0,
) dict[str, Any]#
_setup_vllm_refit_server() None#