diff --git a/conf/cmx-daria/hook/daria_cmx_host_prep.toml b/conf/cmx-daria/hook/daria_cmx_host_prep.toml new file mode 100644 index 000000000..ee4eb3043 --- /dev/null +++ b/conf/cmx-daria/hook/daria_cmx_host_prep.toml @@ -0,0 +1,6 @@ +name = "daria_cmx_host_prep" + +[[Tests]] +id = "daria-cmx-host-prep" +test_name = "daria-cmx-host-prep" +time_limit = "00:10:00" diff --git a/conf/cmx-daria/hook/test/daria_cmx_host_prep.toml b/conf/cmx-daria/hook/test/daria_cmx_host_prep.toml new file mode 100644 index 000000000..651086dcf --- /dev/null +++ b/conf/cmx-daria/hook/test/daria_cmx_host_prep.toml @@ -0,0 +1,27 @@ +name = "daria-cmx-host-prep" +description = "Daria CMX host-side pre-test checks for Lustre, RDMA UCM, and cached images" +test_template_name = "BashCmd" + +[cmd_args] +cmd = """ +set -euo pipefail + +echo Daria CMX host prep on $(hostname) +echo SLURM_JOB_NODELIST=${SLURM_JOB_NODELIST:-} + +test -d /mnt/lustre +test -r /mnt/lustre +test -d /mnt/lustre/users/spothula/cloudai/.install +find /mnt/lustre/users/spothula/cloudai/.install -maxdepth 1 -name '*.sqsh' -print -quit | grep -q . + +srun --overlap --ntasks=${SLURM_JOB_NUM_NODES} --ntasks-per-node=1 bash -lc ' + set -euo pipefail + echo node=\\$(hostname) + sudo -n modprobe rdma_ucm + test -e /dev/infiniband/rdma_cm + test -d /mnt/lustre + test -r /mnt/lustre + test -d /dev/hugepages + ls -l /dev/ng* 2>/dev/null || true +' +""" diff --git a/conf/cmx-daria/system/daria_b300_scctl_spothula.toml b/conf/cmx-daria/system/daria_b300_scctl_spothula.toml new file mode 100644 index 000000000..d13aa8309 --- /dev/null +++ b/conf/cmx-daria/system/daria_b300_scctl_spothula.toml @@ -0,0 +1,46 @@ +# Daria system configuration for launches through a Daria compute node. +# Personal-path variant of Juntao Wang's CMS b300_scctl.toml. + +name = "daria-b300-scctl-spothula" +scheduler = "slurm" + +install_path = "/mnt/lustre/users/spothula/cloudai/.install" +output_path = "/mnt/lustre/users/spothula/cloudai/results" +hf_home_path = "/mnt/lustre/cloudai/huggingface" + +default_partition = "daria" +gpus_per_node = 4 +mpi = "pmix" +monitor_interval = 10 +extra_srun_args = "--container-remap-root --container-writable" +extra_sbatch_args = ["--cpus-per-task=256"] + +cache_docker_images_locally = true + +[[partitions]] +name = "daria" + + [[partitions.groups]] + name = "b300_nodes" + nodes = [ + "daria01", + "daria02", + "daria03", + "daria04", + "daria05", + "daria06", + "daria07", + "daria08", + ] + +[global_env_vars] +HF_HOME = "/root/.cache/huggingface/" +HF_HUB_CACHE = "/root/.cache/huggingface/hub" + +MELLANOX_VISIBLE_DEVICES = "all" + +NCCL_DEBUG = "WARN" +NCCL_IB_DISABLE = "0" +NCCL_NET_GDR_LEVEL = "5" +NCCL_IB_GID_INDEX = "3" +NCCL_SOCKET_IFNAME = "^lo,docker" diff --git a/conf/cmx-daria/test/agentic_gpt_oss20_doca_memos_single.toml b/conf/cmx-daria/test/agentic_gpt_oss20_doca_memos_single.toml new file mode 100644 index 000000000..4c31a875b --- /dev/null +++ b/conf/cmx-daria/test/agentic_gpt_oss20_doca_memos_single.toml @@ -0,0 +1,98 @@ +name = "agentic-gpt-oss20-doca-memos-single" +description = "GPT-OSS-20B agentic workload: LMCache local CPU plus CMX DOCA_MEMOS, single node" +test_template_name = "AIDynamo" +extra_container_mounts = [ + "/dev:/dev", + "/sys/class/nvme:/sys/class/nvme", +] + +[cmd_args] +docker_image_url = "artifactory.nvidia.com/sw-nbu-swx-storage-devops-docker-local/nixl/dynamo-vllm-nixl-cmx:1.2.0-cuda13-lmcache-389b9cfc-doca-41a1f020-nixl-cd2f6d49-hnogit-vllm0.20.2-build24-x86_64" +workloads = "aiperf.sh" + + [cmd_args.doca-memos-preflight] + enabled = true + health-check = true + setup-hugepages = true + skip-data-path-check = false + + [cmd_args.lmcache] + chunk_size = 512 + local_cpu = true + max_local_cpu_size = 25.0 + local_cpu_use_hugepages = true + nixl_buffer_device = "cpu" + + [cmd_args.lmcache.extra_config] + enable_nixl_storage = true + nixl_presence_cache = true + nixl_presence_cache_only = true + nixl_async_put = true + nixl_sync_mode = "STRICT" + nixl_pool_size = 0 + nixl_backend = "DOCA_MEMOS" + + [cmd_args.lmcache.extra_config.nixl_backend_params] + query_mem_mode = "actual" + + [cmd_args.dynamo] + backend = "vllm" + model = "openai/gpt-oss-20b" + ingress-cmd = "python3 -m dynamo.frontend --no-kv-events --router-reset-states --router-mode kv --router-ttl 86400 --kv-overlap-score-weight 1000" + + [cmd_args.dynamo.prefill_worker] + num-nodes = 0 + cmd = "python3 -m dynamo.vllm --is-prefill-worker" + worker-initialized-regex = "VllmWorker.*has.been.initialized" + + [cmd_args.dynamo.prefill_worker.args] + tensor-parallel-size = 4 + pipeline-parallel-size = 1 + gpu-memory-utilization = 0.4 + kv-transfer-config = '{"kv_connector":"LMCacheConnectorV1","kv_role":"kv_both"}' + + [cmd_args.dynamo.decode_worker] + num-nodes = 1 + cmd = "python3 -m dynamo.vllm" + worker-initialized-regex = "VllmWorker.*has.been.initialized" + extra-args = "--kv-cache-metrics --kv-cache-metrics-sample 0.01" + + [cmd_args.dynamo.decode_worker.args] + tensor-parallel-size = 4 + pipeline-parallel-size = 1 + gpu-memory-utilization = 0.4 + kv-transfer-config = '{"kv_connector":"LMCacheConnectorV1","kv_role":"kv_both"}' + + [cmd_args.aiperf] + [cmd_args.aiperf.args] + server-metrics = "auto" + endpoint-type = "chat" + streaming = true + export-level = "records" + slice-duration = 30 + server-metrics-formats = "csv,json,jsonl" + extra-inputs = "ignore_eos:true" + dataset-sampling-strategy = "sequential" + export-http-trace = true + user-centric-rate = 0.6666666667 + num-users = 40 + session-turns-mean = 15 + user-context-prompt-length = 64000 + shared-system-prompt-length = 1000 + synthetic-input-tokens-mean = 2048 + synthetic-input-tokens-stddev = 0 + osl = 64 + num-dataset-entries = 40 + conversation-num = 40 + concurrency = 40 + max-workers = 100 + warmup-request-count = 5 + +[extra_env_vars] +AIPERF_DATASET_CONFIGURATION_TIMEOUT = "900" +AIPERF_SERVICE_PROFILE_CONFIGURE_TIMEOUT = "900" +HF_HUB_OFFLINE = "1" +PYTHONHASHSEED = "0" +PROMETHEUS_MULTIPROC_DIR = "/tmp/lmcache_prometheus" +UCX_TLS = "all" +UCX_HANDLE_ERRORS = "bt,freeze" diff --git a/conf/cmx-daria/test_scenario/agentic_gpt_oss20_doca_memos_single.toml b/conf/cmx-daria/test_scenario/agentic_gpt_oss20_doca_memos_single.toml new file mode 100644 index 000000000..3cd470a6d --- /dev/null +++ b/conf/cmx-daria/test_scenario/agentic_gpt_oss20_doca_memos_single.toml @@ -0,0 +1,12 @@ +name = "agentic-gpt-oss20-doca-memos-single-review" + +[[Tests]] +id = "agentic-gpt-oss20-doca-memos-single-review" +test_name = "agentic-gpt-oss20-doca-memos-single" +nodes = ["daria04"] +time_limit = "24:00:00" + + [Tests.cmd_args] + + [Tests.cmd_args.dynamo.decode_worker] + nodes = "daria04" diff --git a/src/cloudai/workloads/ai_dynamo/ai_dynamo.py b/src/cloudai/workloads/ai_dynamo/ai_dynamo.py index d3726fda1..c5f90eb18 100644 --- a/src/cloudai/workloads/ai_dynamo/ai_dynamo.py +++ b/src/cloudai/workloads/ai_dynamo/ai_dynamo.py @@ -17,7 +17,7 @@ import csv import logging from pathlib import Path -from typing import Literal, Optional, cast +from typing import Any, Literal, Optional, cast from pydantic import ( AliasChoices, @@ -162,6 +162,7 @@ class AIDynamoArgs(BaseModel): model: str = "Qwen/Qwen3-0.6B" backend: Literal["vllm", "sglang", "sglang_dsr1"] = "vllm" + mode: Literal["disaggregated", "aggregate"] = "disaggregated" endpoint: str = Field(default="v1/chat/completions") connector: Optional[str | list[str]] = None @@ -381,6 +382,46 @@ class Constraints(BaseModel): tp_times_pp_le_gpus_per_node: bool = True +class DocaMemosPreflight(BaseModel): + """Optional DOCA_MEMOS-specific preflight checks for AI Dynamo Slurm runs.""" + + model_config = ConfigDict(extra="forbid", populate_by_name=True) + + enabled: bool = False + health_check: bool = Field( + default=True, + serialization_alias="health-check", + validation_alias=AliasChoices("health-check", "health_check"), + ) + setup_hugepages: bool = Field( + default=False, + serialization_alias="setup-hugepages", + validation_alias=AliasChoices("setup-hugepages", "setup_hugepages"), + ) + skip_data_path_check: bool = Field( + default=False, + serialization_alias="skip-data-path-check", + validation_alias=AliasChoices("skip-data-path-check", "skip_data_path_check"), + ) + probe_size_bytes: int = Field( + default=6 * 1024 * 1024, + gt=0, + serialization_alias="probe-size-bytes", + validation_alias=AliasChoices("probe-size-bytes", "probe_size_bytes"), + ) + transfer_timeout_sec: float = Field( + default=45.0, + gt=0, + serialization_alias="transfer-timeout-sec", + validation_alias=AliasChoices("transfer-timeout-sec", "transfer_timeout_sec"), + ) + srun_timeout_sec: int = Field( + default=90, + gt=0, + serialization_alias="srun-timeout-sec", + validation_alias=AliasChoices("srun-timeout-sec", "srun_timeout_sec"), + ) + class LMCacheController(BaseModel): """Optional LMCache controller process to launch on the frontend node.""" @@ -402,6 +443,11 @@ class AIDynamoCmdArgs(CmdArgs): hicache: dict | None = None lmcache: dict | None = None lmcache_controller: LMCacheController | None = None + doca_memos_preflight: DocaMemosPreflight = Field( + default_factory=DocaMemosPreflight, + serialization_alias="doca-memos-preflight", + validation_alias=AliasChoices("doca-memos-preflight", "doca_memos_preflight"), + ) genai_perf: GenAIPerf = Field(default_factory=GenAIPerf) aiperf: AIPerf = Field(default_factory=AIPerf) aiperf_phases: list[AIPerfPhase] | None = None @@ -418,6 +464,46 @@ def validate_workloads(cls, v: str) -> str: raise ValueError(f"Invalid workload: {workload}. Available workloads: {allowed_workloads}") return ",".join(values) + @staticmethod + def _is_enabled_value(value: Any) -> bool: + if isinstance(value, list): + return all(AIDynamoCmdArgs._is_enabled_value(item) for item in value) + return value is True + + @staticmethod + def _is_doca_memos_backend(value: Any) -> bool: + if isinstance(value, list): + return all(AIDynamoCmdArgs._is_doca_memos_backend(item) for item in value) + return isinstance(value, str) and value.casefold() == "doca_memos" + + @model_validator(mode="after") + def validate_doca_memos_preflight(self) -> "AIDynamoCmdArgs": + """Validate DOCA_MEMOS preflight is only enabled for DOCA_MEMOS LMCache configurations.""" + preflight = self.doca_memos_preflight + if not preflight.enabled: + return self + if not preflight.health_check and not preflight.setup_hugepages: + raise ValueError("doca_memos_preflight.enabled requires health_check or setup_hugepages") + if self.lmcache is None: + raise ValueError("doca_memos_preflight requires cmd_args.lmcache") + + extra_config = self.lmcache.get("extra_config") + if not isinstance(extra_config, dict): + raise ValueError("doca_memos_preflight requires cmd_args.lmcache.extra_config") + if not self._is_enabled_value(extra_config.get("enable_nixl_storage")): + raise ValueError("doca_memos_preflight requires enable_nixl_storage=true") + if not self._is_doca_memos_backend(extra_config.get("nixl_backend")): + raise ValueError('doca_memos_preflight requires nixl_backend="DOCA_MEMOS"') + + backend_params = extra_config.get("nixl_backend_params") or {} + if not isinstance(backend_params, dict): + raise ValueError("doca_memos_preflight requires nixl_backend_params to be a mapping when set") + device_name = str(backend_params.get("device_name", "")).strip() + if not preflight.health_check and (not device_name or device_name.casefold() == "auto"): + raise ValueError("automatic DOCA_MEMOS device discovery requires doca_memos_preflight.health_check=true") + + return self + @property def workloads_list(self) -> list[str]: return [w.strip() for w in self.workloads.split(",")] diff --git a/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh b/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh index f0612b2db..4342aefe8 100644 --- a/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh +++ b/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh @@ -43,6 +43,7 @@ SHARED_NODE_DISAGG="false" declare -A dynamo_args dynamo_args["backend"]="vllm" +dynamo_args["mode"]="disaggregated" dynamo_args["node-setup-cmd"]="" dynamo_args["ingress-cmd"]="python -m dynamo.frontend --router-mode kv" dynamo_args["port"]=$((8080 + SLURM_JOBID % 100)) @@ -307,7 +308,11 @@ _apply_connector_settings() { _patch_dynamo_args() { if [[ -z "${dynamo_args["frontend-node"]}" ]]; then - dynamo_args["frontend-node"]=$(echo "${decode_config["node-list"]}" | cut -d',' -f1) + if [[ -n "${decode_config["node-list"]}" ]]; then + dynamo_args["frontend-node"]=$(echo "${decode_config["node-list"]}" | cut -d',' -f1) + else + dynamo_args["frontend-node"]=$(echo "${prefill_config["node-list"]}" | cut -d',' -f1) + fi fi dynamo_args["url"]="http://${dynamo_args["frontend-node"]}:${dynamo_args["port"]}" @@ -564,6 +569,19 @@ _current_node_name() { echo "${SLURMD_NODENAME:-$(hostname)}" } +_select_node_lmcache_config() { + [[ -n "${SLURM_JOB_ID:-}" ]] || return 0 + + local node_name + local node_config + node_name="$(_current_node_name)" + node_config="${RESULTS_DIR}/lmcache-config-${SLURM_JOB_ID}-${node_name}.yaml" + if [[ -f "$node_config" ]]; then + export LMCACHE_CONFIG_FILE="$node_config" + log "Using node-local LMCache config: $LMCACHE_CONFIG_FILE" + fi +} + _is_frontend_node() { local name="$(_current_node_name)" [[ ",${dynamo_args["frontend-node"]}," == *",$name,"* ]] @@ -594,6 +612,7 @@ _is_aiperf_accuracy_enabled() { _init_runtime_env() { if _is_vllm || _is_sglang; then export HF_HOME="${HUGGINGFACE_HOME}" + export HF_HUB_CACHE="${HUGGINGFACE_HOME%/}/hub" hf cache scan || echo "HF cache scan failed" fi export NATS_SERVER="nats://${dynamo_args["frontend-node"]}:${dynamo_args["nats-port"]}" @@ -1432,6 +1451,7 @@ function main() cd "$RESULTS_DIR" || { log "ERROR: Failed to cd to $RESULTS_DIR"; exit 1; } + _select_node_lmcache_config render_lmcache_config log_gpu_utilization & diff --git a/src/cloudai/workloads/ai_dynamo/doca_memos_health_check.py b/src/cloudai/workloads/ai_dynamo/doca_memos_health_check.py new file mode 100644 index 000000000..289e9f894 --- /dev/null +++ b/src/cloudai/workloads/ai_dynamo/doca_memos_health_check.py @@ -0,0 +1,362 @@ +# SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES +# Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +"""Run a small synchronous DOCA_MEMOS write/query/read health check.""" + +from __future__ import annotations + +import argparse +import contextlib +import ctypes +import importlib +import json +import mmap +import os +import re +import stat +import sys +import time +import traceback +import uuid +from pathlib import Path +from typing import Any + +_NVME_GENERIC_DEVICE = re.compile(r"^ng(?P\d+)n\d+$") +_NVME_BLOCK_DEVICE = re.compile(r"^nvme(?P\d+)n\d+$") +_NVME_CONTROLLER = re.compile(r"^nvme(?P\d+)$") +_RUNTIME_COMPONENT = re.compile(r"^[A-Za-z0-9_.-]+$") +_AUTO_DEVICE_NAME = "auto" +_NODE_CONFIG_PREFIX = "lmcache-config" +_HUGEPAGE_SIZE = 2 * 1024 * 1024 +_MAP_HUGETLB = 0x40000 +_MAP_HUGE_SHIFT = 26 +_MAP_HUGE_2MB = 21 << _MAP_HUGE_SHIFT + + +@contextlib.contextmanager +def _hugepage_buffer(size: int) -> Any: + """Allocate a Linux 2 MiB hugetlb mapping and yield its address.""" + if size <= 0: + raise ValueError("hugepage buffer size must be positive") + + allocation_size = ((size + _HUGEPAGE_SIZE - 1) // _HUGEPAGE_SIZE) * _HUGEPAGE_SIZE + libc = ctypes.CDLL(None, use_errno=True) + mmap_fn = libc.mmap + mmap_fn.argtypes = [ + ctypes.c_void_p, + ctypes.c_size_t, + ctypes.c_int, + ctypes.c_int, + ctypes.c_int, + ctypes.c_long, + ] + mmap_fn.restype = ctypes.c_void_p + munmap_fn = libc.munmap + munmap_fn.argtypes = [ctypes.c_void_p, ctypes.c_size_t] + munmap_fn.restype = ctypes.c_int + + flags = mmap.MAP_PRIVATE | mmap.MAP_ANONYMOUS | _MAP_HUGETLB | _MAP_HUGE_2MB + address = mmap_fn( + None, + allocation_size, + mmap.PROT_READ | mmap.PROT_WRITE, + flags, + -1, + 0, + ) + if address == ctypes.c_void_p(-1).value: + error = ctypes.get_errno() + raise OSError(error, f"Unable to allocate {allocation_size} bytes from the 2 MiB hugetlb pool") + + print( + f"DOCA_MEMOS_HUGEPAGE_BUFFER_OK address=0x{address:x} size={size} allocation_size={allocation_size}", + flush=True, + ) + try: + yield address + finally: + if munmap_fn(address, allocation_size) != 0: + error = ctypes.get_errno() + raise OSError(error, f"Unable to unmap hugetlb buffer at 0x{address:x}") + + +def _controller_name(device: Path) -> str: + """Return the sysfs controller name for an NVMe namespace device.""" + for pattern in (_NVME_GENERIC_DEVICE, _NVME_BLOCK_DEVICE): + match = pattern.fullmatch(device.name) + if match: + return f"nvme{match.group('controller')}" + raise ValueError(f"Cannot derive an NVMe controller from device name {device}") + + +def discover_doca_memos_device( + sysfs_root: Path = Path("/sys/class/nvme"), + dev_root: Path = Path("/dev"), +) -> str: + """Return the sole generic namespace for a live DOCA/SNAP controller.""" + candidates: list[Path] = [] + for controller_path in sorted(sysfs_root.glob("nvme*"), key=lambda path: path.name): + match = _NVME_CONTROLLER.fullmatch(controller_path.name) + if not match: + continue + try: + state = (controller_path / "state").read_text().strip() + model = (controller_path / "model").read_text().strip() + except OSError: + continue + if state != "live" or ("DOCA" not in model and "SNAP" not in model): + continue + + controller_index = match.group("controller") + for device in sorted(dev_root.glob(f"ng{controller_index}n*"), key=lambda path: path.name): + if _NVME_GENERIC_DEVICE.fullmatch(device.name): + candidates.append(device) + + if len(candidates) != 1: + devices = ", ".join(str(path) for path in candidates) or "none" + raise RuntimeError( + "Expected exactly one generic namespace for a live DOCA/SNAP controller, " + f"found {len(candidates)}: {devices}" + ) + + device = str(candidates[0]) + print(f"DOCA_MEMOS_DEVICE_DISCOVERED device={device}", flush=True) + return device + + +def _runtime_identity() -> tuple[str, str]: + """Return safe Slurm job and node identifiers for a generated config.""" + job_id = os.environ.get("SLURM_JOB_ID", "").strip() + node_name = os.environ.get("SLURMD_NODENAME", "").strip() or os.uname().nodename + for label, value in (("SLURM_JOB_ID", job_id), ("node name", node_name)): + if not value or not _RUNTIME_COMPONENT.fullmatch(value): + raise RuntimeError(f"Cannot generate node-local LMCache config: invalid {label} {value!r}") + return job_id, node_name + + +def materialize_node_lmcache_config( + base_config_path: Path, + output_dir: Path, + device_name: str, +) -> Path: + """Write a job- and node-specific LMCache config using the discovered device.""" + config = json.loads(base_config_path.read_text()) + if not isinstance(config, dict): + raise ValueError("LMCache base config must be a JSON object") + extra_config = config.get("extra_config") + if not isinstance(extra_config, dict): + raise ValueError("LMCache base config extra_config must be a mapping") + backend_params = extra_config.get("nixl_backend_params") + if backend_params is None: + backend_params = {} + extra_config["nixl_backend_params"] = backend_params + elif not isinstance(backend_params, dict): + raise ValueError("LMCache base config nixl_backend_params must be a mapping") + backend_params["device_name"] = device_name + + job_id, node_name = _runtime_identity() + output_dir.mkdir(parents=True, exist_ok=True) + output_path = output_dir / f"{_NODE_CONFIG_PREFIX}-{job_id}-{node_name}.yaml" + temporary_path = output_dir / f".{output_path.name}.{os.getpid()}.tmp" + try: + temporary_path.write_text(json.dumps(config, indent=2, sort_keys=False) + "\n") + os.replace(temporary_path, output_path) + finally: + with contextlib.suppress(FileNotFoundError): + temporary_path.unlink() + print(f"DOCA_MEMOS_LMCACHE_CONFIG_WRITTEN path={output_path} device={device_name}", flush=True) + return output_path + + +def validate_device(device_name: str) -> None: + """Fail unless the configured DOCA namespace and controller look usable.""" + device = Path(device_name) + device_stat = device.stat() + if not (stat.S_ISCHR(device_stat.st_mode) or stat.S_ISBLK(device_stat.st_mode)): + raise RuntimeError(f"{device} is not a character or block device") + if not os.access(device, os.R_OK | os.W_OK): + raise PermissionError(f"{device} must be readable and writable") + + controller = _controller_name(device) + controller_path = Path("/sys/class/nvme") / controller + state = (controller_path / "state").read_text().strip() + if state != "live": + raise RuntimeError(f"{controller} is not live (state={state!r})") + + model = (controller_path / "model").read_text().strip() + if "DOCA" not in model and "SNAP" not in model: + raise RuntimeError(f"{controller} model {model!r} is not a DOCA/SNAP controller") + print(f"DOCA_MEMOS_DEVICE_OK device={device} controller={controller} model={model!r}", flush=True) + + +def _load_nixl() -> tuple[Any, Any, Any]: + """Load either the standard or CUDA-versioned NIXL Python package.""" + errors: list[str] = [] + for package in ("nixl", "nixl_cu13"): + try: + utils = importlib.import_module(f"{package}._utils") + api = importlib.import_module(f"{package}._api") + return utils, api.nixl_agent, api.nixl_agent_config + except ImportError as exc: + errors.append(f"{package}: {exc}") + raise ImportError("Unable to import a NIXL Python package: " + "; ".join(errors)) + + +def _wait_for_transfer(agent: Any, handle: Any, timeout_seconds: float) -> None: + """Wait for one NIXL transfer and fail on timeout or backend error.""" + state = agent.transfer(handle) + deadline = time.monotonic() + timeout_seconds + while state not in {"DONE", "ERR"}: + if time.monotonic() >= deadline: + raise TimeoutError(f"NIXL transfer did not finish within {timeout_seconds:.1f}s") + time.sleep(0.001) + state = agent.check_xfer_state(handle) + if state != "DONE": + raise RuntimeError(f"NIXL transfer finished in state {state!r}") + + +def _backend_params(raw_params: str) -> dict[str, str]: + parsed = json.loads(raw_params) + if not isinstance(parsed, dict): + raise ValueError("backend parameters must be a JSON object") + params = {str(key): str(value) for key, value in parsed.items()} + params["query_mem_mode"] = "actual" + params["num_tasks"] = "1" + return params + + +def run_health_check(backend_params: dict[str, str], size: int, transfer_timeout: float) -> None: + """Write, query, and read back one object through DOCA_MEMOS.""" + device_name = backend_params.get("device_name") + if not device_name: + raise ValueError("DOCA_MEMOS backend parameters require device_name") + if size <= 0: + raise ValueError("probe size must be positive") + if transfer_timeout <= 0: + raise ValueError("transfer timeout must be positive") + + validate_device(device_name) + _, nixl_agent, nixl_agent_config = _load_nixl() + agent_name = f"CloudAIDocaMemosHealth{os.getpid()}_{time.time_ns()}" + agent = nixl_agent( + agent_name, + nixl_agent_config( + backends=[], + enable_prog_thread=True, + enable_listen_thread=False, + ), + ) + agent.create_backend("DOCA_MEMOS", backend_params) + + with _hugepage_buffer(size) as src_addr, _hugepage_buffer(size) as dst_addr: + ctypes.memset(src_addr, 0xA5, size) + ctypes.memset(dst_addr, 0, size) + + local_reg = None + object_reg = None + handles: list[Any] = [] + key = uuid.uuid4().hex + try: + local_reg = agent.register_memory( + [ + (src_addr, size, 0, ""), + (dst_addr, size, 0, ""), + ], + "DRAM", + backends=["DOCA_MEMOS"], + ) + object_reg = agent.register_memory( + [(0, size, 1, key)], + "OBJ", + backends=["DOCA_MEMOS"], + ) + + object_xfer = object_reg.trim() + src_xfer = agent.get_xfer_descs([(src_addr, size, 0)], "DRAM") + dst_xfer = agent.get_xfer_descs([(dst_addr, size, 0)], "DRAM") + + write_handle = agent.initialize_xfer("WRITE", src_xfer, object_xfer, agent_name) + handles.append(write_handle) + started = time.monotonic() + _wait_for_transfer(agent, write_handle, transfer_timeout) + print(f"DOCA_MEMOS_WRITE_OK seconds={time.monotonic() - started:.6f}", flush=True) + + started = time.monotonic() + query = agent.query_memory( + [(0, 0, 1, key)], + "DOCA_MEMOS", + mem_type="OBJ", + ) + if not query or query[0] is None: + raise RuntimeError("DOCA_MEMOS query did not find the object just written") + print(f"DOCA_MEMOS_QUERY_OK seconds={time.monotonic() - started:.6f}", flush=True) + + read_handle = agent.initialize_xfer("READ", dst_xfer, object_xfer, agent_name) + handles.append(read_handle) + started = time.monotonic() + _wait_for_transfer(agent, read_handle, transfer_timeout) + if ctypes.string_at(src_addr, size) != ctypes.string_at(dst_addr, size): + raise RuntimeError("DOCA_MEMOS read-back data does not match the write") + print(f"DOCA_MEMOS_READ_OK seconds={time.monotonic() - started:.6f}", flush=True) + validate_device(device_name) + finally: + for handle in handles: + with contextlib.suppress(Exception): + agent.release_xfer_handle(handle) + if object_reg is not None: + with contextlib.suppress(Exception): + agent.deregister_memory(object_reg, backends=["DOCA_MEMOS"]) + if local_reg is not None: + with contextlib.suppress(Exception): + agent.deregister_memory(local_reg, backends=["DOCA_MEMOS"]) + + print(f"DOCA_MEMOS_HEALTH_CHECK_OK size={size}", flush=True) + + +def _parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--backend-params-json", required=True) + parser.add_argument("--base-lmcache-config-json", type=Path) + parser.add_argument("--output-config-dir", type=Path) + parser.add_argument("--skip-data-path-check", action="store_true") + parser.add_argument("--size", type=int, default=6 * 1024 * 1024) + parser.add_argument("--transfer-timeout", type=float, default=45.0) + return parser.parse_args() + + +def main() -> int: + args = _parse_args() + try: + backend_params = _backend_params(args.backend_params_json) + device_name = backend_params.get("device_name", "").strip() + if not device_name or device_name.casefold() == _AUTO_DEVICE_NAME: + device_name = discover_doca_memos_device() + backend_params["device_name"] = device_name + if args.base_lmcache_config_json is None or args.output_config_dir is None: + raise ValueError( + "automatic DOCA_MEMOS discovery requires --base-lmcache-config-json and --output-config-dir" + ) + materialize_node_lmcache_config( + args.base_lmcache_config_json, + args.output_config_dir, + device_name, + ) + + if args.skip_data_path_check: + validate_device(device_name) + print("DOCA_MEMOS_DATA_PATH_CHECK_SKIPPED", flush=True) + else: + run_health_check( + backend_params, + size=args.size, + transfer_timeout=args.transfer_timeout, + ) + except Exception as exc: + print(f"DOCA_MEMOS_HEALTH_CHECK_FAILED: {exc}", file=sys.stderr, flush=True) + traceback.print_exc() + return 1 + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/cloudai/workloads/ai_dynamo/slurm_command_gen_strategy.py b/src/cloudai/workloads/ai_dynamo/slurm_command_gen_strategy.py index 0c09c1f7b..e5164fa06 100644 --- a/src/cloudai/workloads/ai_dynamo/slurm_command_gen_strategy.py +++ b/src/cloudai/workloads/ai_dynamo/slurm_command_gen_strategy.py @@ -14,8 +14,12 @@ # See the License for the specific language governing permissions and # limitations under the License. +import json import logging +import math +import os import shlex +import shutil import textwrap from pathlib import Path from typing import Any, List, cast @@ -38,11 +42,18 @@ ) AIPERF_SCRIPT_FILE_NAME = "aiperf.sh" +DOCA_MEMOS_HEALTH_CHECK_FILE_NAME = "doca_memos_health_check.py" +DOCA_MEMOS_BASE_LMCACHE_CONFIG_FILE_NAME = "lmcache-config.base.json" +HUGEPAGE_SIZE_MIB = 2 +LOCAL_CPU_HUGEPAGE_BUFFER_RATIO = 1.10 +NORMAL_HUGEPAGES_ENV = "CLOUDAI_NORMAL_NR_HUGEPAGES" class AIDynamoSlurmCommandGenStrategy(SlurmCommandGenStrategy): """Command generation strategy for AI Dynamo on Slurm systems.""" + CONTAINER_MOUNT_HF_HOME = "/root/.cache/huggingface" + @property def td(self) -> AIDynamoTestDefinition: return cast(AIDynamoTestDefinition, self.test_run.test) @@ -59,6 +70,11 @@ def _container_mounts(self) -> list[str]: @property def final_env_vars(self) -> dict[str, str | list[str]]: env_vars = super().final_env_vars + _, node_list = self.get_cached_nodes_spec() + if node_list: + env_vars["DYNAMO_NODELIST"] = ",".join(node_list) + else: + env_vars["DYNAMO_NODELIST"] = "$(scontrol show hostname $SLURM_JOB_NODELIST | paste -sd, -)" if self.td.cmd_args.hicache is not None: env_vars["HICACHE_CONFIG_FILE"] = f"{self.CONTAINER_MOUNT_OUTPUT}/{HICACHE_CONFIG_FILE_NAME}" if self.td.cmd_args.lmcache is not None: @@ -167,6 +183,61 @@ def _prepare_lmcache_config(self): (self.test_run.output_path / LMCACHE_CONFIG_FILE_NAME).write_text(config) (self.test_run.output_path / LMCACHE_CONFIG_BACKUP_FILE_NAME).write_text(config) + def _lmcache_config(self) -> dict[str, Any] | None: + config = self.td.cmd_args.lmcache + if config is None: + return None + if not isinstance(config, dict): + raise ValueError("LMCache config must be a mapping") + return config + + def _doca_memos_preflight_enabled(self) -> bool: + return self.td.cmd_args.doca_memos_preflight.enabled + + def _doca_memos_health_check_enabled(self) -> bool: + preflight = self.td.cmd_args.doca_memos_preflight + return preflight.enabled and preflight.health_check + + def _doca_memos_hugepage_setup_enabled(self) -> bool: + preflight = self.td.cmd_args.doca_memos_preflight + return preflight.enabled and preflight.setup_hugepages + + def _doca_memos_backend_params(self) -> dict[str, Any]: + config = self._lmcache_config() + if config is None: + raise ValueError("DOCA_MEMOS preflight requires LMCache config") + + extra_config = config.get("extra_config") + if not isinstance(extra_config, dict): + raise ValueError("DOCA_MEMOS preflight requires LMCache extra_config") + + backend_params = extra_config.get("nixl_backend_params") or {} + if not isinstance(backend_params, dict): + raise ValueError("DOCA_MEMOS preflight requires nixl_backend_params to be a mapping when set") + return dict(backend_params) + + def _doca_memos_device_is_auto(self) -> bool: + device_name = str(self._doca_memos_backend_params().get("device_name", "")).strip() + return not device_name or device_name.casefold() == "auto" + + def _prepare_doca_memos_health_check_script(self) -> Path: + self.test_run.output_path.mkdir(parents=True, exist_ok=True) + source_path = Path(__file__).parent / DOCA_MEMOS_HEALTH_CHECK_FILE_NAME + destination_path = self.test_run.output_path / DOCA_MEMOS_HEALTH_CHECK_FILE_NAME + shutil.copyfile(source_path, destination_path) + destination_path.chmod(0o755) + return destination_path + + def _prepare_doca_memos_base_lmcache_config(self) -> Path: + config = self._lmcache_config() + if config is None: + raise ValueError("DOCA_MEMOS automatic device discovery requires LMCache config") + + self.test_run.output_path.mkdir(parents=True, exist_ok=True) + base_config_path = self.test_run.output_path / DOCA_MEMOS_BASE_LMCACHE_CONFIG_FILE_NAME + base_config_path.write_text(json.dumps(config, indent=2, sort_keys=False) + "\n") + return base_config_path + def _prepare_hicache_config(self) -> None: if self.td.cmd_args.hicache is None: return @@ -488,6 +559,174 @@ def _gen_srun_command(self) -> str: return f"{startup_command}\n{main_command}" return main_command + @staticmethod + def _max_int(value: Any, default: int = 1) -> int: + if value is None: + return default + if isinstance(value, list): + return max((AIDynamoSlurmCommandGenStrategy._max_int(item, default) for item in value), default=default) + return int(value) + + @staticmethod + def _max_float(value: Any) -> float: + if isinstance(value, list): + if not value: + raise ValueError("Expected a non-empty numeric list") + return max(AIDynamoSlurmCommandGenStrategy._max_float(item) for item in value) + return float(value) + + @staticmethod + def _normal_hugepages() -> int: + raw_value = os.environ.get(NORMAL_HUGEPAGES_ENV, "0") + try: + value = int(raw_value) + except ValueError as exc: + raise ValueError(f"{NORMAL_HUGEPAGES_ENV} must be an integer, got {raw_value!r}") from exc + if value < 0: + raise ValueError(f"{NORMAL_HUGEPAGES_ENV} must be >= 0, got {value}") + return value + + def _desired_hugepages(self) -> int | None: + config = self._lmcache_config() + if not config: + return None + + if config.get("local_cpu_use_hugepages"): + max_local_cpu_size = config.get("max_local_cpu_size") + if max_local_cpu_size is None: + raise ValueError("local_cpu_use_hugepages requires max_local_cpu_size in the LMCache config") + + prefill_worker = self.td.cmd_args.dynamo.prefill_worker + max_tp = self._max_int(prefill_worker.args.tensor_parallel_size if prefill_worker else None, default=1) + size_gib = self._max_float(max_local_cpu_size) + if max_tp <= 0 or size_gib <= 0: + raise ValueError( + f"local_cpu_use_hugepages requires positive TP and max_local_cpu_size, got {max_tp=} {size_gib=}" + ) + return math.ceil(max_tp * size_gib * LOCAL_CPU_HUGEPAGE_BUFFER_RATIO * 1024 / HUGEPAGE_SIZE_MIB) + + extra_config = config.get("extra_config") or {} + if isinstance(extra_config, dict) and extra_config.get("nixl_use_hugepages"): + return None + + return self._normal_hugepages() + + def _gen_doca_memos_hugepage_setup_block(self) -> list[str]: + if not self._doca_memos_hugepage_setup_enabled(): + return [] + + desired_hugepages = self._desired_hugepages() + if desired_hugepages is None: + return [] + + num_nodes, node_list = self.get_cached_nodes_spec() + out_dir = self.test_run.output_path.absolute() + setup_command = ( + f"target={desired_hugepages}; " + 'sudo -n /usr/sbin/sysctl -w "vm.nr_hugepages=$target" >/dev/null; ' + "actual=$(cat /proc/sys/vm/nr_hugepages); " + 'if [ "$actual" -ne "$target" ]; then ' + 'echo "nr_hugepages=$actual, expected=$target" >&2; exit 1; ' + "fi; " + 'printf "node=%s nr_hugepages=%s\\n" "$(hostname)" "$actual"; ' + "grep -E 'HugePages_Total|HugePages_Free|Hugepagesize' /proc/meminfo" + ) + srun_parts = [ + "srun", + "--export=ALL", + f"--nodes={num_nodes}", + *([] if not node_list else [f"--nodelist={','.join(node_list)}"]), + f"--ntasks={num_nodes}", + "--ntasks-per-node=1", + "--kill-on-bad-exit=1", + f"--output={out_dir / 'doca-memos-hugepages-node-%n-stdout.txt'}", + f"--error={out_dir / 'doca-memos-hugepages-node-%n-stderr.txt'}", + "bash", + "-lc", + shlex.quote(setup_command), + ] + srun_command = " \\\n ".join(srun_parts) + return [ + "# Set host hugepages for DOCA_MEMOS LMCache.", + *self._wrap_fail_closed(srun_command, "DOCA_MEMOS hugepage setup failed; refusing to start Dynamo"), + ] + + def _gen_doca_memos_health_check_block(self) -> list[str]: + if not self._doca_memos_health_check_enabled(): + return [] + + self._prepare_doca_memos_health_check_script() + preflight = self.td.cmd_args.doca_memos_preflight + backend_params = self._doca_memos_backend_params() + health_cmd = [ + "timeout", + f"{preflight.srun_timeout_sec}s", + "python3", + f"{self.CONTAINER_MOUNT_OUTPUT}/{DOCA_MEMOS_HEALTH_CHECK_FILE_NAME}", + "--backend-params-json", + json.dumps(backend_params, separators=(",", ":")), + "--size", + str(preflight.probe_size_bytes), + "--transfer-timeout", + str(preflight.transfer_timeout_sec), + ] + if self._doca_memos_device_is_auto(): + self._prepare_doca_memos_base_lmcache_config() + health_cmd.extend( + [ + "--base-lmcache-config-json", + f"{self.CONTAINER_MOUNT_OUTPUT}/{DOCA_MEMOS_BASE_LMCACHE_CONFIG_FILE_NAME}", + "--output-config-dir", + self.CONTAINER_MOUNT_OUTPUT, + ] + ) + if preflight.skip_data_path_check: + health_cmd.append("--skip-data-path-check") + + num_nodes, node_list = self.get_cached_nodes_spec() + out_dir = self.test_run.output_path.absolute() + srun_parts = [ + *self._gen_aux_srun_prefix(self.image_path()), + f"--nodes={num_nodes}", + *([] if not node_list else [f"--nodelist={','.join(node_list)}"]), + f"--ntasks={num_nodes}", + "--ntasks-per-node=1", + "--kill-on-bad-exit=1", + f"--output={out_dir / 'doca-memos-health-node-%n-stdout.txt'}", + f"--error={out_dir / 'doca-memos-health-node-%n-stderr.txt'}", + "bash", + "-lc", + shlex.quote(shlex.join(health_cmd)), + ] + srun_command = " \\\n ".join(srun_parts) + return [ + "# Run DOCA_MEMOS/NIXL health check on every node before AI Dynamo.", + *self._wrap_fail_closed(srun_command, "DOCA_MEMOS preflight failed; refusing to start Dynamo"), + ] + + @staticmethod + def _wrap_fail_closed(command: str, failure_message: str) -> list[str]: + lines = command.splitlines() + if not lines: + return [] + lines[0] = f"if ! {lines[0]}" + lines[-1] = f"{lines[-1]}; then" + return [ + *lines, + f" echo {shlex.quote(failure_message)} >&2", + " exit 1", + "fi", + ] + + def _gen_doca_memos_preflight_block(self) -> list[str]: + if not self._doca_memos_preflight_enabled(): + return [] + return [ + "ulimit -l unlimited", + *self._gen_doca_memos_hugepage_setup_block(), + *self._gen_doca_memos_health_check_block(), + ] + def _gen_dcgm_launcher_block(self) -> list[str]: dcgm_image = self.td.dcgm_exporter_image if not dcgm_image: @@ -593,6 +832,10 @@ def gen_exec_command(self) -> str: command_list.extend([pre_test_command, "if [ $PRE_TEST_SUCCESS -eq 1 ]; then"]) indent = " " + doca_memos_preflight_block = self._gen_doca_memos_preflight_block() + if doca_memos_preflight_block: + command_list.extend(f"{indent}{line}" for line in doca_memos_preflight_block) + dcgm_block = self._gen_dcgm_launcher_block() if dcgm_block: command_list.extend(f"{indent}{line}" for line in dcgm_block) diff --git a/tests/ref_data/ai-dynamo.sbatch b/tests/ref_data/ai-dynamo.sbatch index c00906d40..a62e9e1d9 100644 --- a/tests/ref_data/ai-dynamo.sbatch +++ b/tests/ref_data/ai-dynamo.sbatch @@ -10,13 +10,13 @@ export SLURM_JOB_MASTER_NODE=$(scontrol show hostname $SLURM_JOB_NODELIST | head -n 1) -srun --export=ALL --mpi=pmix -N2 --container-image=nvcr.io/nvidia/ai-dynamo:24.09 --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/cloudai_install/huggingface,/tmp:/tmp --output=__OUTPUT_DIR__/output/mapping-stdout.txt --error=__OUTPUT_DIR__/output/mapping-stderr.txt bash -c "echo \$(date): \$(hostname):node \${SLURM_NODEID}:rank \${SLURM_PROCID}." +srun --export=ALL --mpi=pmix -N2 --container-image=nvcr.io/nvidia/ai-dynamo:24.09 --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/root/.cache/huggingface,/tmp:/tmp --output=__OUTPUT_DIR__/output/mapping-stdout.txt --error=__OUTPUT_DIR__/output/mapping-stderr.txt bash -c "echo \$(date): \$(hostname):node \${SLURM_NODEID}:rank \${SLURM_PROCID}." -srun --export=ALL --mpi=pmix -N2 --container-image=nvcr.io/nvidia/ai-dynamo:24.09 --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/cloudai_install/huggingface,/tmp:/tmp --ntasks=2 --ntasks-per-node=1 --output=__OUTPUT_DIR__/output/metadata/node-%N.toml --error=__OUTPUT_DIR__/output/metadata/nodes.err bash /cloudai_install/slurm-metadata.sh +srun --export=ALL --mpi=pmix -N2 --container-image=nvcr.io/nvidia/ai-dynamo:24.09 --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/root/.cache/huggingface,/tmp:/tmp --ntasks=2 --ntasks-per-node=1 --output=__OUTPUT_DIR__/output/metadata/node-%N.toml --error=__OUTPUT_DIR__/output/metadata/nodes.err bash /cloudai_install/slurm-metadata.sh # Start DCGM exporter on each node. echo "Starting DCGM exporter..." -srun --export=ALL --mpi=pmix --container-image=nvcr.io/nvidia/k8s/dcgm-exporter:4.5.2-4.8.1-distroless --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/cloudai_install/huggingface,/tmp:/tmp --overlap --job-name=cloudai-dcgm-exporter -N2 --ntasks=2 --ntasks-per-node=1 --output=__OUTPUT_DIR__/output/dcgm-node-%n-stdout.txt --error=__OUTPUT_DIR__/output/dcgm-node-%n-stderr.txt bash -lc 'DCGM_EXPORTER_LISTEN=:9501 dcgm-exporter' & +srun --export=ALL --mpi=pmix --container-image=nvcr.io/nvidia/k8s/dcgm-exporter:4.5.2-4.8.1-distroless --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/root/.cache/huggingface,/tmp:/tmp --overlap --job-name=cloudai-dcgm-exporter -N2 --ntasks=2 --ntasks-per-node=1 --output=__OUTPUT_DIR__/output/dcgm-node-%n-stdout.txt --error=__OUTPUT_DIR__/output/dcgm-node-%n-stderr.txt bash -lc 'DCGM_EXPORTER_LISTEN=:9501 dcgm-exporter' & DCGM_EXPORTER_SRUN_PID=$! echo "DCGM exporter srun PID: ${DCGM_EXPORTER_SRUN_PID}" DCGM_EXPORTER_STEP_ID= @@ -66,7 +66,7 @@ srun \ --mpi=pmix \ -N2 \ --container-image=nvcr.io/nvidia/ai-dynamo:24.09 \ - --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/cloudai_install/huggingface,/tmp:/tmp \ + --container-mounts=__OUTPUT_DIR__/output:/cloudai_run_results,__INSTALL_DIR__:/cloudai_install,__OUTPUT_DIR__/output,__INSTALL_DIR__/huggingface:/root/.cache/huggingface,/tmp:/tmp \ --nodes=2 \ --ntasks=2 \ --ntasks-per-node=1 \ @@ -78,7 +78,7 @@ srun \ --install-dir /cloudai_install \ --results-dir /cloudai_run_results \ --dynamo-repo /cloudai_install/dynamo__f7e468c7e8ff0d1426db987564e60572167e8464 \ - --hf-home /cloudai_install/huggingface \ + --hf-home /root/.cache/huggingface \ --workloads aiperf.sh \ --failure-marker /cloudai_run_results/failure-marker.txt \ --success-marker /cloudai_run_results/success-marker.txt \ diff --git a/tests/workloads/ai_dynamo/test_command_gen_strategy_slurm.py b/tests/workloads/ai_dynamo/test_command_gen_strategy_slurm.py index e0c3d8146..076bd07ad 100644 --- a/tests/workloads/ai_dynamo/test_command_gen_strategy_slurm.py +++ b/tests/workloads/ai_dynamo/test_command_gen_strategy_slurm.py @@ -14,6 +14,7 @@ # See the License for the specific language governing permissions and # limitations under the License. +import json import shlex from pathlib import Path from typing import cast @@ -41,6 +42,7 @@ WorkerBaseArgs, WorkerConfig, ) +from cloudai.workloads.ai_dynamo.ai_dynamo import DocaMemosPreflight @pytest.fixture @@ -651,6 +653,109 @@ def test_gen_script_args_writes_lmcache_object_as_yaml(strategy: AIDynamoSlurmCo assert "--lmcache" not in result +def test_doca_memos_preflight_requires_doca_memos_lmcache(cmd_args: AIDynamoCmdArgs) -> None: + data = cmd_args.model_dump() + data["doca_memos_preflight"] = {"enabled": True} + + with pytest.raises(ValueError, match=r"requires cmd_args\.lmcache"): + AIDynamoCmdArgs.model_validate(data) + + data["lmcache"] = { + "extra_config": { + "enable_nixl_storage": True, + "nixl_backend": "POSIX", + } + } + + with pytest.raises(ValueError, match='requires nixl_backend="DOCA_MEMOS"'): + AIDynamoCmdArgs.model_validate(data) + + +def test_gen_exec_command_includes_doca_memos_preflight(strategy: AIDynamoSlurmCommandGenStrategy) -> None: + td = cast(AIDynamoTestDefinition, strategy.test_run.test) + td.cmd_args.lmcache = { + "chunk_size": 512, + "local_cpu": True, + "local_cpu_use_hugepages": True, + "max_local_cpu_size": 25.0, + "extra_config": { + "enable_nixl_storage": True, + "nixl_backend": "DOCA_MEMOS", + "nixl_backend_params": {"device_name": "auto"}, + }, + } + td.cmd_args.doca_memos_preflight = DocaMemosPreflight(enabled=True, setup_hugepages=True) + + strategy.gen_exec_command() + + script = (strategy.test_run.output_path / "cloudai_sbatch_script.sh").read_text() + base_config = json.loads((strategy.test_run.output_path / "lmcache-config.base.json").read_text()) + assert (strategy.test_run.output_path / "doca_memos_health_check.py").exists() + assert base_config["chunk_size"] == 512 + assert base_config["extra_config"]["nixl_backend_params"]["device_name"] == "auto" + assert "doca-memos-hugepages-node-%n-stdout.txt" in script + assert "target=112641;" in script + assert 'sudo -n /usr/sbin/sysctl -w "vm.nr_hugepages=$target"' in script + assert "doca_memos_health_check.py" in script + assert script.count("if ! srun \\") >= 2 + assert "DOCA_MEMOS hugepage setup failed; refusing to start Dynamo" in script + assert "DOCA_MEMOS preflight failed; refusing to start Dynamo" in script + assert "--base-lmcache-config-json" in script + assert "--output-config-dir /cloudai_run_results" in script + assert "--skip-data-path-check" not in script + assert "--nodelist=n0,n1" in script + assert "/dev:/dev" not in script + assert "/sys/class/nvme:/sys/class/nvme" not in script + assert script.index("doca-memos-hugepages-node-%n-stdout.txt") < script.index("doca_memos_health_check.py") + assert script.index("doca_memos_health_check.py") < script.index("ai_dynamo.sh") + + +def test_doca_memos_preflight_with_explicit_device_skips_node_config( + strategy: AIDynamoSlurmCommandGenStrategy, +) -> None: + td = cast(AIDynamoTestDefinition, strategy.test_run.test) + td.cmd_args.lmcache = { + "chunk_size": 512, + "extra_config": { + "enable_nixl_storage": True, + "nixl_backend": "DOCA_MEMOS", + "nixl_backend_params": {"device_name": "/dev/ng3n1"}, + }, + } + td.cmd_args.doca_memos_preflight = DocaMemosPreflight(enabled=True) + + strategy.gen_exec_command() + + script = (strategy.test_run.output_path / "cloudai_sbatch_script.sh").read_text() + assert "doca_memos_health_check.py" in script + assert "/dev/ng3n1" in script + assert "--base-lmcache-config-json" not in script + assert not (strategy.test_run.output_path / "lmcache-config.base.json").exists() + + +def test_doca_memos_preflight_uses_dse_expanded_lmcache_config( + slurm_system: SlurmSystem, + test_run: TestRun, +) -> None: + td = cast(AIDynamoTestDefinition, test_run.test) + td.cmd_args.lmcache = { + "chunk_size": [256, 512], + "extra_config": { + "enable_nixl_storage": True, + "nixl_backend": "DOCA_MEMOS", + "nixl_backend_params": {"device_name": "auto"}, + }, + } + td.cmd_args.doca_memos_preflight = DocaMemosPreflight(enabled=True) + + expanded_run = test_run.apply_params_set({"lmcache.chunk_size": 512}) + strategy = AIDynamoSlurmCommandGenStrategy(slurm_system, expanded_run) + strategy.gen_exec_command() + + base_config = json.loads((strategy.test_run.output_path / "lmcache-config.base.json").read_text()) + assert base_config["chunk_size"] == 512 + + def test_gen_script_args_writes_hicache_object_as_toml(strategy: AIDynamoSlurmCommandGenStrategy) -> None: td = cast(AIDynamoTestDefinition, strategy.test_run.test) td.cmd_args.dynamo.backend = "sglang"