diff --git a/src/cloudai/workloads/ai_dynamo/ai_dynamo.py b/src/cloudai/workloads/ai_dynamo/ai_dynamo.py index d3726fda1..2b347c773 100644 --- a/src/cloudai/workloads/ai_dynamo/ai_dynamo.py +++ b/src/cloudai/workloads/ai_dynamo/ai_dynamo.py @@ -196,6 +196,21 @@ def validate_connector(cls, v: str | list[str] | None) -> str | list[str] | None serialization_alias="node-setup-cmd", validation_alias=AliasChoices("node-setup-cmd", "node_setup_cmd"), ) + aiperf_phase_restart_services: bool = Field( + default=False, + serialization_alias="aiperf-phase-restart-services", + validation_alias=AliasChoices("aiperf-phase-restart-services", "aiperf_phase_restart_services"), + ) + aiperf_phase_setup_scope: Literal["frontend", "all"] = Field( + default="all", + serialization_alias="aiperf-phase-setup-scope", + validation_alias=AliasChoices("aiperf-phase-setup-scope", "aiperf_phase_setup_scope"), + ) + aiperf_phase_setup_cmd_scope: Literal["frontend", "all"] = Field( + default="frontend", + serialization_alias="aiperf-phase-setup-cmd-scope", + validation_alias=AliasChoices("aiperf-phase-setup-cmd-scope", "aiperf_phase_setup_cmd_scope"), + ) port: int = Field( default=8000, description="Dynamo frontend HTTP API port", diff --git a/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh b/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh index f0612b2db..1fdce47e4 100644 --- a/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh +++ b/src/cloudai/workloads/ai_dynamo/ai_dynamo.sh @@ -37,6 +37,8 @@ declare -A aiperf_args declare -A aiperf_config declare -A aiperf_accuracy_args declare -A aiperf_accuracy_config +declare -a DYNAMO_DECODE_PIDS=() +declare -a DYNAMO_PREFILL_PIDS=() lmcache_controller_cmd="" SHARED_NODE_DISAGG="false" @@ -45,6 +47,9 @@ declare -A dynamo_args dynamo_args["backend"]="vllm" dynamo_args["node-setup-cmd"]="" dynamo_args["ingress-cmd"]="python -m dynamo.frontend --router-mode kv" +dynamo_args["aiperf-phase-restart-services"]="False" +dynamo_args["aiperf-phase-setup-scope"]="all" +dynamo_args["aiperf-phase-setup-cmd-scope"]="frontend" dynamo_args["port"]=$((8080 + SLURM_JOBID % 100)) dynamo_args["endpoint"]="v1/chat/completions" dynamo_args["model"]="Qwen/Qwen3-0.6B" @@ -101,6 +106,11 @@ _csv_lists_overlap() { return 1 } +_truthy() { + local value="${1:-}" + [[ "${value,,}" == "true" || "${value}" == "1" || "${value,,}" == "yes" ]] +} + _gpus_per_node() { local n=$(echo "${CUDA_VISIBLE_DEVICES:-}" | tr ',' '\n' | grep -c . || true) [[ "$n" -gt 0 ]] && echo "$n" || echo "1" @@ -507,11 +517,11 @@ _total_workers_decode() { } _count_initialized_prefill() { - grep -i -l -E "${prefill_config["worker-initialized-regex"]}" "${RESULTS_DIR}"/dynamo_*prefill* 2>/dev/null | wc -l + grep -i -l -E "${prefill_config["worker-initialized-regex"]}" $(_worker_log_glob_for_role "prefill") 2>/dev/null | wc -l } _count_initialized_decode() { - grep -i -l -E "${decode_config["worker-initialized-regex"]}" "${RESULTS_DIR}"/dynamo_*decode* 2>/dev/null | wc -l + grep -i -l -E "${decode_config["worker-initialized-regex"]}" $(_worker_log_glob_for_role "decode") 2>/dev/null | wc -l } _expected_ready_prefill() { @@ -549,9 +559,22 @@ _gpu_list_for_worker_offset() { _log_file_for_worker() { local role="$1" local idx="$2" + if _aiperf_phase_restart_services_enabled && [[ -n "${DYNAMO_PHASE_GENERATION:-}" ]]; then + echo "${RESULTS_DIR}/dynamo_${role}_${SLURM_NODEID}_${idx}.r${DYNAMO_PHASE_GENERATION}.log" + return + fi echo "${RESULTS_DIR}/dynamo_${role}_${SLURM_NODEID}_${idx}.log" } +_worker_log_glob_for_role() { + local role="$1" + if _aiperf_phase_restart_services_enabled && [[ -n "${DYNAMO_PHASE_GENERATION:-}" ]]; then + echo "${RESULTS_DIR}/dynamo_${role}_"*"_"*".r${DYNAMO_PHASE_GENERATION}.log" + return + fi + echo "${RESULTS_DIR}/dynamo_"*"${role}"*"" +} + function log_node_role() { local node_name=$1 @@ -591,6 +614,10 @@ _is_aiperf_accuracy_enabled() { [[ -n "${aiperf_accuracy_config["--script"]:-}" ]] } +_aiperf_phase_restart_services_enabled() { + _truthy "${dynamo_args["aiperf-phase-restart-services"]:-False}" +} + _init_runtime_env() { if _is_vllm || _is_sglang; then export HF_HOME="${HUGGINGFACE_HOME}" @@ -800,9 +827,14 @@ validate_environment() { function wait_for_frontend_marker() { while [ ! -f "$DONE_MARKER" ]; do + handle_aiperf_phase_setup_requests exit_on_error - log "Waiting for frontend completion marker by polling $DONE_MARKER" - sleep 30 + if _aiperf_phase_restart_services_enabled; then + sleep 1 + else + log "Waiting for frontend completion marker by polling $DONE_MARKER" + sleep 30 + fi done log "Done marker found." @@ -849,7 +881,7 @@ function write_routerctl() export ROUTER_HEALTH_MODEL="${dynamo_args["model"]}" export ROUTER_PID_FILE="${RESULTS_DIR}/router.pid" export ROUTER_LOG_FILE="${RESULTS_DIR}/dynamo_ingress.log" - export ROUTER_START_TIMEOUT="${ROUTER_START_TIMEOUT:-120}" + export ROUTER_START_TIMEOUT="${ROUTER_START_TIMEOUT:-300}" export ROUTER_STOP_TIMEOUT="${ROUTER_STOP_TIMEOUT:-30}" cat > "${RESULTS_DIR}/routerctl.sh" <<'EOF' @@ -864,7 +896,7 @@ log() { echo "[$(date +%F\ %T) $(hostname)]: $*"; } : "${ROUTER_HEALTH_MODEL:?ROUTER_HEALTH_MODEL is not set}" : "${ROUTER_PID_FILE:?ROUTER_PID_FILE is not set}" : "${ROUTER_LOG_FILE:?ROUTER_LOG_FILE is not set}" -: "${ROUTER_START_TIMEOUT:=120}" +: "${ROUTER_START_TIMEOUT:=300}" : "${ROUTER_STOP_TIMEOUT:=30}" router_pid() { @@ -970,6 +1002,176 @@ function start_router() "${RESULTS_DIR}/routerctl.sh" start } +_stop_pid() { + local pid="$1" + local name="$2" + local timeout="${DYNAMO_PHASE_RESTART_STOP_TIMEOUT_SEC:-${DYNAMO_PHASE_STOP_TIMEOUT:-120}}" + if [[ -z "${pid}" ]] || ! kill -0 "${pid}" 2>/dev/null; then + return + fi + + log "Stopping ${name} pid=${pid}" + kill -TERM "${pid}" 2>/dev/null || true + + local deadline=$((SECONDS + timeout)) + while kill -0 "${pid}" 2>/dev/null; do + if (( SECONDS >= deadline )); then + log "WARN: ${name} pid=${pid} did not stop within ${timeout}s; sending SIGKILL" + kill -KILL "${pid}" 2>/dev/null || true + break + fi + sleep 1 + done + + wait "${pid}" 2>/dev/null || true +} + +_stop_pid_array() { + local name="$1" + shift + + local pid + for pid in "$@"; do + _stop_pid "${pid}" "${name}" + done +} + +_kill_residual_phase_processes() { + if _is_frontend_node; then + pkill -TERM -f "^python[0-9.]* -m dynamo.frontend" 2>/dev/null || true + if _has_connector "kvbm"; then + pkill -TERM -f "^cargo run" 2>/dev/null || true + pkill -TERM -f "^sample-registry" 2>/dev/null || true + fi + fi + + if _is_vllm && { _is_prefill_node || _is_decode_node; }; then + pkill -TERM -f "^python[0-9.]* -m dynamo.vllm" 2>/dev/null || true + sleep "${DYNAMO_PHASE_RESTART_GRACE_SEC:-5}" + pkill -KILL -f "^python[0-9.]* -m dynamo.vllm" 2>/dev/null || true + fi +} + +stop_phase_managed_dynamo_services() { + if _is_frontend_node && [[ -x "${RESULTS_DIR}/routerctl.sh" ]]; then + "${RESULTS_DIR}/routerctl.sh" stop || true + fi + + if _is_decode_node; then + _stop_pid_array "decode worker" "${DYNAMO_DECODE_PIDS[@]:-}" + DYNAMO_DECODE_PIDS=() + fi + + if _is_prefill_node; then + _stop_pid_array "prefill worker" "${DYNAMO_PREFILL_PIDS[@]:-}" + DYNAMO_PREFILL_PIDS=() + fi + + _kill_residual_phase_processes +} + +start_phase_managed_dynamo_services() { + local phase_index="$1" + local phase_name="$2" + + export DYNAMO_PHASE_GENERATION=$((phase_index + 1)) + log "Starting phase-managed Dynamo services for [${phase_name}] with generation ${DYNAMO_PHASE_GENERATION}" + + if _is_decode_node; then + launch_decode || return 1 + fi + + if _is_prefill_node; then + launch_prefill || return 1 + fi + + if _is_frontend_node; then + launch_ingress || return 1 + if _is_sglang_dsr1; then + launch_sgl_http_server || return 1 + fi + fi +} + +_wait_for_aiperf_phase_markers() { + local prefix="$1" + local suffix="$2" + local timeout="${AIPERF_PHASE_SETUP_TIMEOUT:-900}" + local deadline=$((SECONDS + timeout)) + local node + + while :; do + local missing="" + for node in $(echo "${DYNAMO_NODELIST}" | tr ',' ' '); do + if [[ ! -f "${prefix}_${node}.${suffix}" ]]; then + missing="${missing} ${node}" + fi + done + if [[ -z "${missing}" ]]; then + return 0 + fi + if (( SECONDS >= deadline )); then + mark_failed "Timed out waiting for AIPerf phase ${suffix} marker(s):${missing}" + return 1 + fi + sleep 1 + done +} + +_run_aiperf_phase_setup_cmd() { + local cmd_file="$1" + local cmd_scope="${dynamo_args["aiperf-phase-setup-cmd-scope"]:-frontend}" + + [[ -s "${cmd_file}" ]] || return 0 + if [[ "${cmd_scope}" == "all" ]] || { [[ "${cmd_scope}" == "frontend" ]] && _is_frontend_node; }; then + log "Running AIPerf phase setup command from ${cmd_file}" + bash -lc "$(cat "${cmd_file}")" + fi +} + +handle_aiperf_phase_setup_requests() { + _aiperf_phase_restart_services_enabled || return 0 + + local request + for request in "${RESULTS_DIR}"/aiperf_phase_setup_*.request; do + [[ -f "${request}" ]] || continue + + local prefix="${request%.request}" + local node_name="$(_current_node_name)" + local done_marker="${prefix}_${node_name}.done" + local stopped_marker="${prefix}_${node_name}.stopped" + [[ -f "${done_marker}" ]] && continue + + local phase_index="${prefix##*_}" + local phase_name + phase_name="$(cat "${prefix}.name" 2>/dev/null || echo "${phase_index}")" + + local setup_scope="${dynamo_args["aiperf-phase-setup-scope"]:-all}" + local participates=false + if [[ "${setup_scope}" == "all" ]] || { [[ "${setup_scope}" == "frontend" ]] && _is_frontend_node; }; then + participates=true + fi + + if [[ "${participates}" == "true" ]]; then + log "Stopping phase-managed Dynamo services for [${phase_name}]" + stop_phase_managed_dynamo_services + fi + touch "${stopped_marker}" + _wait_for_aiperf_phase_markers "${prefix}" "stopped" || return 1 + + if [[ "${participates}" == "true" ]]; then + _run_aiperf_phase_setup_cmd "${prefix}.cmd" || return 1 + if ! start_phase_managed_dynamo_services "${phase_index}" "${phase_name}"; then + mark_failed "Failed to start phase-managed Dynamo services for [${phase_name}]" + stop_phase_managed_dynamo_services + return 1 + fi + fi + touch "${done_marker}" + log "AIPerf phase setup completed for [${phase_name}]" + done +} + launch_sgl_http_server() { local script_path="${dynamo_args["repo"]}/components/backends/sglang/src/dynamo/sglang/utils/sgl_http_server.py" local port="${dynamo_args["sgl-http-port"]}" @@ -1031,6 +1233,9 @@ function launch_decode() ${decode_config["cmd"]} \ ${args_arr[@]} \ ${decode_config["extra-args"]} > $log_file 2>&1 & + local pid=$! + DYNAMO_DECODE_PIDS+=("${pid}") + log "Decode worker $i PID: ${pid}" done } @@ -1105,6 +1310,9 @@ function launch_prefill() ${prefill_config["cmd"]} \ ${args_arr[@]} \ ${prefill_config["extra-args"]} > $log_file 2>&1 & + local pid=$! + DYNAMO_PREFILL_PIDS+=("${pid}") + log "Prefill worker $i PID: ${pid}" done } @@ -1357,6 +1565,7 @@ function launch_workload() export AIPERF_ENDPOINT="${dynamo_args["endpoint"]}" export AIPERF_FAILURE_MARKER="${FATAL_ERROR_MARKER}" export AIPERF_SERVER_METRICS_URLS="$(_resolve_aiperf_server_metrics_urls)" + export AIPERF_PHASE_SETUP_PREFIX="${RESULTS_DIR}/aiperf_phase_setup" # Build config and workload args as proper bash arrays to preserve # multi-word values (e.g. --cmd "genai-perf profile") through word splitting. @@ -1397,7 +1606,15 @@ function launch_workload() function launch_workloads() { - wait_for_dynamo_frontend + if _aiperf_phase_restart_services_enabled; then + if _is_genai_perf_workload || _is_aiperf_accuracy_enabled; then + mark_failed "aiperf-phase-restart-services currently supports aiperf.sh-only runs" + return 1 + fi + log "AIPerf phase restart mode enabled: services will be started by each phase setup barrier" + else + wait_for_dynamo_frontend + fi if _is_genai_perf_workload; then launch_workload genai_perf_config genai_perf_args || return $? @@ -1450,27 +1667,34 @@ function main() # Workers launch BEFORE the ingress: launch_ingress blocks in # wait_for_router, and the router only becomes ready once a worker # registers — on a combined frontend+worker node the old order serialized - # the whole ROUTER_START_TIMEOUT (120 s of failing readiness curls) in + # the whole ROUTER_START_TIMEOUT of failing readiness curls in # front of every worker start. Workers only need etcd/nats (waited above) # and the lmcache config from setup_lmcache; they never talk to the router. - if _is_decode_node; then + local phase_restart_services=false + if _aiperf_phase_restart_services_enabled && _is_aiperf_workload; then + phase_restart_services=true + fi + + if [[ "${phase_restart_services}" != "true" ]] && _is_decode_node; then log "Node ID: $SLURM_NODEID, Role: decode" log_node_role "$(_current_node_name)" "decode" launch_decode & fi - if _is_prefill_node; then + if [[ "${phase_restart_services}" != "true" ]] && _is_prefill_node; then log "Node ID: $SLURM_NODEID, Role: prefill" log_node_role "$(_current_node_name)" "prefill" launch_prefill & fi if _is_frontend_node; then - launch_ingress - if _is_sglang_dsr1; then - launch_sgl_http_server + if [[ "${phase_restart_services}" != "true" ]]; then + launch_ingress || { mark_failed "Failed to start Dynamo ingress"; exit 1; } + if _is_sglang_dsr1; then + launch_sgl_http_server || { mark_failed "Failed to start SGL HTTP server"; exit 1; } + fi + sleep 10 fi - sleep 10 launch_workloads & fi 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..ecb7908f1 100644 --- a/src/cloudai/workloads/ai_dynamo/slurm_command_gen_strategy.py +++ b/src/cloudai/workloads/ai_dynamo/slurm_command_gen_strategy.py @@ -59,6 +59,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: @@ -258,6 +263,54 @@ def _render_aiperf_setup_blocks(self, log_message: str, setup_cmd: str | None) - ).rstrip() ] + def _render_aiperf_phase_restart_helpers(self) -> str: + return textwrap.dedent( + """\ + phase_expected_nodes() { + echo "${DYNAMO_NODELIST:?DYNAMO_NODELIST is not set}" | tr ',' ' ' + } + + wait_for_phase_markers() { + local prefix="$1" + local suffix="$2" + local timeout="${AIPERF_PHASE_SETUP_TIMEOUT:-900}" + local deadline=$((SECONDS + timeout)) + local missing="" + while :; do + missing="" + for node in $(phase_expected_nodes); do + if [[ ! -f "${prefix}_${node}.${suffix}" ]]; then + missing="${missing} ${node}" + fi + done + if [[ -z "${missing}" ]]; then + return 0 + fi + if (( SECONDS >= deadline )); then + log "FATAL: timed out waiting for AIPerf phase ${suffix} marker(s):${missing}" + return 1 + fi + sleep 1 + done + } + + request_aiperf_phase_setup() { + local phase_index="$1" + local phase_name="$2" + local setup_cmd="${3:-}" + local prefix="${AIPERF_PHASE_SETUP_PREFIX:-/cloudai_run_results/aiperf_phase_setup}_${phase_index}" + + rm -f "${prefix}.request" "${prefix}.name" "${prefix}.cmd" "${prefix}"_*.stopped "${prefix}"_*.done + printf '%s\\n' "${phase_name}" > "${prefix}.name" + printf '%s' "${setup_cmd}" > "${prefix}.cmd" + log "Requesting AIPerf phase setup for ${phase_name}" + touch "${prefix}.request" + wait_for_phase_markers "${prefix}" "done" + rm -f "${prefix}.request" + } + """ + ).rstrip() + def _render_between_aiperf_phases_block( self, phase_name: str, @@ -278,9 +331,25 @@ def _render_between_aiperf_phases_block( .splitlines() ) + def _render_aiperf_phase_setup_lines( + self, + phase_index: int, + phase: AIPerfPhase, + phase_restart_services: bool, + ) -> list[str]: + phase_setup = phase.setup_cmd if "setup_cmd" in phase.model_fields_set else None + if phase_restart_services: + return [ + "request_aiperf_phase_setup " + f"{phase_index} {shlex.quote(phase.name)} {shlex.quote(phase_setup or '')}" + ] + + return self._render_aiperf_setup_blocks(f"Running AIPerf phase setup for {phase.name}", phase_setup) + def _render_aiperf_script(self) -> str: phases = self.td.cmd_args.aiperf_phases or [AIPerfPhase.model_validate({"name": "aiperf"})] single_phase = len(phases) == 1 + phase_restart_services = self.td.cmd_args.dynamo.aiperf_phase_restart_services blocks = [ textwrap.dedent( f"""\ @@ -297,6 +366,9 @@ def _render_aiperf_script(self) -> str: ).rstrip() ] + if phase_restart_services: + blocks.append(self._render_aiperf_phase_restart_helpers()) + blocks.extend(self._render_aiperf_setup_blocks("Running aiperf setup", self.td.cmd_args.aiperf.setup_cmd)) write_phase_logs = not single_phase @@ -319,8 +391,7 @@ def _render_aiperf_script(self) -> str: else: run_cmd = cmd log_message = f"Running {phase.name}: {cmd}" - phase_setup = phase.setup_cmd if "setup_cmd" in phase.model_fields_set else None - phase_lines = self._render_aiperf_setup_blocks(f"Running AIPerf phase setup for {phase.name}", phase_setup) + phase_lines = self._render_aiperf_phase_setup_lines(idx, phase, phase_restart_services) phase_lines.append( textwrap.dedent( f"""\ 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..f119d356f 100644 --- a/tests/workloads/ai_dynamo/test_command_gen_strategy_slurm.py +++ b/tests/workloads/ai_dynamo/test_command_gen_strategy_slurm.py @@ -378,6 +378,39 @@ def test_generated_aiperf_script_supports_core_overrides_and_server_metrics_auto assert "--no-server-metrics" not in script +def test_aiperf_phase_restart_services_renders_barrier_and_dynamo_args( + strategy: AIDynamoSlurmCommandGenStrategy, +) -> None: + td = cast(AIDynamoTestDefinition, strategy.test_run.test) + td.cmd_args.workloads = "aiperf.sh" + td.cmd_args.dynamo.aiperf_phase_restart_services = True + td.cmd_args.dynamo.aiperf_phase_setup_scope = "all" + td.cmd_args.dynamo.aiperf_phase_setup_cmd_scope = "all" + td.cmd_args.aiperf_phases = [ + AIPerfPhase.model_validate({"name": "round_1", "args": {"concurrency": 1}}), + AIPerfPhase.model_validate( + { + "name": "round_2", + "setup-cmd": "rm -rf /tmp/lmcache/*", + "args": {"concurrency": 2}, + } + ), + ] + + result = strategy._gen_script_args(td) + + assert '--dynamo-aiperf-phase-restart-services "True"' in result + assert '--dynamo-aiperf-phase-setup-scope "all"' in result + assert '--dynamo-aiperf-phase-setup-cmd-scope "all"' in result + assert strategy.final_env_vars["DYNAMO_NODELIST"] == "n0,n1" + + script = (strategy.test_run.output_path / "aiperf.sh").read_text() + assert "request_aiperf_phase_setup()" in script + assert "request_aiperf_phase_setup 0 round_1 ''" in script + assert "request_aiperf_phase_setup 1 round_2 'rm -rf /tmp/lmcache/*'" in script + assert "Running AIPerf phase setup for round_2" not in script + + def test_generated_aiperf_script_rejects_list_args(strategy: AIDynamoSlurmCommandGenStrategy) -> None: td = cast(AIDynamoTestDefinition, strategy.test_run.test) td.cmd_args.workloads = "aiperf.sh"