diff --git a/miles/utils/external_utils/command_utils/helm_backend/launcher/command_wrapper.py b/miles/utils/external_utils/command_utils/helm_backend/launcher/command_wrapper.py index 4e831ec0cc..2fc29d1f86 100644 --- a/miles/utils/external_utils/command_utils/helm_backend/launcher/command_wrapper.py +++ b/miles/utils/external_utils/command_utils/helm_backend/launcher/command_wrapper.py @@ -2,6 +2,7 @@ from __future__ import annotations import json import subprocess +from collections.abc import Sequence from pathlib import Path from typing import Any, TypeVar @@ -96,17 +97,33 @@ class Helm: _run(["helm", "dependency", "build", str(chart)], capture_output=False) @staticmethod - def list_releases(*, namespace: str, selector: str) -> list[str]: - listed = _run( - ["helm", "list", "--namespace", namespace, "--selector", selector, "--output", "json"], - capture_output=True, - ) + def list_releases(*, namespace: str, selector: str | None = None, name_filter: str | None = None) -> list[str]: + command = ["helm", "list", "--namespace", namespace, "--output", "json"] + if selector is not None: + command += ["--selector", selector] + if name_filter is not None: + command += ["--filter", name_filter] + listed = _run(command, capture_output=True) return [release["name"] for release in json.loads(listed.stdout or "[]")] @staticmethod def uninstall(*, release: str, namespace: str) -> None: _run(["helm", "uninstall", release, "--namespace", namespace], capture_output=False) + @staticmethod + def uninstall_if_present(*, release: str, namespace: str) -> None: + result = Helm.run_raw("uninstall", release, "--namespace", namespace) + if result.returncode == 0: + return + + reason = (result.stderr or result.stdout or "").strip() + if "not found" in reason.lower(): + return + raise RuntimeError( + f"Could not uninstall Helm release {release!r} from namespace {namespace!r}: " + f"exit code {result.returncode}: {reason}" + ) + @staticmethod def upgrade_command( release: str, namespace: str, chart: str | Path, values_files: list[str | Path], *, ci_run: bool @@ -233,6 +250,12 @@ class Kubectl: def release_selector(release: str) -> str: return f"{INSTANCE_LABEL}={release}" + @staticmethod + def releases_selector(releases: Sequence[str]) -> str: + if len(releases) == 1: + return Kubectl.release_selector(releases[0]) + return f"{INSTANCE_LABEL} in ({','.join(releases)})" + @staticmethod def job_selector(name: str) -> str: return f"{_JOB_NAME_LABEL}={name}" diff --git a/miles/utils/external_utils/command_utils/helm_backend/launcher/entrypoint.py b/miles/utils/external_utils/command_utils/helm_backend/launcher/entrypoint.py index 5dc9468973..670e58b56d 100644 --- a/miles/utils/external_utils/command_utils/helm_backend/launcher/entrypoint.py +++ b/miles/utils/external_utils/command_utils/helm_backend/launcher/entrypoint.py @@ -392,7 +392,13 @@ def _assert_upgrade_is_allowed( def _belongs_to_run(release: str, *, run_id: str) -> bool: - return (parsed := ReleaseName.parse(release)) is not None and parsed.run_id == run_id + try: + return (parsed := ReleaseName.parse(release)) is not None and parsed.run_id == run_id + except Exception: + logger.info( + f"Release {release} does not follow the Miles naming rules, so it belongs to no run", exc_info=True + ) + return False def _uninstall_leftover_ci_releases(namespace: str, *, keep_run_id: str) -> list[str]: @@ -422,11 +428,18 @@ def _defuse_previous_generation( def _collect_diagnosis(*, release: str, namespace: str, state_file: Path) -> None: + releases = _releases_of_run(release=release, namespace=namespace) + if len(releases) > 1: + logger.info( + f"This run is deployed as {len(releases)} releases, and a deployment other than {release} can be what " + f"failed, so the diagnosis covers all of them: {', '.join(releases)}" + ) + try: diagnosis = collect_diagnosis( namespace=namespace, output_dir=state_file.parent, - selector=Kubectl.release_selector(release), + selector=Kubectl.releases_selector(releases), state_file=state_file, ) except Exception: @@ -438,6 +451,23 @@ def _collect_diagnosis(*, release: str, namespace: str, state_file: Path) -> Non logger.warning(f"The diagnosis is incomplete, these could not be collected: {', '.join(diagnosis.missing)}") +def _releases_of_run(*, release: str, namespace: str) -> list[str]: + parsed = ReleaseName.parse(release) + if parsed is None: + return [release] + + try: + listed = Helm.list_releases( + namespace=namespace, name_filter=f"^{re.escape(ReleaseName.run_prefix(run_id=parsed.run_id))}" + ) + except Exception: + logger.warning(f"Could not list what else run {parsed.run_id} installed; diagnosing {release} alone") + logger.debug("Listing the releases of the run failed", exc_info=True) + return [release] + + return sorted({one for one in listed if _belongs_to_run(one, run_id=parsed.run_id)} | {release}) + + def _write_helm_values(path: Path, values: dict[str, Any]) -> None: path.parent.mkdir(parents=True, exist_ok=True) path.write_text(yaml.dump(values, Dumper=_HelmValuesDumper, default_flow_style=False, sort_keys=True)) diff --git a/miles/utils/external_utils/command_utils/helm_backend/naming.py b/miles/utils/external_utils/command_utils/helm_backend/naming.py index b2787a8da8..c82b6009ea 100644 --- a/miles/utils/external_utils/command_utils/helm_backend/naming.py +++ b/miles/utils/external_utils/command_utils/helm_backend/naming.py @@ -64,6 +64,10 @@ class ReleaseName(FrozenStrictBaseModel): parts.append(self.deploy_instance_id) return "-".join(parts) + @staticmethod + def run_prefix(*, run_id: str) -> str: + return f"{CHART_NAME}-{run_id}-" + @classmethod def parse(cls, release: str) -> ReleaseName | None: if not release.startswith(f"{CHART_NAME}-"): diff --git a/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_command_wrapper.py b/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_command_wrapper.py index f8b5365bce..39529afaa6 100644 --- a/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_command_wrapper.py +++ b/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_command_wrapper.py @@ -1,5 +1,6 @@ import json import subprocess +from pathlib import Path import pytest @@ -11,6 +12,7 @@ from miles.utils.external_utils.command_utils.helm_backend.launcher.command_wrap Kubectl, ) from miles.utils.external_utils.command_utils.helm_backend.naming import RUN_ID_MAX_LENGTH, ReleaseName +from miles.utils.workers.k8s_types import Pod from miles.utils.workers.types import DeployComponent from miles.utils.workers.worker_provider.kubernetes.helm import naming from miles.utils.workers.worker_provider.kubernetes.helm.env import INSTANCE_LABEL @@ -86,6 +88,67 @@ class TestUpgradeCommand: assert command.index("--labels") > command.index("--namespace") +class TestBuildDependencies: + def test_chart_dependencies_are_rebuilt_only_when_a_locked_dependency_is_missing( + self, monkeypatch: pytest.MonkeyPatch, tmp_path: Path + ) -> None: + """Helm rebuilds a locked chart only after one of its cached dependencies disappears.""" + chart = tmp_path / "chart" + charts = chart / "charts" + charts.mkdir(parents=True) + (chart / "Chart.lock").write_text("dependencies:\n - name: worker\n - name: runtime\n") + (charts / "worker").mkdir() + runtime = charts / "runtime" + runtime.mkdir() + commands: list[list[str]] = [] + + def fake_run_process( + argv: list[str], + *, + capture_output: bool, + check: bool, + input: str | None = None, + timeout: float | None = None, + ) -> subprocess.CompletedProcess[str]: + commands.append(argv) + return subprocess.CompletedProcess(args=argv, returncode=0, stdout="", stderr="") + + monkeypatch.setattr(command_wrapper, "run_process", fake_run_process) + + Helm.build_dependencies(chart) + runtime.rmdir() + Helm.build_dependencies(chart) + + assert commands == [["helm", "dependency", "build", str(chart)]] + + +class TestGetManifest: + def test_a_missing_release_is_reported_as_absent(self, monkeypatch: pytest.MonkeyPatch) -> None: + """Helm reports a missing release on stdout in some versions, and that is genuine absence.""" + monkeypatch.setattr( + command_wrapper, + "run_process", + lambda *_args, **_kwargs: subprocess.CompletedProcess( + [], returncode=1, stdout="Error: release: not found", stderr="" + ), + ) + + assert Helm.get_manifest("miles-run-a", "ci") is None + + def test_a_manifest_read_failure_other_than_absence_is_reported(self, monkeypatch: pytest.MonkeyPatch) -> None: + """An authorization or network failure must remain distinguishable from a missing release.""" + monkeypatch.setattr( + command_wrapper, + "run_process", + lambda *_args, **_kwargs: subprocess.CompletedProcess( + [], returncode=1, stdout="", stderr="Kubernetes cluster unreachable: access forbidden" + ), + ) + + with pytest.raises(RuntimeError, match="Kubernetes cluster unreachable: access forbidden"): + Helm.get_manifest("miles-run-a", "ci") + + def _kubectl_answering(monkeypatch, *, returncode: int, stdout: str = "", stderr: str = "") -> list[list[str]]: commands: list[list[str]] = [] @@ -118,6 +181,15 @@ class TestRequestBounds: assert [timeout for _, timeout in calls] == [30.0, 60.0] +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 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.""" @@ -252,6 +324,51 @@ class TestCiCleanup: assert commands[1] == ["helm", "uninstall", "miles-run-a", "--namespace", "ci"] +def _listed_releases(monkeypatch: pytest.MonkeyPatch, **kwargs) -> 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([{"name": "miles-run-a-all"}]), stderr="" + ) + + monkeypatch.setattr(command_wrapper, "_run", fake_run) + Helm.list_releases(**kwargs) + return commands[0] + + +class TestListReleases: + def test_a_name_filter_is_passed_to_helm_rather_than_applied_afterwards(self, monkeypatch): + """helm truncates to its 256-release maximum before returning, and a filter is applied before that.""" + command = _listed_releases(monkeypatch, namespace="rl", name_filter="^miles-run-a-") + + assert command[command.index("--filter") + 1] == "^miles-run-a-" + + def test_a_listing_that_names_no_filter_asks_for_the_whole_namespace(self, monkeypatch): + """The ci cleanup selects on a label instead, and a filter here would hide the releases it removes.""" + assert "--filter" not in _listed_releases(monkeypatch, namespace="rl") + + def test_a_filter_and_a_label_selector_narrow_the_same_listing(self, monkeypatch): + """Neither is a replacement for the other, so passing one may not drop the other.""" + command = _listed_releases(monkeypatch, namespace="rl", selector="ci=true", name_filter="^miles-run-a-") + + assert command[command.index("--selector") + 1] == "ci=true" + assert command[command.index("--filter") + 1] == "^miles-run-a-" + + def test_it_reads_back_the_names_helm_reported(self, monkeypatch): + """The filter is only worth passing if what comes back is still the list of release names.""" + monkeypatch.setattr( + command_wrapper, + "_run", + lambda command, capture_output: subprocess.CompletedProcess( + args=command, returncode=0, stdout=json.dumps([{"name": "miles-run-a-all"}]), stderr="" + ), + ) + + assert Helm.list_releases(namespace="rl", name_filter="^miles-run-a-") == ["miles-run-a-all"] + + class TestChartDir: def test_finds_the_chart_inside_the_checkout(self): """The launcher installs the chart of the code it runs, not one from a registry.""" @@ -270,6 +387,22 @@ class TestReleaseName: """The launcher finds a run's release again from the run id alone, so the rule is fixed.""" assert _unsplit("260101-000000-000") == "miles-run-260101-000000-000-all" + def test_the_run_prefix_is_what_every_release_of_that_run_starts_with(self): + """A filter built from it has to match the trainer and inference releases of the run as well.""" + prefix = ReleaseName.run_prefix(run_id="260101-000000-000") + + assert _unsplit("260101-000000-000").startswith(prefix) + engines = ReleaseName( + run_id="260101-000000-000", deploy_component=DeployComponent.INFERENCE, deploy_instance_id="dc1" + ) + assert engines.serialize().startswith(prefix) + + def test_the_run_prefix_ends_where_the_component_begins(self): + """Without the separator the prefix of one run also matches a run whose id merely starts with it.""" + assert not ReleaseName.run_prefix(run_id="260101-000000-000").startswith( + ReleaseName.run_prefix(run_id="260101-000000-00") + ) + 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 _unsplit(LONGEST_RUN_ID) == _unsplit(LONGEST_RUN_ID) diff --git a/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_run_diagnosis.py b/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_run_diagnosis.py new file mode 100644 index 0000000000..dcc0aa2a09 --- /dev/null +++ b/tests/fast/utils/external_utils/command_utils/helm_backend/launcher/test_run_diagnosis.py @@ -0,0 +1,131 @@ +import json +import re +import subprocess + +import pytest + +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 Kubectl +from miles.utils.external_utils.command_utils.helm_backend.naming import ReleaseName +from miles.utils.workers.types import DeployComponent +from miles.utils.workers.worker_provider.kubernetes.helm.env import INSTANCE_LABEL + +RUN_ID = "260101-000000-000" +OTHER_RUN_ID = "260101-111111-111" + + +def _release(component: DeployComponent, instance_id: str | None = None, run_id: str = RUN_ID) -> str: + return ReleaseName(run_id=run_id, deploy_component=component, deploy_instance_id=instance_id).serialize() + + +PRIMARY = _release(DeployComponent.PRIMARY) +TRAINER = _release(DeployComponent.TRAINER) +ENGINE = _release(DeployComponent.INFERENCE, "e0") +STRANGER = _release(DeployComponent.ALL, run_id=OTHER_RUN_ID) + + +def _helm_listing(monkeypatch: pytest.MonkeyPatch, releases: list[str]) -> list[list[str]]: + commands: list[list[str]] = [] + + def fake_run(command: list[str], capture_output: bool) -> subprocess.CompletedProcess: + commands.append(command) + body = json.dumps([{"name": release} for release in releases]) + return subprocess.CompletedProcess(args=command, returncode=0, stdout=body, stderr="") + + monkeypatch.setattr(command_wrapper, "_run", fake_run) + return commands + + +class TestTheReleasesADiagnosisCovers: + def test_covers_every_deployment_the_run_installed(self, monkeypatch): + """A split run fails in whichever deployment broke, and that is rarely the one being followed.""" + _helm_listing(monkeypatch, [PRIMARY, TRAINER, ENGINE, STRANGER]) + + found = entrypoint._releases_of_run(release=PRIMARY, namespace="rl") + + assert found == sorted([PRIMARY, TRAINER, ENGINE]) + + def test_leaves_another_run_in_the_namespace_alone(self, monkeypatch): + """Namespaces hold more than one run, and describing someone else's pods reads their experiment.""" + _helm_listing(monkeypatch, [PRIMARY, STRANGER]) + + assert STRANGER not in entrypoint._releases_of_run(release=PRIMARY, namespace="rl") + + def test_asks_helm_for_the_whole_namespace_because_a_run_carries_no_label(self, monkeypatch): + """Releases are labelled per release, so the siblings can only be found by listing and parsing.""" + commands = _helm_listing(monkeypatch, [PRIMARY]) + + entrypoint._releases_of_run(release=PRIMARY, namespace="rl") + + assert "--selector" not in commands[0] + assert commands[0][commands[0].index("--namespace") + 1] == "rl" + + def test_falls_back_to_the_followed_release_when_helm_cannot_be_asked(self, monkeypatch): + """A diagnosis of one deployment beats no diagnosis at all when the run has already failed.""" + + def fake_run(command: list[str], capture_output: bool) -> subprocess.CompletedProcess: + raise RuntimeError("helm is not reachable") + + monkeypatch.setattr(command_wrapper, "_run", fake_run) + + assert entrypoint._releases_of_run(release=PRIMARY, namespace="rl") == [PRIMARY] + + def test_keeps_the_followed_release_even_when_helm_forgot_it(self, monkeypatch): + """The release being diagnosed is the one known to have failed, listed or not.""" + _helm_listing(monkeypatch, []) + + assert entrypoint._releases_of_run(release=PRIMARY, namespace="rl") == [PRIMARY] + + def test_asks_helm_only_for_the_releases_of_this_run(self, monkeypatch): + """helm truncates a busy namespace to 256 releases, which can drop the sibling that actually failed.""" + commands = _helm_listing(monkeypatch, [PRIMARY]) + + entrypoint._releases_of_run(release=PRIMARY, namespace="rl") + + assert commands[0][commands[0].index("--filter") + 1] == f"^{re.escape(ReleaseName.run_prefix(run_id=RUN_ID))}" + + def test_the_filter_is_a_regex_the_run_id_cannot_break_out_of(self): + """helm reads it as a regular expression, and a run id is not written to be one.""" + assert re.fullmatch(f"^{re.escape(ReleaseName.run_prefix(run_id=RUN_ID))}.*", PRIMARY) + + +class TestTheSelectorItDescribesWith: + def test_matches_every_release_of_the_run_at_once(self): + """One kubectl call per describe keeps the failure path short while covering all deployments.""" + selector = Kubectl.releases_selector([PRIMARY, TRAINER]) + + assert selector == f"{INSTANCE_LABEL} in ({PRIMARY},{TRAINER})" + + def test_stays_an_equality_for_an_unsplit_run(self): + """An all-in-one run has one release, and the plain selector is what its snapshots already carry.""" + assert Kubectl.releases_selector([PRIMARY]) == Kubectl.release_selector(PRIMARY) + + +# Legal for helm, but not for miles: the run id is one character longer than the naming rules allow. +UNPARSABLE_NEIGHBOUR = "miles-run-" + "a" * 34 + "-all" + + +class TestAReleaseTheMilesRulesCannotRead: + def test_it_belongs_to_no_run(self): + """ReleaseName.parse builds a validated model, so such a name raised out of the candidate loop.""" + assert entrypoint._belongs_to_run(UNPARSABLE_NEIGHBOUR, run_id=RUN_ID) is False + + def test_a_release_of_another_naming_scheme_belongs_to_no_run(self): + """A namespace holds releases of other charts, and none of them is part of this run.""" + assert entrypoint._belongs_to_run("someone-elses-release", run_id=RUN_ID) is False + + def test_a_release_of_this_run_still_belongs_to_it(self): + """Reading an unreadable name as belonging to nothing may not cost the readable ones.""" + assert entrypoint._belongs_to_run(PRIMARY, run_id=RUN_ID) is True + + def test_such_a_neighbour_does_not_replace_the_diagnosis_of_the_failed_run(self, monkeypatch): + """That loop runs outside both try blocks, so one neighbour replaced the diagnosis with a traceback.""" + _helm_listing(monkeypatch, [PRIMARY, TRAINER, UNPARSABLE_NEIGHBOUR]) + + assert entrypoint._releases_of_run(release=PRIMARY, namespace="rl") == sorted([PRIMARY, TRAINER]) + + def test_such_a_neighbour_is_not_diagnosed_as_part_of_the_run(self, monkeypatch): + """Describing someone else's pods reads an experiment this run has nothing to do with.""" + _helm_listing(monkeypatch, [PRIMARY, UNPARSABLE_NEIGHBOUR]) + + assert UNPARSABLE_NEIGHBOUR not in entrypoint._releases_of_run(release=PRIMARY, namespace="rl")