Pin what a typical run renders, and the commands one launch runs (#2460)

This commit is contained in:
fzyzcjy
2026-09-26 18:37:21 +08:00
committed by GitHub
parent 353570f894
commit defbe3d80c
10 changed files with 2637 additions and 10 deletions
+1
View File
@@ -15,6 +15,7 @@ repos:
rev: v4.5.0
hooks:
- id: check-yaml
args: ['--allow-multiple-documents']
exclude: '^charts/.*/templates/'
- id: check-case-conflict
- id: detect-private-key
@@ -0,0 +1,303 @@
import os
import re
import shlex
import subprocess
import sys
from collections.abc import Iterator
from pathlib import Path
from typing import Any
import pytest
import yaml
from megatron.training import arguments as megatron_arguments
from tests.fast.charts.utils import NAMESPACE, RUN_CHART_DIR, RUN_ID, RUN_RELEASE_NAME, requires_helm
from tests.fast.launch_scripts.sh_harness import REPO_ROOT, SANDBOX_PLACEHOLDER, assert_matches_snapshot
from miles.ray.specs.entrypoint import compute_specs
from miles.utils.arguments import parse_args
from miles.utils.external_utils.command_utils.common import rsync_cmd
from miles.utils.external_utils.command_utils.helm_backend.launcher.values.builder import build_values
from miles.utils.external_utils.command_utils.helm_backend.launcher.values.misc import LaunchPlan
from miles.utils.external_utils.model_args_utils import load_model_args
from miles.utils.workers.serving.utils import override_argv
SNAPSHOT_DIR = REPO_ROOT / "tests" / "snapshots" / "charts" / "miles-run"
FIXTURE_DIR = Path(__file__).resolve().parent
INFRA_VALUES = FIXTURE_DIR / "typical-infra.yaml"
SGLANG_CONFIG = FIXTURE_DIR / "typical-sglang.yaml"
HF_CHECKPOINT = FIXTURE_DIR / "typical-model"
PYTHON_PLACEHOLDER = "<PYTHON>"
FIXTURE_DIR_PLACEHOLDER = "<FIXTURE_DIR>"
RANDOM_SEED_PLACEHOLDER = "<RANDOM_SEED>"
RANDOM_SEED_FLAG = "--random-seed"
SCENARIOS = ("typical-values", "typical")
MODEL_TYPE = "qwen3-4B"
ROTARY_BASE = "1000000"
PREFILL_GPUS = 16
DECODE_GPUS = 16
GPUS_PER_NODE = 8
TRAINER_NODES = 4
ORCHESTRATOR_COMMAND = ["python", "scripts/run_qwen3_4b.py", "train", "--cluster-backend", "kubernetes"]
WORKER_ARGV = ["--cluster-backend", "kubernetes", "--rollout-num-gpus", str(PREFILL_GPUS + DECODE_GPUS)]
PREPARE_CMD = {"trainer": rsync_cmd("/cluster-storage/models/Qwen3-4B", "/scratch/Qwen3-4B")}
PARSER_ENV = {"CUDA_DEVICE_MAX_CONNECTIONS": "1"}
SCENARIO_ARGV = [
*shlex.split(load_model_args(MODEL_TYPE, rotary_base=ROTARY_BASE)),
# named rather than left to sglang's own probe, which reads the launcher's accelerator and
# has none to read on the cpu lane that renders this snapshot
"--sglang-device",
"cuda",
"--hf-checkpoint",
str(HF_CHECKPOINT),
"--load",
"/cluster-storage/models/Qwen3-4B_torch_dist",
"--save",
"/cluster-storage/myteam/miles_data/miles-runs/myrun/checkpoints",
"--save-interval",
"20",
"--prompt-data",
"/cluster-storage/datasets/dapo-math-17k/dapo-math-17k.jsonl",
"--input-key",
"prompt",
"--label-key",
"label",
"--apply-chat-template",
"--rollout-shuffle",
"--rm-type",
"math",
"--num-rollout",
"16",
"--rollout-batch-size",
"32",
"--n-samples-per-prompt",
"8",
"--rollout-max-response-len",
"8192",
"--rollout-temperature",
"1",
"--global-batch-size",
"256",
"--balance-data",
"--optimizer",
"adam",
"--lr",
"1e-6",
"--lr-decay-style",
"constant",
"--weight-decay",
"0.1",
"--adam-beta1",
"0.9",
"--adam-beta2",
"0.98",
"--advantage-estimator",
"grpo",
"--eps-clip",
"0.2",
"--eps-clip-high",
"0.28",
"--use-dynamic-batch-size",
"--max-tokens-per-gpu",
"9216",
"--sglang-config",
str(SGLANG_CONFIG),
"--rollout-num-gpus",
str(PREFILL_GPUS + DECODE_GPUS),
"--sglang-chunked-prefill-size",
"4096",
"--sglang-mem-fraction-static",
"0.7",
"--tensor-model-parallel-size",
"2",
"--sequence-parallel",
"--pipeline-model-parallel-size",
"1",
"--context-parallel-size",
"4",
"--cp-comm-type",
"a2a",
"--expert-model-parallel-size",
"1",
"--expert-tensor-parallel-size",
"1",
"--recompute-granularity",
"full",
"--recompute-method",
"uniform",
"--recompute-num-layers",
"1",
"--attention-dropout",
"0.0",
"--hidden-dropout",
"0.0",
"--accumulate-allreduce-grads-in-fp32",
"--attention-softmax-in-fp32",
"--attention-backend",
"flash",
"--actor-num-nodes",
str(TRAINER_NODES),
"--actor-num-gpus-per-node",
str(GPUS_PER_NODE),
"--num-gpus-per-node",
str(GPUS_PER_NODE),
"--colocate",
"--use-session-server",
"--cluster-backend",
"kubernetes",
"--run-uuid",
"0123456789abcdef",
]
@pytest.fixture(autouse=True)
def parser_process_env(monkeypatch: pytest.MonkeyPatch) -> Iterator[None]:
for name, value in PARSER_ENV.items():
monkeypatch.setenv(name, value)
tuning_env_name = "SGLANG_OPT_USE_CUSTOM_ALL_REDUCE_V2"
tuning_env_value = os.environ.pop(tuning_env_name, None)
# megatron validates the arguments against the local accelerator, which the cpu lane that
# runs this snapshot has none of; the rendered values do not depend on the answer
monkeypatch.setattr(megatron_arguments, "get_device_arch_version", lambda: _DEVICE_ARCH_VERSION)
yield
if tuning_env_value is None:
os.environ.pop(tuning_env_name, None)
else:
os.environ[tuning_env_name] = tuning_env_value
_DEVICE_ARCH_VERSION = 9
def _dump_values(values: dict[str, Any]) -> str:
# unwrapped: yaml folds long lines against the real interpreter path, which the snapshot
# only replaces afterwards, so a machine whose path is a different length folds elsewhere
return yaml.safe_dump(values, default_flow_style=False, sort_keys=True, width=_NO_WRAP)
_NO_WRAP = 1 << 30
def synthetic_specs() -> list[Any]:
with override_argv(SCENARIO_ARGV):
return compute_specs(parse_args())
def synthetic_run_values() -> dict[str, Any]:
return build_values(
synthetic_specs(),
LaunchPlan(
run_id=RUN_ID,
release=RUN_RELEASE_NAME,
namespace="rl",
state_file=f"/cluster-storage/myteam/miles_data/miles-runs/{RUN_ID}/state/orchestrator-260101-000000-000001.state",
orchestrator_command=ORCHESTRATOR_COMMAND,
worker_argv=WORKER_ARGV,
env={"PYTHONUNBUFFERED": "1", **PARSER_ENV},
colocate=True,
prepare_cmd=PREPARE_CMD,
),
).as_values()
def render_from(values_file: Path) -> str:
result = subprocess.run(
[
"helm",
"template",
RUN_RELEASE_NAME,
str(RUN_CHART_DIR),
"-n",
NAMESPACE,
"-f",
str(INFRA_VALUES),
"-f",
str(values_file),
],
capture_output=True,
text=True,
)
assert result.returncode == 0, result.stderr
return result.stdout
def freeze(text: str, sandbox: Path) -> str:
# only this test's own fixtures really move from machine to machine. the chart mounts the repo at
# a fixed container path, so masking the whole checkout instead would mask that constant away on
# any machine that happens to be checked out there, and record a snapshot only it can reproduce
masked = text.replace(str(sandbox), SANDBOX_PLACEHOLDER).replace(str(FIXTURE_DIR), FIXTURE_DIR_PLACEHOLDER)
return mask_random_seeds(masked.replace(sys.executable, PYTHON_PLACEHOLDER))
def mask_random_seeds(text: str) -> str:
lines = text.split("\n")
masked = [
(
re.sub(r"\d+", RANDOM_SEED_PLACEHOLDER, line)
if index and _yaml_scalar(lines[index - 1]) == RANDOM_SEED_FLAG
else line
)
for index, line in enumerate(lines)
]
return "\n".join(masked)
def _yaml_scalar(line: str) -> str:
return line.strip().removeprefix("- ").strip("'\"")
@requires_helm
class TestGeneratedValuesSnapshot:
def test_the_launcher_turns_the_specs_into_exactly_the_recorded_values(self, tmp_path):
"""The spec to values transform decides a run's whole shape, so it is pinned end to end."""
values = _dump_values(synthetic_run_values())
assert_matches_snapshot(
SNAPSHOT_DIR / "typical-values.yaml", freeze(values, sandbox=tmp_path), "miles-run generated values"
)
def test_those_values_render_exactly_the_recorded_manifests(self, tmp_path):
"""Rendering the file the launcher really writes is what pins the two halves to each other."""
values_file = tmp_path / "run-values.yaml"
values_file.write_text(_dump_values(synthetic_run_values()))
assert_matches_snapshot(
SNAPSHOT_DIR / "typical.yaml", freeze(render_from(values_file), sandbox=tmp_path), "miles-run manifests"
)
class TestRandomSeedMasking:
def test_a_seed_the_engine_drew_becomes_a_placeholder_in_the_generated_values(self):
"""sglang draws a fresh seed per render, so the values a run generates cannot record the number."""
values = " - --tp-size\n - '8'\n - --random-seed\n - '379064976'\n - --enable-metrics\n"
assert mask_random_seeds(values) == (
" - --tp-size\n - '8'\n - --random-seed\n - '<RANDOM_SEED>'\n - --enable-metrics\n"
)
def test_a_seed_the_engine_drew_becomes_a_placeholder_in_the_rendered_manifests(self):
"""The manifests quote their argv differently from the values, and must be masked all the same."""
manifests = ' - "--random-seed"\n - "723999131"\n'
assert mask_random_seeds(manifests) == (' - "--random-seed"\n - "<RANDOM_SEED>"\n')
def test_a_number_that_no_seed_flag_introduces_is_left_alone(self):
"""Masking every number would hide the real argv, so only the seed's own value is replaced."""
argv = " - --tp-size\n - '8'\n"
assert mask_random_seeds(argv) == argv
class TestSnapshotFiles:
def test_the_recorded_files_are_exactly_the_declared_ones(self):
"""A renamed or deleted scenario must not leave an orphan golden, nor a new one go unrecorded."""
recorded = {path.stem for path in SNAPSHOT_DIR.glob("*.yaml")}
assert recorded == set(SCENARIOS)
@@ -0,0 +1,30 @@
infra:
image:
repository: myregistry.example/myteam/miles
tag: v0.9.3-cu128
pullPolicy: IfNotPresent
pullSecrets:
- myregistry-pull-secret
sharedStorage:
type: pvc
pvcClaimName: miles-shared-nvme
mountPath: /cluster-storage
paths:
runsSubPath: myteam/miles_data
repos:
miles: myuser/miles
megatron: myuser/Megatron-LM
nodeLocalStorage:
hostPath: /mnt/local-nvme/miles-staging
mountPath: /scratch
scheduling:
nodeSelector:
nvidia.com/gpu.product: NVIDIA-H100-80GB-HBM3
tolerations:
- key: nvidia.com/gpu
operator: Exists
effect: NoSchedule
env:
NCCL_SOCKET_IFNAME: bond0
NCCL_IB_HCA: mlx5_0,mlx5_1
HF_ENDPOINT: https://hf-mirror.example
@@ -0,0 +1,19 @@
{
"architectures": [
"Qwen3ForCausalLM"
],
"model_type": "qwen3",
"hidden_size": 2560,
"intermediate_size": 9728,
"num_hidden_layers": 36,
"num_attention_heads": 32,
"num_key_value_heads": 8,
"head_dim": 128,
"hidden_act": "silu",
"rms_norm_eps": 1e-06,
"rope_theta": 1000000.0,
"tie_word_embeddings": true,
"vocab_size": 151936,
"max_position_embeddings": 40960,
"torch_dtype": "bfloat16"
}
@@ -0,0 +1,9 @@
sglang:
- name: default
server_groups:
- worker_type: prefill
num_gpus: 16
num_gpus_per_engine: 16
- worker_type: decode
num_gpus: 16
num_gpus_per_engine: 8
@@ -1,11 +1,17 @@
import json
import subprocess
from pathlib import Path
import pytest
from miles.utils.external_utils.command_utils.common import chart_dir
from miles.utils.external_utils.command_utils.helm_backend.launcher import command_wrapper
from miles.utils.external_utils.command_utils.helm_backend.launcher.command_wrapper import Helm, Kubectl
from miles.utils.external_utils.command_utils.helm_backend.launcher import command_wrapper, entrypoint
from miles.utils.external_utils.command_utils.helm_backend.launcher.command_wrapper import (
_JOB_COMPLETION_JSONPATH,
Helm,
Kubectl,
)
from miles.utils.external_utils.command_utils.helm_backend.naming import RunNames
from miles.utils.workers.k8s_types import Pod
from miles.utils.workers.worker_provider.kubernetes.helm import naming
from miles.utils.workers.worker_provider.kubernetes.helm.env import INSTANCE_LABEL
@@ -43,21 +49,56 @@ class TestLogCommands:
"""An engine and a trainer share only this label, so anything narrower misses half the run."""
assert Kubectl.release_selector("miles-run-x") == f"{INSTANCE_LABEL}=miles-run-x"
def test_selects_a_job_by_the_label_the_job_controller_stamps_on_its_pods(self):
"""A command job's pods are named after their job, and only this label survives a pod restart."""
assert Kubectl.job_selector("miles-run-command-gpu") == "batch.kubernetes.io/job-name=miles-run-command-gpu"
class TestUpgradeCommand:
def test_installs_a_missing_release_and_updates_an_existing_one(self):
"""Relaunching a run id must update it in place, which plain upgrade would refuse to do."""
command = Helm.upgrade_command("r", "rl", "/c", [], ci_run=False)
assert command[:4] == ["helm", "upgrade", "--install", "r"]
def test_keeps_the_user_values_ahead_of_the_computed_ones(self):
"""A run value must win over a cluster default, and helm lets the later file win."""
command = Helm.upgrade_command("r", "rl", "/c", ["/infra.yaml", "/run.yaml"], ci_run=False)
assert command[command.index("/infra.yaml") - 1] == "--values"
assert command.index("/infra.yaml") < command.index("/run.yaml")
def test_installs_a_missing_release_and_updates_an_existing_one_without_ci_run_argument(self):
"""Relaunching a run id must update it in place, which plain upgrade would refuse to do."""
command = Helm.upgrade_command("r", "myns", "/c", [])
command = Helm.upgrade_command("r", "rl", "/c", [], ci_run=False)
assert command[:4] == ["helm", "upgrade", "--install", "r"]
def test_keeps_the_user_values_ahead_of_the_computed_ones_without_ci_run_argument(self):
"""A run value must win over a cluster default, and helm lets the later file win."""
command = Helm.upgrade_command("r", "myns", "/c", ["/infra.yaml", "/run.yaml"])
command = Helm.upgrade_command("r", "rl", "/c", ["/infra.yaml", "/run.yaml"], ci_run=False)
assert command[command.index("/infra.yaml") - 1] == "--values"
assert command.index("/infra.yaml") < command.index("/run.yaml")
def test_labels_a_ci_release_so_the_next_run_can_clean_it_up(self):
"""The cleanup selects on this label, and an unlabelled CI release is one nothing will ever remove."""
command = Helm.upgrade_command("r", "rl", "/c", [], ci_run=True)
assert command[command.index("--labels") + 1] == f"{command_wrapper.CI_LABEL}=true"
def test_leaves_a_human_release_unlabelled(self):
"""A developer's run carrying the CI label would be uninstalled by the next CI job in that namespace."""
command = Helm.upgrade_command("r", "rl", "/c", [], ci_run=False)
assert "--labels" not in command
def test_labels_the_release_rather_than_its_objects(self):
"""helm --labels records release metadata; a values-level label would not be selectable by helm list."""
command = Helm.upgrade_command("r", "rl", "/c", [], ci_run=True)
assert command.index("--labels") > command.index("--namespace")
class TestBuildDependencies:
def test_chart_dependencies_are_rebuilt_only_when_a_locked_dependency_is_missing(
@@ -167,13 +208,156 @@ def _kubectl_answering(monkeypatch, *, returncode: int, stdout: str = "", stderr
return commands
class TestGetJson:
def test_a_failed_get_is_not_reported_as_an_absent_object(self, monkeypatch: pytest.MonkeyPatch) -> None:
"""A failed lookup must expose its exit code and stderr instead of looking like an absent object."""
_kubectl_answering(monkeypatch, returncode=23, stderr="the api server refused the request")
class TestRequestBounds:
def test_reads_and_creates_bound_the_subprocess_instead_of_passing_a_kubectl_override(self, monkeypatch):
"""A `--request-timeout` override stops kubectl from falling back to the in-cluster config a workbench relies on."""
calls: list[tuple[list[str], float | None]] = []
with pytest.raises(RuntimeError, match="code 23: the api server refused the request"):
Kubectl.get_json("pod", return_type=Pod, name="trainer-0", namespace="rl")
def fake_run(argv: list[str], **kwargs) -> subprocess.CompletedProcess:
calls.append((argv[1:], kwargs.get("timeout")))
return subprocess.CompletedProcess(args=argv, returncode=0, stdout="", stderr="")
monkeypatch.setattr(command_wrapper, "run_process", fake_run)
Kubectl.get_json("pods", return_type=Pod, namespace="training")
Kubectl.create_if_absent("/etc/miles/job.yaml")
assert all("--request-timeout" not in argv for argv, _ in calls)
assert [timeout for _, timeout in calls] == [30.0, 60.0]
class TestCreateIfAbsent:
def test_creates_the_objects_of_a_rendered_manifest(self, monkeypatch):
"""kubectl apply would adopt an object helm owns; create is what refuses to touch one."""
commands = _kubectl_answering(monkeypatch, returncode=0, stderr="")
assert Kubectl.create_if_absent("/etc/miles/job.yaml")
assert commands == [["create", "-f", "/etc/miles/job.yaml"]]
def test_reports_an_object_that_was_already_there_without_failing(self, monkeypatch):
"""Its callers retry after a restart, and the whole point is that the second attempt is harmless."""
_kubectl_answering(
monkeypatch, returncode=1, stderr='Error from server (AlreadyExists): jobs "u" already exists'
)
assert not Kubectl.create_if_absent("/etc/miles/job.yaml")
def test_refuses_to_read_any_other_failure_as_idempotence(self, monkeypatch):
"""A forbidden create means the object is missing, and pretending otherwise loses it silently."""
_kubectl_answering(monkeypatch, returncode=1, stderr="Error from server (Forbidden): jobs is forbidden")
with pytest.raises(RuntimeError, match="Could not create"):
Kubectl.create_if_absent("/etc/miles/job.yaml")
class TestJobsFinished:
def test_reports_a_job_that_already_succeeded_or_failed(self, monkeypatch):
"""A job holding its name after it ran never runs again, so its caller has to recreate it."""
commands = _kubectl_answering(monkeypatch, returncode=0, stdout="1,")
assert Kubectl.jobs_finished("/etc/miles/job.yaml")
assert commands == [["get", "-f", "/etc/miles/job.yaml", "--output", _JOB_COMPLETION_JSONPATH]]
def test_reports_a_job_that_has_not_run_yet_as_unfinished(self, monkeypatch):
"""An armed job still does its work, and replacing it would throw that work away."""
_kubectl_answering(monkeypatch, returncode=0, stdout=",")
assert not Kubectl.jobs_finished("/etc/miles/job.yaml")
def test_refuses_to_read_an_unanswered_query_as_unfinished(self, monkeypatch):
"""Guessing here adopts a job that already ran and leaves the release installed forever."""
_kubectl_answering(monkeypatch, returncode=1, stderr="Error from server: etcdserver: request timed out")
with pytest.raises(RuntimeError, match="Could not read"):
Kubectl.jobs_finished("/etc/miles/job.yaml")
class TestReplace:
def test_forces_the_objects_of_a_rendered_manifest_over_the_ones_holding_their_names(self, monkeypatch):
"""A job's pod template is immutable, so only a forced replace can put a fresh job under that name."""
commands = _kubectl_answering(monkeypatch, returncode=0)
Kubectl.replace("/etc/miles/job.yaml")
assert commands == [["replace", "--force", "-f", "/etc/miles/job.yaml"]]
def test_lets_a_caller_that_cannot_go_on_without_the_replacement_fail(self, monkeypatch):
"""Reporting a replace that never happened would tell the run that it uninstalls itself when it does not."""
_kubectl_answering(monkeypatch, returncode=1, stderr="Error from server (Forbidden): jobs is forbidden")
with pytest.raises(RuntimeError, match="Could not replace"):
Kubectl.replace("/etc/miles/job.yaml")
class TestDeleteJob:
def test_treats_a_job_that_is_not_there_as_deleted(self, monkeypatch):
"""The launcher deletes a job that usually does not exist, which is the outcome it wants anyway."""
commands = _kubectl_answering(monkeypatch, returncode=0, stderr="")
Kubectl.delete_job("miles-run-x-uninstall", namespace="rl", check=True)
assert "--ignore-not-found" in commands[0]
def test_waits_for_the_pods_of_the_job_it_deletes(self, monkeypatch):
"""The launcher installs over the deleted job, and a pod that outlives it uninstalls the new release."""
commands = _kubectl_answering(monkeypatch, returncode=0, stderr="")
Kubectl.delete_job("miles-run-x-uninstall", namespace="rl", check=True)
assert commands[0][commands[0].index("--cascade") + 1] == "foreground"
def test_lets_a_caller_that_cannot_go_on_without_the_deletion_fail(self, monkeypatch):
"""Installing over a job that is still armed hands the new release to the old run's uninstall."""
_kubectl_answering(monkeypatch, returncode=1, stderr="the api server refused")
with pytest.raises(subprocess.CalledProcessError):
Kubectl.delete_job("miles-run-x-uninstall", namespace="rl", check=True)
def test_stays_tolerant_for_the_cleanup_of_a_command_job(self, monkeypatch):
"""That caller deletes the same job twice around a run, and neither call is worth failing over."""
_kubectl_answering(monkeypatch, returncode=1, stderr="the api server refused")
Kubectl.delete_job("miles-run-command-convert", namespace="rl")
LAUNCHING_RUN_ID = "260101-000000-000"
def _recorded_ci_cleanup(
monkeypatch: pytest.MonkeyPatch, namespace: str, *, listed: list[dict] | None = None
) -> list[list[str]]:
commands: list[list[str]] = []
def fake_run(command: list[str], capture_output: bool) -> subprocess.CompletedProcess:
commands.append(command)
return subprocess.CompletedProcess(args=command, returncode=0, stdout=json.dumps(listed or []), stderr="")
monkeypatch.setattr(command_wrapper, "_run", fake_run)
entrypoint._uninstall_leftover_ci_releases(namespace, keep_run_id=LAUNCHING_RUN_ID)
return commands
class TestCiCleanup:
def test_narrows_the_search_by_both_namespace_and_label(self, monkeypatch):
"""Deleting another user's run would kill a live experiment, so neither filter may be dropped."""
command = _recorded_ci_cleanup(monkeypatch, "ci-runner-3")[0]
assert command[command.index("--namespace") + 1] == "ci-runner-3"
assert command[command.index("--selector") + 1] == f"{command_wrapper.CI_LABEL}=true"
def test_reads_the_release_names_helm_reports(self):
"""The names drive uninstall, so a parse that silently returns nothing would leave releases behind."""
output = json.dumps([{"name": "miles-run-a", "namespace": "ci"}, {"name": "miles-run-b"}])
assert [release["name"] for release in json.loads(output or "[]")] == ["miles-run-a", "miles-run-b"]
def test_treats_no_output_as_nothing_to_clean(self):
"""helm prints nothing when no release matches, and that is not an error."""
assert json.loads("" or "[]") == []
def test_uninstalls_inside_the_namespace_it_was_told(self, monkeypatch):
"""A release name exists per namespace, so a missing namespace could hit a different one."""
commands = _recorded_ci_cleanup(monkeypatch, "ci", listed=[{"name": "miles-run-a"}])
assert commands[1] == ["helm", "uninstall", "miles-run-a", "--namespace", "ci"]
class TestChartDir:
@@ -182,6 +366,19 @@ class TestChartDir:
assert chart_dir(repo_base_dir="/repo").as_posix() == "/repo/charts/miles-run"
LONGEST_RUN_ID = "a" * 40
class TestReleaseName:
def test_a_release_is_the_chart_name_and_the_run_id(self):
"""The launcher finds a run's release again from the run id alone, so the rule is fixed."""
assert RunNames.release(run_id="260101-000000-000") == "miles-run-260101-000000-000"
def test_the_same_run_id_always_names_the_same_release(self):
"""Relaunching a run upgrades its release; a fresh name would deploy a second copy instead."""
assert RunNames.release(run_id=LONGEST_RUN_ID) == RunNames.release(run_id=LONGEST_RUN_ID)
class TestComponentName:
def test_an_object_is_the_release_the_chart_name_and_the_component(self):
"""Every object of a run is traceable to the release that made it."""
@@ -239,6 +436,15 @@ class TestComponentName:
)
class TestGetJson:
def test_a_failed_get_is_not_reported_as_an_absent_object(self, monkeypatch: pytest.MonkeyPatch) -> None:
"""A failed lookup must expose its exit code and stderr instead of looking like an absent object."""
_kubectl_answering(monkeypatch, returncode=23, stderr="the api server refused the request")
with pytest.raises(RuntimeError, match="code 23: the api server refused the request"):
Kubectl.get_json("pod", return_type=Pod, name="trainer-0", namespace="rl")
class TestStaticWorkerHost:
def test_a_static_cell_is_reached_through_its_own_pod_of_the_headless_service(self):
"""A pool of session servers is several addresses, and pod zero can answer only one of them."""
@@ -0,0 +1,203 @@
import subprocess
import sys
from pathlib import Path
from types import SimpleNamespace
from typing import Any
import yaml
from tests.fast.launch_scripts.sh_harness import REPO_ROOT, assert_matches_snapshot, sanitize
from miles.ray.specs.inference import POOL_CATEGORY_INFERENCE_ENGINE
from miles.ray.specs.train import POOL_CATEGORY_TRAINER_ENGINE
from miles.utils.external_utils.command_utils.base_backend import ExecuteTrainConfig, ExecuteTrainRequest
from miles.utils.external_utils.command_utils.helm_backend import naming
from miles.utils.external_utils.command_utils.helm_backend.launcher import command_wrapper, entrypoint
from miles.utils.external_utils.command_utils.helm_backend.launcher.command_wrapper import Helm
from miles.utils.external_utils.command_utils.helm_backend.launcher.values.misc import MooncakeInfo
from miles.utils.workers.worker_spec import CommandWorkerSpec, PortInfo, SchedulingSpec, ServeWorkerSpec
SNAPSHOT_DIR = REPO_ROOT / "tests" / "snapshots" / "helm_backend"
FROZEN_RUN_ID = "260101-000000-000"
FROZEN_LAUNCH_TOKEN = "260101-000000-000001"
NAMESPACE = "rl"
PYTHON_PLACEHOLDER = "<PYTHON>"
def _router() -> CommandWorkerSpec:
return CommandWorkerSpec(
name="inference-router-0",
port_infos=[PortInfo(name="primary", static_port=30000)],
env_var=lambda ctx: {},
scheduling=SchedulingSpec.single(num_gpus_per_worker=0),
launch_command=lambda ctx: (
f"python -m sglang_router.launch_router --host {ctx.self_addrs['primary'].host} --port 30000"
),
)
def _engine() -> CommandWorkerSpec:
return CommandWorkerSpec(
name="inference-engine-0-0",
category=POOL_CATEGORY_INFERENCE_ENGINE,
port_infos=[
PortInfo(name="primary", static_port=8000),
PortInfo(name="dist_init", static_port=9000, mode="master"),
],
env_var=lambda ctx: {"NVSHMEM_DISABLE_NCCL": "1"},
scheduling=SchedulingSpec(
num_cells=2,
num_workers_per_cell=2,
num_gpus_per_worker=0.2,
num_gpu_slots_per_worker=8,
num_gpus_per_node=8,
),
launch_command=lambda ctx: (
f"python -m sglang.launch_server --node-rank {ctx.worker_in_cell_index} "
f"--dist-init-addr {ctx.self_addrs['dist_init'].host}:{ctx.self_addrs['dist_init'].port}"
),
)
def _trainer() -> ServeWorkerSpec:
return ServeWorkerSpec(
name="trainer-engine-actor",
category=POOL_CATEGORY_TRAINER_ENGINE,
port_infos=[PortInfo(name="master", static_port=9000, mode="master")],
env_var=lambda ctx: {"NCCL_CUMEM_ENABLE": "0"},
scheduling=SchedulingSpec(
num_cells=2,
num_workers_per_cell=8,
num_gpus_per_worker=0.4,
num_gpu_slots_per_worker=1,
num_gpus_per_node=8,
),
worker_class="miles.backends.megatron_utils.actor.MegatronTrainRayActor",
ctor_kwargs=lambda ctx: {},
)
def _request() -> ExecuteTrainRequest:
return ExecuteTrainRequest(
train_args="--rollout-num-gpus 8",
num_gpus_per_node=8,
megatron_model_type=None,
train_script="/repo/train.py",
train_backend_fsdp=False,
extra_env_vars={},
megatron_path="/root/Megatron-LM",
before_ray_job_submit=None,
prepare_cmd={},
)
def helm_values_file(sandbox: Path) -> Path:
values_file = sandbox / "infra.yaml"
values_file.write_text(
yaml.safe_dump(
{
"infra": {
"image": {"repository": "myregistry.example/miles", "tag": "v1"},
"sharedStorage": {
"type": "hostPath",
"hostPath": f"{sandbox}/cluster-storage",
"mountPath": f"{sandbox}/cluster-storage",
},
"paths": {"runsSubPath": "miles_data"},
}
}
)
)
return values_file
def record_launch(monkeypatch, sandbox: Path) -> list[str]:
recorded: list[str] = []
def fake_run(command: list[str], **kwargs: Any) -> Any:
recorded.append(" ".join(str(part) for part in command))
return subprocess.CompletedProcess(args=command, returncode=0, stdout="", stderr="")
monkeypatch.setattr(command_wrapper, "run_process", fake_run)
monkeypatch.setattr(Helm, "get_manifest", staticmethod(lambda release, namespace: None))
monkeypatch.setattr(entrypoint, "repo_base_dir", str(REPO_ROOT))
monkeypatch.setattr(naming, "_new_launch_token", lambda: FROZEN_LAUNCH_TOKEN)
_stub_launch_inputs(monkeypatch, specs=[_router(), _engine(), _trainer()])
entrypoint.execute_train(
request=_request(),
config=ExecuteTrainConfig(
namespace=NAMESPACE, run_id=FROZEN_RUN_ID, helm_values=(str(helm_values_file(sandbox)),)
),
)
return recorded
def freeze(text: str, sandbox: Path) -> str:
return sanitize(text.replace(sys.executable, PYTHON_PLACEHOLDER), sandbox=sandbox)
def format_launch(commands: list[str], values_text: str, sandbox: Path) -> str:
lines: list[str] = []
for index, command in enumerate(commands):
lines.append(f"### {index}")
lines.append(freeze(command, sandbox=sandbox))
lines.append("")
lines.append("### pseudo file 1")
lines.append(freeze(values_text, sandbox=sandbox))
lines.append("")
return "\n".join(lines)
def _stub_launch_inputs(monkeypatch, *, specs, colocate: bool = False) -> None:
monkeypatch.setattr(entrypoint, "compute_specs", lambda args: specs)
monkeypatch.setattr(
entrypoint,
"parse_args",
lambda: SimpleNamespace(colocate=colocate, argv=[], use_wandb=False, wandb_run_id=None),
)
monkeypatch.setattr(MooncakeInfo, "plan_of_args", staticmethod(lambda args: None))
monkeypatch.setattr(entrypoint, "_follow_until_finished", lambda **kwargs: None)
class TestKubernetesLaunchSnapshot:
def test_the_helm_argv_and_the_generated_values_match_the_recording(self, monkeypatch, tmp_path):
"""The values file is the whole training recipe, so a snapshot of only the argv would prove little."""
commands = record_launch(monkeypatch, tmp_path)
values_file = (
tmp_path
/ "cluster-storage"
/ "miles_data"
/ "miles-runs"
/ FROZEN_RUN_ID
/ "values"
/ f"values-{FROZEN_LAUNCH_TOKEN}.yaml"
)
recorded = format_launch(commands, values_file.read_text(), tmp_path)
assert_matches_snapshot(SNAPSHOT_DIR / "kubernetes_launch.txt", recorded, "kubernetes launcher recording")
def test_the_generated_values_carry_no_infra_section(self, monkeypatch, tmp_path):
"""infra is the user's half of the contract; the launcher writing it would silently override a cluster."""
record_launch(monkeypatch, tmp_path)
values_file = (
tmp_path
/ "cluster-storage"
/ "miles_data"
/ "miles-runs"
/ FROZEN_RUN_ID
/ "values"
/ f"values-{FROZEN_LAUNCH_TOKEN}.yaml"
)
assert set(yaml.safe_load(values_file.read_text())) == {"run"}
class TestSnapshotFiles:
def test_the_recorded_files_are_exactly_the_declared_ones(self):
"""A renamed or deleted case must not leave an orphan golden behind."""
recorded = {path.name for path in SNAPSHOT_DIR.glob("*.txt")}
assert recorded == {"kubernetes_launch.txt"}
@@ -0,0 +1,42 @@
from miles.utils.external_utils.command_utils.helm_backend.naming import RunFiles, _orchestrator_state_path
from miles.utils.external_utils.command_utils.helm_backend.orchestrator.state import (
OrchestratorState,
OrchestratorStatus,
)
def _write(path, status: OrchestratorStatus, *, exit_code: int | None = None) -> None:
OrchestratorState(status=status, exit_code=exit_code).write(path)
def _state_file(tmp_path):
return _orchestrator_state_path(tmp_path, "260101-000000-000001")
class TestRunDir:
def test_places_a_run_under_the_shared_root(self):
"""Every pod resolves the same run directory from the shared storage mount and the run id."""
assert str(RunFiles.run_dir(shared_root="/cluster-storage/miles_data", run_id="260101-000000-000")).endswith(
"/cluster-storage/miles_data/miles-runs/260101-000000-000"
)
def test_keeps_the_state_file_in_a_state_subdirectory(self):
"""Grouping the machine-written state keeps it out of the way of a run's own outputs."""
path = _orchestrator_state_path("/runs/abc", "abc123")
assert path.as_posix() == "/runs/abc/state/orchestrator-abc123.state"
class TestLatestExitFile:
def test_names_no_file_before_a_launch_has_written_one(self, tmp_path):
"""A run directory a launch has only just created holds no verdict to collect."""
assert RunFiles.latest_state_file(run_directory=tmp_path) is None
def test_picks_the_newest_launch_rather_than_the_newest_write(self, tmp_path):
"""An earlier launch torn down after a later one started writes last, and its verdict is not the run's."""
later = _orchestrator_state_path(tmp_path, "260101-000200-000001")
earlier = _orchestrator_state_path(tmp_path, "260101-000100-000002")
_write(later, OrchestratorStatus.EXITED, exit_code=0)
_write(earlier, OrchestratorStatus.EXITED, exit_code=1)
assert RunFiles.latest_state_file(run_directory=tmp_path) == later
@@ -0,0 +1,328 @@
run:
colocate:
enabled: true
enginePool: inference-engine-0-1
trainerPool: trainer-engine-actor
env:
CUDA_DEVICE_MAX_CONNECTIONS: '1'
PYTHONUNBUFFERED: '1'
stateFile: /cluster-storage/myteam/miles_data/miles-runs/260101-000000-000/state/orchestrator-260101-000000-000001.state
id: 260101-000000-000
inferenceEngines:
- command:
- <PYTHON>
- -m
- sglang.launch_server
- --model-path
- <FIXTURE_DIR>/typical-model
- --host
- 0.0.0.0
- --port
- '8000'
- --disaggregation-mode
- prefill
- --trust-remote-code
- --mem-fraction-static
- '0.7'
- --chunked-prefill-size
- '4096'
- --nccl-port
- '10000'
- --dist-init-addr
- $(LWS_LEADER_ADDRESS):9000
- --nnodes
- '2'
- --node-rank
- $(LWS_WORKER_INDEX)
- --tp-size
- '16'
- --load-balance-method
- round_robin
- --random-seed
- <RANDOM_SEED>
- --skip-server-warmup
- --enable-metrics
- --cuda-graph-backend-prefill
- disabled
- --lora-use-virtual-experts
- --disaggregation-bootstrap-port
- '11000'
- --engine-info-bootstrap-port
- '12000'
- --enable-draft-weights-cpu-backup
env:
NVSHMEM_DISABLE_NCCL: '1'
RAY_EXPERIMENTAL_NOSET_ASCEND_RT_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_CUDA_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_HABANA_VISIBLE_MODULES: '1'
RAY_EXPERIMENTAL_NOSET_HIP_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_NEURON_RT_VISIBLE_CORES: '1'
RAY_EXPERIMENTAL_NOSET_ONEAPI_DEVICE_SELECTOR: '1'
RAY_EXPERIMENTAL_NOSET_TPU_VISIBLE_CHIPS: '1'
SGLANG_BATCH_INVARIANT_OPS_ENABLE_MM_FALLBACK_VARIANT: 'true'
SGLANG_DG_CACHE_DIR_PER_PROCESS: '1'
SGLANG_ENABLE_HEALTH_ENDPOINT_GENERATION: 'false'
SGLANG_ENABLE_STRICT_MEM_CHECK_DURING_IDLE: 'false'
SGLANG_ENABLE_TP_MEMORY_INBALANCE_CHECK: 'false'
SGLANG_JIT_DEEPGEMM_PRECOMPILE: 'false'
SGLANG_MEMORY_SAVER_CUDA_GRAPH: 'true'
SGLANG_OPT_USE_CUSTOM_ALL_REDUCE_V2: '1'
meta:
gpu_ids: 0,1,2,3,4,5,6,7
name: inference-engine-0-0
objectName: myrun-miles-run-inference-engine-0-0
poolId: inference-engine-0-0
ports:
- name: primary
port: 8000
- name: dist-init
port: 9000
- name: nccl
port: 10000
- name: disaggregation
port: 11000
- name: engine-info-boo
port: 12000
- name: gate
port: 13000
replicas: 1
resources:
limits:
nvidia.com/gpu: 8
size: 2
- command:
- <PYTHON>
- -m
- sglang.launch_server
- --model-path
- <FIXTURE_DIR>/typical-model
- --host
- 0.0.0.0
- --port
- '8000'
- --disaggregation-mode
- decode
- --trust-remote-code
- --mem-fraction-static
- '0.7'
- --chunked-prefill-size
- '4096'
- --nccl-port
- '10000'
- --dist-init-addr
- $(LWS_LEADER_ADDRESS):9000
- --node-rank
- $(LWS_WORKER_INDEX)
- --tp-size
- '8'
- --random-seed
- <RANDOM_SEED>
- --skip-server-warmup
- --enable-metrics
- --cuda-graph-backend-prefill
- disabled
- --lora-use-virtual-experts
- --engine-info-bootstrap-port
- '12000'
- --enable-memory-saver
- --enable-draft-weights-cpu-backup
env:
NVSHMEM_DISABLE_NCCL: '1'
RAY_EXPERIMENTAL_NOSET_ASCEND_RT_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_CUDA_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_HABANA_VISIBLE_MODULES: '1'
RAY_EXPERIMENTAL_NOSET_HIP_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_NEURON_RT_VISIBLE_CORES: '1'
RAY_EXPERIMENTAL_NOSET_ONEAPI_DEVICE_SELECTOR: '1'
RAY_EXPERIMENTAL_NOSET_TPU_VISIBLE_CHIPS: '1'
SGLANG_BATCH_INVARIANT_OPS_ENABLE_MM_FALLBACK_VARIANT: 'true'
SGLANG_DG_CACHE_DIR_PER_PROCESS: '1'
SGLANG_ENABLE_HEALTH_ENDPOINT_GENERATION: 'false'
SGLANG_ENABLE_STRICT_MEM_CHECK_DURING_IDLE: 'false'
SGLANG_ENABLE_TP_MEMORY_INBALANCE_CHECK: 'false'
SGLANG_JIT_DEEPGEMM_PRECOMPILE: 'false'
SGLANG_MEMORY_SAVER_CUDA_GRAPH: 'true'
SGLANG_OPT_USE_CUSTOM_ALL_REDUCE_V2: '1'
meta:
gpu_ids: 0,1,2,3,4,5,6,7
name: inference-engine-0-1
objectName: myrun-miles-run-inference-engine-0-1
poolId: inference-engine-0-1
ports:
- name: primary
port: 8000
- name: dist-init
port: 9000
- name: nccl
port: 10000
- name: engine-info-boo
port: 12000
- name: gate
port: 13000
replicas: 4
resources:
limits:
nvidia.com/gpu: 8
objectNames:
colocatePairing: myrun-miles-run-colocate-pairing
mooncakeMaster: myrun-miles-run-mooncake-master
orchestrator: myrun-miles-run-orchestrator
orchestrator:
command:
- python
- scripts/run_qwen3_4b.py
- train
- --cluster-backend
- kubernetes
staticWorkers:
- command:
- <PYTHON>
- -m
- miles.utils.workers.serving.serve
- --worker
- miles.ray.rollout.inference_controller.InferenceController
- --pool-id
- inference-controller
- --ctor-kwargs-fn
- miles.ray.specs.bootstrap.compute_ctor_kwargs
- --ranks-per-pod
- '1'
- --gpu-slots-per-rank
- '0'
- --
- --cluster-backend
- kubernetes
- --rollout-num-gpus
- '48'
name: inference-controller
objectName: myrun-miles-run-inference-controller
poolId: inference-controller
ports:
- name: rpc
port: 8000
replicas: 1
- command:
- <PYTHON>
- -m
- sglang_router.launch_router
- --port
- '8000'
- --pd-disaggregation
- --policy
- manual
- --assignment-mode
- min_load
- --log-level
- warn
- --prometheus-port
- '9000'
- --request-timeout-secs
- '14400'
name: inference-router-0
objectName: myrun-miles-run-inference-router-0
poolId: inference-router-0
ports:
- name: primary
port: 8000
- name: prometheus
port: 9000
replicas: 1
- command:
- <PYTHON>
- -m
- miles.rollout.session.server
- --config-json
- '{"host":"0.0.0.0","port":8000,"instance_id":"0123456789abcdef-0","backend_url":"http://myrun-miles-run-inference-router-0-0.myrun-miles-run-inference-router-0:8000","timeout":null,"hf_checkpoint":"<FIXTURE_DIR>/typical-model","chat_template_path":null,"tito_model":"default","apply_chat_template_kwargs":{},"use_rollout_routing_replay":false,"use_rollout_indexer_replay":false,"use_sampling_support_replay":false,"sglang_speculative_algorithm":null,"num_layers":36,"moe_router_topk":2,"save_debug_trajectory_data":null,"lora_rank":0,"lora_adapter_path":null,"lora_train_only":false}'
name: session-server
objectName: myrun-miles-run-session-server
poolId: session-server
ports:
- name: primary
port: 8000
replicas: 32
- command:
- <PYTHON>
- -m
- miles.utils.workers.serving.serve
- --worker
- miles.ray.rollout.rollout_executor.RolloutExecutor
- --pool-id
- rollout-executor
- --ctor-kwargs-fn
- miles.ray.specs.bootstrap.compute_ctor_kwargs
- --ranks-per-pod
- '1'
- --gpu-slots-per-rank
- '0'
- --
- --cluster-backend
- kubernetes
- --rollout-num-gpus
- '48'
name: rollout-executor
objectName: myrun-miles-run-rollout-executor
poolId: rollout-executor
ports:
- name: rpc
port: 8000
replicas: 1
- command:
- <PYTHON>
- -m
- miles.utils.workers.serving.serve
- --worker
- miles.ray.train.group.TrainerController
- --pool-id
- trainer-controller-actor
- --ctor-kwargs-fn
- miles.ray.specs.bootstrap.compute_ctor_kwargs
- --ranks-per-pod
- '1'
- --gpu-slots-per-rank
- '0'
- --
- --cluster-backend
- kubernetes
- --rollout-num-gpus
- '48'
name: trainer-controller-actor
objectName: myrun-miles-run-trainer-controller-actor
poolId: trainer-controller-actor
ports:
- name: rpc
port: 8000
replicas: 1
trainerEngines:
- command:
- bash
- -c
- mkdir -p /scratch/Qwen3-4B && rsync -a --info=progress2 /cluster-storage/models/Qwen3-4B/ /scratch/Qwen3-4B && exec <PYTHON> -m miles.utils.workers.process_supervisor --num-subprocesses 8 -- <PYTHON> -m miles.utils.workers.serving.serve --worker miles.backends.megatron_utils.actor.MegatronTrainRayActor --pool-id trainer-engine-actor --ctor-kwargs-fn miles.ray.specs.bootstrap.compute_ctor_kwargs --ranks-per-pod 8 --gpu-slots-per-rank 1 -- --cluster-backend kubernetes --rollout-num-gpus 48
env:
LD_PRELOAD: /usr/local/lib/python3.12/dist-packages/torch_memory_saver_hook_mode_preload_cu13.abi3.so
NCCL_CUMEM_ENABLE: '0'
NVSHMEM_DISABLE_NCCL: '1'
NVTE_FP8_BLOCK_SCALING_FP32_SCALES: '1'
RAY_EXPERIMENTAL_NOSET_ASCEND_RT_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_CUDA_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_HABANA_VISIBLE_MODULES: '1'
RAY_EXPERIMENTAL_NOSET_HIP_VISIBLE_DEVICES: '1'
RAY_EXPERIMENTAL_NOSET_NEURON_RT_VISIBLE_CORES: '1'
RAY_EXPERIMENTAL_NOSET_ONEAPI_DEVICE_SELECTOR: '1'
RAY_EXPERIMENTAL_NOSET_TPU_VISIBLE_CHIPS: '1'
TMS_INIT_ENABLE: '1'
TMS_INIT_ENABLE_CPU_BACKUP: '1'
meta:
gpu_ids: 0,1,2,3,4,5,6,7
name: trainer-engine-actor
objectName: myrun-miles-run-trainer-engine-actor
poolId: trainer-engine-actor
ports:
- name: master
port: 9000
- name: rpc
port: 8000
replicas: 1
resources:
limits:
nvidia.com/gpu: 8
size: 4
File diff suppressed because it is too large Load Diff