diff --git a/openapi/ga/individual/platform.openapi.yaml b/openapi/ga/individual/platform.openapi.yaml index 5376947a7e..33a0d3c022 100644 --- a/openapi/ga/individual/platform.openapi.yaml +++ b/openapi/ga/individual/platform.openapi.yaml @@ -4907,7 +4907,16 @@ paths: tags: - Jobs summary: Get Execution Profiles - description: Get all currently configured execution profiles. + description: 'Get all currently configured execution profiles. + + + Returns the capability-filtered merge from jobs config. In local standalone + + the controller may prune the shared list further after registry boot; in + + split topologies the API advertises its own merge result (not controller + + process memory).' operationId: get_execution_profiles_apis_jobs_v2_execution_profiles_get responses: '200': diff --git a/openapi/ga/openapi.yaml b/openapi/ga/openapi.yaml index 5376947a7e..33a0d3c022 100644 --- a/openapi/ga/openapi.yaml +++ b/openapi/ga/openapi.yaml @@ -4907,7 +4907,16 @@ paths: tags: - Jobs summary: Get Execution Profiles - description: Get all currently configured execution profiles. + description: 'Get all currently configured execution profiles. + + + Returns the capability-filtered merge from jobs config. In local standalone + + the controller may prune the shared list further after registry boot; in + + split topologies the API advertises its own merge result (not controller + + process memory).' operationId: get_execution_profiles_apis_jobs_v2_execution_profiles_get responses: '200': diff --git a/openapi/openapi.yaml b/openapi/openapi.yaml index 5376947a7e..33a0d3c022 100644 --- a/openapi/openapi.yaml +++ b/openapi/openapi.yaml @@ -4907,7 +4907,16 @@ paths: tags: - Jobs summary: Get Execution Profiles - description: Get all currently configured execution profiles. + description: 'Get all currently configured execution profiles. + + + Returns the capability-filtered merge from jobs config. In local standalone + + the controller may prune the shared list further after registry boot; in + + split topologies the API advertises its own merge result (not controller + + process memory).' operationId: get_execution_profiles_apis_jobs_v2_execution_profiles_get responses: '200': diff --git a/packages/nemo_platform/pyproject.toml b/packages/nemo_platform/pyproject.toml index dc7a0c0122..6a7aa8af0b 100644 --- a/packages/nemo_platform/pyproject.toml +++ b/packages/nemo_platform/pyproject.toml @@ -197,6 +197,7 @@ jobs-service = [ "pyyaml>=6.0.2", "hvac>=2.3.0", "nmp-common", + "nemo-platform-plugin", "aiosqlite>=0.20.0", "duckdb<2.0.0,>=1.1.3", "pandas>=1.5.3", diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/services/cli.py b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/services/cli.py index e71a09da26..5c58af463d 100644 --- a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/services/cli.py +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/services/cli.py @@ -14,6 +14,7 @@ import httpx import typer from nemo_platform_ext.cli.core.help_formatter import create_typer_app +from nemo_platform_ext.cli.docker_preflight import require_docker_for_default_local from nemo_platform_ext.local.process import ( ForegroundInstanceError, InstanceAlreadyRunningError, @@ -236,16 +237,6 @@ def run_services( _ensure_port_available(host, port, scope, base_dir=base_dir) - try: - lock_fd = acquire_lock(scope, base_dir=base_dir) - except InstanceAlreadyRunningError: - _fail_already_running(scope, base_dir) - - # _NMP_LAUNCH_MODE is set by start_background() when this process was - # spawned via ``nemo services start``. Without it we default to - # "foreground", which protects interactive ``run`` sessions from being - # killed by ``stop``. - mode = "background" if os.environ.get("_NMP_LAUNCH_MODE") == "background" else "foreground" platform_config = PlatformAppConfig( services=_parse_csv_option(services), service_group=service_group, @@ -259,6 +250,18 @@ def run_services( keep_alive_timeout_seconds=keep_alive_timeout_seconds, state_root=base_dir, ) + require_docker_for_default_local(platform_config) + + try: + lock_fd = acquire_lock(scope, base_dir=base_dir) + except InstanceAlreadyRunningError: + _fail_already_running(scope, base_dir) + + # _NMP_LAUNCH_MODE is set by start_background() when this process was + # spawned via ``nemo services start``. Without it we default to + # "foreground", which protects interactive ``run`` sessions from being + # killed by ``stop``. + mode = "background" if os.environ.get("_NMP_LAUNCH_MODE") == "background" else "foreground" desc = InstanceDescriptor.from_config( platform_config, @@ -382,6 +385,7 @@ def start_services( keep_alive_timeout_seconds=keep_alive_timeout_seconds, state_root=base_dir, ) + require_docker_for_default_local(platform_config) typer.echo("Starting platform services...") proc = start_background(platform_config) @@ -565,11 +569,6 @@ def restart_services( ) raise typer.Exit(1) - typer.echo("Stopping platform services...") - # restart always produces a background instance, so force=True is - # appropriate even for foreground targets. - stop_instance(scope, base_dir=base_dir, force=True) - previous_config = prev.config if prev else None effective_services = _parse_csv_option(services) if services is not None else None if services is None and previous_config is not None: @@ -601,9 +600,6 @@ def restart_services( ) ) - _warn_bind_all(effective_host) - - _ensure_port_available(effective_host, effective_port, scope, base_dir=base_dir) platform_config = PlatformAppConfig( services=effective_services, service_group=effective_service_group, @@ -617,6 +613,17 @@ def restart_services( keep_alive_timeout_seconds=effective_keep_alive_timeout_seconds, state_root=base_dir, ) + # Preflight before stop so a missing Docker daemon does not tear down a healthy instance. + require_docker_for_default_local(platform_config) + + typer.echo("Stopping platform services...") + # restart always produces a background instance, so force=True is + # appropriate even for foreground targets. + stop_instance(scope, base_dir=base_dir, force=True) + + _warn_bind_all(effective_host) + + _ensure_port_available(effective_host, effective_port, scope, base_dir=base_dir) typer.echo("Starting platform services...") proc = start_background(platform_config) diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/setup.py b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/setup.py index 80bc4f10a7..ecfc915d6d 100644 --- a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/setup.py +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/setup.py @@ -27,8 +27,8 @@ import typer import yaml as _yaml from nemo_platform import NeMoPlatform +from nemo_platform_plugin.capabilities import probe_docker from nemo_platform_plugin.client.adapter import client_from_platform -from nemo_platform_plugin.config import validate_docker_available from nemo_platform_plugin.secrets.client import SecretsClient from nemo_platform_plugin.secrets.types import PlatformSecretCreateRequest, PlatformSecretUpdateRequest from nmp.common.config import nmp_user_data_dir @@ -43,6 +43,7 @@ from nemo_platform_ext.cli.commands.skills.registry import get_installer, load_skills from nemo_platform_ext.cli.core.context import CLIContext from nemo_platform_ext.cli.core.errors import handle_errors +from nemo_platform_ext.cli.docker_preflight import DOCKER_PREFLIGHT_MESSAGE, require_docker_for_default_local from nemo_platform_ext.cli.telemetry import emit from nemo_platform_ext.cli.telemetry.events import OnboardingStepEvent, TaskStatusEnum from nemo_platform_ext.client.tls import client_verify_from_env @@ -1057,7 +1058,7 @@ def _should_hint_docker_unavailable(*, exit_code: int | None, log_path: Path | N """ if _services_log_suggests_docker_failure(log_path): return True - if exit_code is not None and not validate_docker_available(): + if exit_code is not None and not probe_docker(use_cache=False).available: return True return False @@ -1158,6 +1159,9 @@ def _maybe_start_services( console.print(" [cyan]PYO3_USE_ABI3_FORWARD_COMPATIBILITY=1 pip install 'nemo-platform\\[all]'[/cyan]") raise typer.Exit(1) + # Fail before stop/spawn when default local needs Docker (NVBug 6537617). + require_docker_for_default_local(console=console) + if already_running: console.print(" Restarting platform services...") _kill_existing_services(base_url) @@ -1186,10 +1190,7 @@ def _maybe_start_services( console.print(f"{CROSS} Platform did not become ready within {timeout}s") console.print(f" Check {log} for details.") if _should_hint_docker_unavailable(exit_code=exit_code, log_path=log): - console.print( - " Docker does not appear to be available. " - "Install and start Docker, or configure non-Docker executors, then retry." - ) + console.print(f" {DOCKER_PREFLIGHT_MESSAGE}") raise typer.Exit(1) console.print(f"{CHECK} Platform running at {base_url} (pid {proc.pid})\n") diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/docker_preflight.py b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/docker_preflight.py new file mode 100644 index 0000000000..c57ffc8c1c --- /dev/null +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/docker_preflight.py @@ -0,0 +1,147 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Config-aware Docker preflight for local platform startup (NVBug 6537617). + +Fails fast before spawn/wait when the resolved run would start the deployments +service/controller with a docker-backed ``default_executor`` while the daemon is +unreachable. Does not treat Docker as a global platform dependency — kubernetes and +reduced selections that never hit that fail-close path are left alone. +""" + +from __future__ import annotations + +from dataclasses import dataclass, replace +from pathlib import Path +from typing import Any + +import typer +import yaml +from nemo_platform_plugin.capabilities import probe_docker +from nmp.platform_runner.config import PlatformAppConfig, default_config_path, resolve_run_configuration +from rich.console import Console + +DOCKER_PREFLIGHT_MESSAGE = ( + "Docker is required for this default local setup (deployments default_executor " + "uses the docker backend) but the Docker daemon is not available. " + "Install and start Docker, or use a kubernetes / non-docker config " + "(and omit the deployments service if you do not need it), then retry." +) + + +@dataclass(frozen=True) +class _DefaultDockerExecutorProbe: + """Whether the default deployments executor is docker, and its optional host override.""" + + is_docker: bool + docker_host: str | None = None + + +def _load_platform_yaml(config_path: str) -> dict[str, Any]: + """Load platform YAML without running Pydantic validators (no soft-downgrade).""" + path = Path(config_path) + if not path.is_file(): + return {} + with path.open(encoding="utf-8") as fh: + data = yaml.safe_load(fh) or {} + return data if isinstance(data, dict) else {} + + +def _intended_runtime_is_kubernetes(raw: dict[str, Any]) -> bool: + platform = raw.get("platform") + if not isinstance(platform, dict): + return False + runtime = platform.get("runtime") + return isinstance(runtime, str) and runtime.strip().lower() == "kubernetes" + + +def _default_docker_executor_probe(raw: dict[str, Any]) -> _DefaultDockerExecutorProbe: + """Resolve deployments.default_executor → docker backend + optional config.docker_host.""" + deployments = raw.get("deployments") + if not isinstance(deployments, dict): + return _DefaultDockerExecutorProbe(is_docker=False) + default_name = deployments.get("default_executor") + if not isinstance(default_name, str) or not default_name: + return _DefaultDockerExecutorProbe(is_docker=False) + executors = deployments.get("executors") + if not isinstance(executors, list): + return _DefaultDockerExecutorProbe(is_docker=False) + for spec in executors: + if not isinstance(spec, dict): + continue + if spec.get("name") != default_name: + continue + backend = spec.get("backend") + if not (isinstance(backend, str) and backend.strip().lower() == "docker"): + return _DefaultDockerExecutorProbe(is_docker=False) + docker_host: str | None = None + config = spec.get("config") + if isinstance(config, dict): + host = config.get("docker_host") + if isinstance(host, str) and host.strip(): + docker_host = host.strip() + return _DefaultDockerExecutorProbe(is_docker=True, docker_host=docker_host) + return _DefaultDockerExecutorProbe(is_docker=False) + + +def _resolve_default_local_docker_probe( + platform_config: PlatformAppConfig | None = None, + *, + config_path: str | None = None, +) -> _DefaultDockerExecutorProbe | None: + """Return the docker probe target when this run needs the default-local gate; else None.""" + app_config = platform_config or PlatformAppConfig() + if config_path is not None: + app_config = replace(app_config, config_path=config_path) + + try: + resolved = resolve_run_configuration(app_config) + except ValueError: + # Invalid selections fail elsewhere; do not block on Docker here. + return None + + starts_deployments = "deployments" in resolved.services or "deployments" in resolved.controllers + if not starts_deployments: + return None + + raw = _load_platform_yaml(resolved.config_path or default_config_path()) + if _intended_runtime_is_kubernetes(raw): + return None + + probe_target = _default_docker_executor_probe(raw) + if not probe_target.is_docker: + return None + return probe_target + + +def default_local_needs_docker( + platform_config: PlatformAppConfig | None = None, + *, + config_path: str | None = None, +) -> bool: + """Return whether this run would fail closed on a missing Docker daemon.""" + return _resolve_default_local_docker_probe(platform_config, config_path=config_path) is not None + + +def require_docker_for_default_local( + platform_config: PlatformAppConfig | None = None, + *, + config_path: str | None = None, + console: Console | None = None, +) -> None: + """Exit with a clear message when default-local Docker is required but missing.""" + probe_target = _resolve_default_local_docker_probe(platform_config, config_path=config_path) + if probe_target is None: + return + if probe_docker(docker_host=probe_target.docker_host, use_cache=False).available: + return + out = console or Console(stderr=True) + out.print(f"[red]✗[/red] {DOCKER_PREFLIGHT_MESSAGE}") + raise typer.Exit(1) + + +__all__ = [ + "DOCKER_PREFLIGHT_MESSAGE", + "default_local_needs_docker", + "require_docker_for_default_local", +] diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/preflight.py b/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/preflight.py index f7f9290657..0cccad3bf2 100644 --- a/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/preflight.py +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/preflight.py @@ -11,6 +11,8 @@ from dataclasses import dataclass from enum import Enum +from nemo_platform_plugin.capabilities import probe_docker + from .config import QuickstartConfig @@ -96,11 +98,9 @@ def run_all(self) -> list[PreflightResult]: def _check_docker_available(self) -> None: """Verify Docker daemon is running and accessible.""" - try: - import docker - - client = docker.from_env() - client.ping() + # Preflight can be re-run after the user starts Docker. + result = probe_docker(use_cache=False) + if result.available: self.results.append( PreflightResult( name="Docker Available", @@ -108,13 +108,13 @@ def _check_docker_available(self) -> None: message="Docker daemon is running", ) ) - except Exception as e: + else: self.results.append( PreflightResult( name="Docker Available", status=CheckStatus.FAIL, message="Docker daemon is not accessible", - details=str(e), + details=result.detail or "Docker daemon unreachable", ) ) diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/validators.py b/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/validators.py index 966c8110f7..216f0a172e 100644 --- a/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/validators.py +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/quickstart/validators.py @@ -9,6 +9,8 @@ from dataclasses import dataclass from pathlib import Path +from nemo_platform_plugin.capabilities import probe_docker + from .config import QuickstartConfig @@ -25,19 +27,16 @@ def __bool__(self) -> bool: def validate_docker_available() -> ValidationResult: - """Check if Docker daemon is available. + """Check if Docker daemon is available via an uncached reachability probe (≤5s). Returns: ValidationResult indicating success or failure. """ - try: - import docker - - client = docker.from_env() - client.ping() + # CLI may re-run after the user starts Docker; do not pin a prior miss. + result = probe_docker(use_cache=False) + if result.available: return ValidationResult(True, "Docker is available") - except Exception as e: - return ValidationResult(False, f"Docker is not available: {e}") + return ValidationResult(False, result.detail or "Docker is not available") def validate_ngc_credentials(api_key: str) -> ValidationResult: diff --git a/packages/nemo_platform_ext/tests/cli/commands/test_setup.py b/packages/nemo_platform_ext/tests/cli/commands/test_setup.py index f6265c2ebe..8997f3be38 100644 --- a/packages/nemo_platform_ext/tests/cli/commands/test_setup.py +++ b/packages/nemo_platform_ext/tests/cli/commands/test_setup.py @@ -702,6 +702,8 @@ def maybe_start_preflight_mocks(): patch(f"{SETUP_MOD}._start_services_background") as mock_start, patch(f"{SETUP_MOD}.prompt_choice", return_value="yes"), patch(f"{SETUP_MOD}._prompt_data_dir", return_value="/tmp/data"), + # Docker preflight is covered separately; keep other start-path tests focused. + patch(f"{SETUP_MOD}.require_docker_for_default_local"), ): yield mock_start @@ -814,6 +816,7 @@ def test_restarts_when_running_and_start_services_true(self): patch(f"{SETUP_MOD}._prompt_data_dir", return_value="/tmp/test-data") as mock_db_prompt, patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._ensure_port_available_for_start", wraps=_ensure_port_available_for_start) as mock_port, + patch(f"{SETUP_MOD}.require_docker_for_default_local"), patch(f"{SETUP_MOD}._pause"), ): mock_start.return_value = MagicMock(pid=999) @@ -834,6 +837,7 @@ def test_restarts_exits_when_port_still_occupied_after_kill(self, capsys): patch(f"{SETUP_MOD}._start_services_background") as mock_start, patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=conflict), patch(f"{SETUP_MOD}._prompt_data_dir", return_value="/tmp/test-data"), + patch(f"{SETUP_MOD}.require_docker_for_default_local"), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), ): @@ -882,7 +886,7 @@ def test_early_exit_names_eaddrinuse_when_startup_log_has_bind_failure( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", return_value=False), patch(f"{SETUP_MOD}.log_path_for", return_value=log), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=True), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=True)), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), ): @@ -911,7 +915,7 @@ def test_early_exit_prints_docker_hint_when_daemon_unavailable(self, maybe_start with ( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", wait), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=False), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=False)), patch(f"{SETUP_MOD}.log_path_for", return_value=MagicMock(__str__=lambda self: "/tmp/services.log")), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), @@ -923,7 +927,20 @@ def test_early_exit_prints_docker_hint_when_daemon_unavailable(self, maybe_start captured = capsys.readouterr() assert "exited early (exit code 3)" in captured.err assert "Check /tmp/services.log for details." in captured.err - assert "Docker does not appear to be available" in captured.err + assert "Docker is required for this default local setup" in captured.err + + def test_docker_preflight_blocks_before_spawn(self, maybe_start_preflight_mocks, capsys): + """Default-local Docker gate exits before start_background (NVBug 6537617).""" + with ( + patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), + patch( + f"{SETUP_MOD}.require_docker_for_default_local", + side_effect=typer.Exit(1), + ), + pytest.raises(ClickExit), + ): + _maybe_start_services("http://localhost:8080", auto=False, start_services=True) + maybe_start_preflight_mocks.assert_not_called() def test_readiness_timeout_does_not_hint_docker_without_evidence( self, maybe_start_preflight_mocks, capsys, tmp_path @@ -935,7 +952,7 @@ def test_readiness_timeout_does_not_hint_docker_without_evidence( with ( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", return_value=False), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=False), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=False)), patch(f"{SETUP_MOD}.log_path_for", return_value=log), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), @@ -954,7 +971,7 @@ def test_docker_hint_from_services_log_markers(self, maybe_start_preflight_mocks with ( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", return_value=False), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=True), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=True)), patch(f"{SETUP_MOD}.log_path_for", return_value=log), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), @@ -962,7 +979,7 @@ def test_docker_hint_from_services_log_markers(self, maybe_start_preflight_mocks maybe_start_preflight_mocks.return_value = alive _maybe_start_services("http://localhost:8080", auto=False, start_services=True) captured = capsys.readouterr() - assert "Docker does not appear to be available" in captured.err + assert "Docker is required for this default local setup" in captured.err def test_startup_port_conflict_prefers_live_port_probe(self, tmp_path): log = tmp_path / "services.log" @@ -1273,6 +1290,7 @@ def test_auto_mode_skips_prompt_and_uses_persisted(self, tmp_path, monkeypatch): patch(f"{SETUP_MOD}._start_services_background") as mock_start, patch(f"{SETUP_MOD}._wait_for_platform", return_value=True), patch(f"{SETUP_MOD}._prompt_data_dir") as mock_prompt, + patch(f"{SETUP_MOD}.require_docker_for_default_local"), patch(f"{SETUP_MOD}._pause"), ): mock_start.return_value = MagicMock(pid=999) diff --git a/packages/nemo_platform_ext/tests/cli/test_docker_preflight.py b/packages/nemo_platform_ext/tests/cli/test_docker_preflight.py new file mode 100644 index 0000000000..d9ac0397a3 --- /dev/null +++ b/packages/nemo_platform_ext/tests/cli/test_docker_preflight.py @@ -0,0 +1,209 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Tests for config-aware Docker preflight (NVBug 6537617).""" + +from __future__ import annotations + +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest +import typer +from nemo_platform_ext.cli.docker_preflight import ( + DOCKER_PREFLIGHT_MESSAGE, + default_local_needs_docker, + require_docker_for_default_local, +) +from nemo_platform_plugin.capabilities import ProbeResult +from nmp.platform_runner.config import PlatformAppConfig + + +@pytest.fixture +def stock_local_yaml(tmp_path: Path) -> Path: + path = tmp_path / "local.yaml" + path.write_text( + """ +platform: + runtime: docker +deployments: + executors: + - name: local-docker + backend: docker + config: {} + default_executor: local-docker +""", + encoding="utf-8", + ) + return path + + +@pytest.fixture +def k8s_yaml(tmp_path: Path) -> Path: + path = tmp_path / "k8s.yaml" + path.write_text( + """ +platform: + runtime: kubernetes +deployments: + executors: + - name: local-k8s + backend: kubernetes + config: {} + default_executor: local-k8s +""", + encoding="utf-8", + ) + return path + + +@pytest.fixture +def non_docker_default_yaml(tmp_path: Path) -> Path: + path = tmp_path / "subprocess.yaml" + path.write_text( + """ +platform: + runtime: docker +deployments: + executors: + - name: local-sandbox + backend: sandbox + config: {} + default_executor: local-sandbox +""", + encoding="utf-8", + ) + return path + + +def test_default_local_needs_docker_for_stock_full_platform(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml)) + with patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments", "entities"}, + controllers={"deployments"}, + config_path=str(stock_local_yaml), + ), + ): + assert default_local_needs_docker(cfg) is True + + +def test_default_local_skips_when_deployments_not_selected(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml), services=["entities"], service_group=None) + with patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"entities"}, + controllers=set(), + config_path=str(stock_local_yaml), + ), + ): + assert default_local_needs_docker(cfg) is False + + +def test_default_local_skips_kubernetes_runtime(k8s_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(k8s_yaml)) + with patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(k8s_yaml), + ), + ): + assert default_local_needs_docker(cfg) is False + + +def test_default_local_skips_non_docker_default_executor(non_docker_default_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(non_docker_default_yaml)) + with patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers=set(), + config_path=str(non_docker_default_yaml), + ), + ): + assert default_local_needs_docker(cfg) is False + + +def test_require_docker_exits_without_spawn_when_probe_false(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml)) + with ( + patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(stock_local_yaml), + ), + ), + patch( + "nemo_platform_ext.cli.docker_preflight.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ) as probe, + pytest.raises(typer.Exit) as exc, + ): + require_docker_for_default_local(cfg) + assert exc.value.exit_code == 1 + probe.assert_called_once_with(docker_host=None, use_cache=False) + + +def test_require_docker_noop_when_probe_true(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml)) + with ( + patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(stock_local_yaml), + ), + ), + patch( + "nemo_platform_ext.cli.docker_preflight.probe_docker", + return_value=ProbeResult(available=True), + ), + ): + require_docker_for_default_local(cfg) + + +def test_require_docker_probes_executor_docker_host(tmp_path: Path) -> None: + path = tmp_path / "remote-docker.yaml" + path.write_text( + """ +platform: + runtime: docker +deployments: + executors: + - name: remote-docker + backend: docker + config: + docker_host: tcp://docker.example:2375 + default_executor: remote-docker +""", + encoding="utf-8", + ) + cfg = PlatformAppConfig(config_path=str(path)) + with ( + patch( + "nemo_platform_ext.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(path), + ), + ), + patch( + "nemo_platform_ext.cli.docker_preflight.probe_docker", + return_value=ProbeResult(available=True), + ) as probe, + ): + require_docker_for_default_local(cfg) + probe.assert_called_once_with(docker_host="tcp://docker.example:2375", use_cache=False) + + +def test_preflight_message_names_docker() -> None: + assert "Docker" in DOCKER_PREFLIGHT_MESSAGE + assert "default local" in DOCKER_PREFLIGHT_MESSAGE.lower() or "default_executor" in DOCKER_PREFLIGHT_MESSAGE diff --git a/packages/nemo_platform_plugin/src/nemo_platform_plugin/capabilities.py b/packages/nemo_platform_plugin/src/nemo_platform_plugin/capabilities.py new file mode 100644 index 0000000000..118b5af00b --- /dev/null +++ b/packages/nemo_platform_plugin/src/nemo_platform_plugin/capabilities.py @@ -0,0 +1,136 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Shared backend capability probes for NeMo Platform. + +Runtime (``platform.runtime``) describes *where the platform process runs*. +Capability probes answer *which backends can run jobs/deployments right now*. + +This module owns the Docker probe used by jobs, deployments, setup, and +customization. GPU/Kubernetes probes are deferred to AIRCORE-972. + +Caching +------- +Results are memoized per Docker endpoint for long-lived server processes +(registry boot). CLI / preflight paths that need a fresh verdict should pass +``use_cache=False``. Tests that construct Docker backends should call +:func:`reset_capability_cache` so a prior miss does not poison later fixtures. +""" + +from __future__ import annotations + +import contextlib +import logging +import os +from dataclasses import dataclass + +logger = logging.getLogger(__name__) + +_DOCKER_PROBE_TIMEOUT_SECONDS = 5 + + +class CapabilityUnavailableError(RuntimeError): + """A required backend capability is unavailable. + + Raised when an optional packaging extra is missing or a runtime substrate + (e.g. Docker daemon/socket) is unreachable. Registries may soft-skip + optional backends that raise this, while still fail-closing when a + configured default still names the unavailable backend. + """ + + +@dataclass(frozen=True) +class ProbeResult: + """Outcome of a capability probe.""" + + available: bool + detail: str | None = None + + +# Cache keyed by docker host (None → default DOCKER_HOST / from_env). +_docker_probe_cache: dict[str | None, ProbeResult] = {} + + +def reset_capability_cache() -> None: + """Clear memoized probe results. + + Call from test fixtures so a prior unavailable verdict does not pin the + process. CLI helpers that need a fresh probe should prefer + ``probe_docker(use_cache=False)`` instead of resetting the whole cache. + """ + _docker_probe_cache.clear() + + +def probe_docker( + *, + docker_host: str | None = None, + use_cache: bool = True, +) -> ProbeResult: + """Probe Docker daemon reachability with a short timeout. + + Args: + docker_host: Optional Docker API URL override (same meaning as + ``DOCKER_HOST`` / deployments ``docker_host``). ``None`` uses the + environment default via ``docker.from_env``. Passed by setting + ``DOCKER_HOST`` in the env dict — docker-py 7.x rejects + ``base_url=`` on ``from_env``. + use_cache: When True (default), reuse a prior result for this host key. + Pass False for CLI retry UX after the user starts Docker. + """ + cache_key = docker_host + if use_cache and cache_key in _docker_probe_cache: + return _docker_probe_cache[cache_key] + + result = _probe_docker_uncached(docker_host=docker_host) + if use_cache: + _docker_probe_cache[cache_key] = result + return result + + +def docker_from_env_kwargs(*, timeout: float, docker_host: str | None = None) -> dict[str, object]: + """Build kwargs for ``docker.from_env`` that docker-py 7.x accepts. + + When ``docker_host`` is set, override ``DOCKER_HOST`` via an ``environment`` + copy. Leave the process environment untouched when no host override is + configured (default socket / existing ``DOCKER_HOST`` still apply). + """ + kwargs: dict[str, object] = {"timeout": timeout} + if docker_host: + kwargs["environment"] = {**os.environ, "DOCKER_HOST": docker_host} + return kwargs + + +def _probe_docker_uncached(*, docker_host: str | None) -> ProbeResult: + try: + from docker.errors import DockerException + from requests.exceptions import ConnectionError as RequestsConnectionError + from requests.exceptions import Timeout as RequestsTimeout + + import docker + except ImportError as exc: + detail = f"Docker Python package is not installed ({exc})" + logger.debug(detail) + return ProbeResult(available=False, detail=detail) + + client = None + try: + client = docker.from_env( + **docker_from_env_kwargs(timeout=_DOCKER_PROBE_TIMEOUT_SECONDS, docker_host=docker_host) + ) + client.ping() + return ProbeResult(available=True, detail=None) + except (DockerException, RequestsConnectionError, RequestsTimeout, OSError) as exc: + detail = f"Docker daemon unreachable ({exc})" + logger.debug(detail) + return ProbeResult(available=False, detail=detail) + finally: + if client is not None: + with contextlib.suppress(Exception): + client.close() + + +def require_docker(*, docker_host: str | None = None, use_cache: bool = True) -> None: + """Raise :class:`CapabilityUnavailableError` when Docker is not available.""" + result = probe_docker(docker_host=docker_host, use_cache=use_cache) + if not result.available: + raise CapabilityUnavailableError(result.detail or "Docker daemon is unavailable") diff --git a/packages/nemo_platform_plugin/src/nemo_platform_plugin/config.py b/packages/nemo_platform_plugin/src/nemo_platform_plugin/config.py index f638c83afe..08ac076a5a 100644 --- a/packages/nemo_platform_plugin/src/nemo_platform_plugin/config.py +++ b/packages/nemo_platform_plugin/src/nemo_platform_plugin/config.py @@ -78,14 +78,10 @@ def test_something(): from typing import Any, ClassVar, Literal, Self, Type, TypeVar import yaml -from docker.errors import DockerException +from nemo_platform_plugin.capabilities import probe_docker from pydantic import BaseModel, Field, field_validator, model_validator from pydantic._internal._model_construction import ModelMetaclass from pydantic_settings import BaseSettings, PydanticBaseSettingsSource, SettingsConfigDict -from requests.exceptions import ConnectionError as RequestsConnectionError -from requests.exceptions import Timeout as RequestsTimeout - -import docker logger = logging.getLogger(__name__) @@ -383,17 +379,19 @@ def determine_loopback_override() -> str | None: def validate_docker_available() -> bool: - """Validate that Docker is available using a lightweight check.""" - client = None - try: - client = docker.from_env(timeout=5) - client.ping() - return True - except (DockerException, RequestsConnectionError, RequestsTimeout, OSError): - return False - finally: - if client: - client.close() + """Validate Docker reachability with an uncached probe (≤5s). + + Thin wrapper over :func:`nemo_platform_plugin.capabilities.probe_docker`. + Prefer ``probe_docker`` when callers need the failure detail or a + ``docker_host`` override. + + Uses ``use_cache=False`` so :meth:`NemoPlatformConfig.validate_runtime` + soft-downgrade does not pin a Docker-unavailable verdict into the + process-wide cache before jobs/deployments registry construction. + Long-lived server paths that want memoization should call + ``probe_docker()`` directly. + """ + return probe_docker(use_cache=False).available class ImagePullSecret(BaseSettings): @@ -653,7 +651,14 @@ def validate_ngc_api_key_secret(self) -> Self: def validate_runtime(self) -> Self: if self.runtime == Runtime.DOCKER: if not validate_docker_available(): - logger.warning("Docker is not available, setting runtime to NONE") + # Deprecated convenience: Runtime is topology, not capability. + # Capability probes (nemo_platform_plugin.capabilities) own Docker + # availability. Soft-downgrade remains for one release; AIRCORE-972 + # removes or shrinks Runtime.NONE as a Docker-absence signal. + logger.warning( + "Docker is not available, setting runtime to NONE " + "(deprecated: prefer capability probes; see AIRCORE-972)" + ) self.runtime = Runtime.NONE return self diff --git a/packages/nemo_platform_plugin/src/nemo_platform_plugin/jobs/docker.py b/packages/nemo_platform_plugin/src/nemo_platform_plugin/jobs/docker.py index a6fcff748e..8805d68a4c 100644 --- a/packages/nemo_platform_plugin/src/nemo_platform_plugin/jobs/docker.py +++ b/packages/nemo_platform_plugin/src/nemo_platform_plugin/jobs/docker.py @@ -6,6 +6,7 @@ import logging from nemo_platform.types.jobs import PlatformJobSpecParam +from nemo_platform_plugin.capabilities import probe_docker from nemo_platform_plugin.config import Configuration, NemoPlatformConfig, Runtime from nemo_platform_plugin.jobs.exceptions import PlatformJobCompilationError from pydantic import ValidationError @@ -36,8 +37,22 @@ def spec_has_gpu_step(job: PlatformJobSpecParam) -> bool: return False +def _docker_gpu_validation_applies(runtime: Runtime) -> bool: + """Whether reserved-GPU checks apply for this platform runtime. + + ``Runtime.DOCKER`` always applies. Soft-downgraded ``Runtime.NONE`` still + runs Docker-backed jobs when the daemon is reachable, so reserved-GPU + checks must not be skipped solely because runtime is not DOCKER. + """ + if runtime == Runtime.DOCKER: + return True + if runtime == Runtime.NONE: + return probe_docker().available + return False + + def validate_gpu_available_for_docker(job: PlatformJobSpecParam) -> None: - """Fail fast when job requires GPU but platform is Docker with no GPUs configured. + """Fail fast when job requires GPU but platform Docker has no GPUs configured. platform_config.docker.get_reserved_gpu_ids() returns: - None when reserved_gpu_device_ids is "all" (auto-detect GPUs) @@ -58,7 +73,7 @@ def validate_gpu_available_for_docker(job: PlatformJobSpecParam) -> None: logger.debug("Skipping GPU availability validation: could not load platform config") return - if platform_config.runtime != Runtime.DOCKER: + if not _docker_gpu_validation_applies(platform_config.runtime): return # None = "all" (auto-detect), [] = "none"/empty (no GPUs), list[int] = explicit IDs. diff --git a/packages/nemo_platform_plugin/tests/test_capabilities.py b/packages/nemo_platform_plugin/tests/test_capabilities.py new file mode 100644 index 0000000000..3ab1347c64 --- /dev/null +++ b/packages/nemo_platform_plugin/tests/test_capabilities.py @@ -0,0 +1,138 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Tests for nemo_platform_plugin.capabilities — shared Docker capability probe.""" + +from __future__ import annotations + +from collections.abc import Iterator +from unittest.mock import MagicMock, patch + +import pytest +from docker.errors import DockerException +from nemo_platform_plugin.capabilities import ( + CapabilityUnavailableError, + probe_docker, + require_docker, + reset_capability_cache, +) +from requests.exceptions import ConnectionError as RequestsConnectionError +from requests.exceptions import Timeout as RequestsTimeout + + +@pytest.fixture(autouse=True) +def _clear_capability_cache() -> Iterator[None]: + reset_capability_cache() + yield + reset_capability_cache() + + +def test_probe_docker_available() -> None: + client = MagicMock() + client.ping.return_value = True + with patch("docker.from_env", return_value=client) as from_env: + result = probe_docker() + assert result.available is True + assert result.detail is None + from_env.assert_called_once_with(timeout=5) + client.ping.assert_called_once() + client.close.assert_called_once() + + +@pytest.mark.parametrize( + "exc", + [ + DockerException("boom"), + RequestsConnectionError("refused"), + RequestsTimeout("timed out"), + OSError("no such file"), + ], +) +def test_probe_docker_unavailable_on_connection_failures(exc: Exception) -> None: + with patch("docker.from_env", side_effect=exc): + result = probe_docker() + assert result.available is False + assert result.detail is not None + assert "unreachable" in result.detail.lower() or "Docker" in result.detail + + +def test_probe_docker_unavailable_on_ping_failure() -> None: + client = MagicMock() + client.ping.side_effect = DockerException("ping failed") + with patch("docker.from_env", return_value=client): + result = probe_docker() + assert result.available is False + client.close.assert_called_once() + + +def test_probe_docker_caches_result() -> None: + client = MagicMock() + with patch("docker.from_env", return_value=client) as from_env: + first = probe_docker() + second = probe_docker() + assert first.available is True + assert second.available is True + assert from_env.call_count == 1 + + +def test_reset_capability_cache_allows_reprobe() -> None: + client = MagicMock() + with patch("docker.from_env", return_value=client) as from_env: + probe_docker() + reset_capability_cache() + probe_docker() + assert from_env.call_count == 2 + + +def test_probe_docker_use_cache_false_bypasses_memo() -> None: + client = MagicMock() + with patch("docker.from_env", return_value=client) as from_env: + probe_docker() + probe_docker(use_cache=False) + assert from_env.call_count == 2 + + +def test_probe_docker_host_keys_are_independent() -> None: + default_client = MagicMock() + remote_client = MagicMock() + remote_client.ping.side_effect = DockerException("remote down") + + def _from_env(**kwargs): + env = kwargs.get("environment") or {} + if env.get("DOCKER_HOST") == "tcp://remote:2375": + return remote_client + return default_client + + with patch("docker.from_env", side_effect=_from_env) as from_env: + default = probe_docker() + remote = probe_docker(docker_host="tcp://remote:2375") + + assert default.available is True + assert remote.available is False + assert from_env.call_count == 2 + from_env.assert_any_call(timeout=5) + remote_call = next(c for c in from_env.call_args_list if c.kwargs.get("environment")) + assert remote_call.kwargs["timeout"] == 5 + assert remote_call.kwargs["environment"]["DOCKER_HOST"] == "tcp://remote:2375" + + +def test_probe_docker_explicit_host_returns_unavailable_without_typeerror() -> None: + """Regression: docker-py 7.x rejects base_url= on from_env; must not raise TypeError.""" + result = probe_docker(docker_host="tcp://127.0.0.1:1", use_cache=False) + assert result.available is False + assert result.detail is not None + + +def test_require_docker_raises_when_unavailable() -> None: + with patch("docker.from_env", side_effect=DockerException("down")): + with pytest.raises(CapabilityUnavailableError, match="unreachable"): + require_docker() + + +def test_validate_docker_available_delegates_to_probe() -> None: + from nemo_platform_plugin.config import validate_docker_available + + with patch("nemo_platform_plugin.config.probe_docker") as probe: + probe.return_value = MagicMock(available=False) + assert validate_docker_available() is False + probe.assert_called_once() diff --git a/packages/nemo_platform_plugin/tests/test_config.py b/packages/nemo_platform_plugin/tests/test_config.py index 995836aa16..f24d4238fb 100644 --- a/packages/nemo_platform_plugin/tests/test_config.py +++ b/packages/nemo_platform_plugin/tests/test_config.py @@ -274,6 +274,7 @@ def test_validate_docker_available_returns_false_on_connection_failures() -> Non from unittest.mock import MagicMock, patch from docker.errors import DockerException + from nemo_platform_plugin.capabilities import reset_capability_cache from nemo_platform_plugin.config import validate_docker_available from requests.exceptions import ConnectionError as RequestsConnectionError @@ -282,11 +283,14 @@ def test_validate_docker_available_returns_false_on_connection_failures() -> Non RequestsConnectionError("refused"), OSError("no such file"), ): - with patch("nemo_platform_plugin.config.docker.from_env", side_effect=exc): + reset_capability_cache() + with patch("docker.from_env", side_effect=exc): assert validate_docker_available() is False + reset_capability_cache() client = MagicMock() client.ping.side_effect = exc - with patch("nemo_platform_plugin.config.docker.from_env", return_value=client): + with patch("docker.from_env", return_value=client): assert validate_docker_available() is False client.close.assert_called() + reset_capability_cache() diff --git a/packages/nmp_common/src/nmp/common/jobs/docker.py b/packages/nmp_common/src/nmp/common/jobs/docker.py index 9be4887aa2..5c4cd326d8 100644 --- a/packages/nmp_common/src/nmp/common/jobs/docker.py +++ b/packages/nmp_common/src/nmp/common/jobs/docker.py @@ -3,40 +3,16 @@ """Docker-specific job validation compatibility wrapper.""" -import logging - from nemo_platform.types.jobs import PlatformJobSpecParam from nemo_platform_plugin.jobs.docker import spec_has_gpu_step as spec_has_gpu_step -from nmp.common.config import Runtime, get_platform_config -from nmp.common.jobs.exceptions import PlatformJobCompilationError -from pydantic import ValidationError - -logger = logging.getLogger(__name__) - -_GPU_UNAVAILABLE_MSG = ( - "This job requires a GPU. The platform is running on Docker with no GPUs " - "configured. Configure GPUs in platform config or use a GPU-enabled environment." -) +from nemo_platform_plugin.jobs.docker import validate_gpu_available_for_docker as _plugin_validate def validate_gpu_available_for_docker(job: PlatformJobSpecParam) -> None: - """Fail fast when job requires GPU but platform is Docker with no GPUs configured.""" - try: - platform_config = get_platform_config() - except (ValueError, OSError, ValidationError): - logger.debug("Skipping GPU availability validation: could not load platform config") - return - - if platform_config.runtime != Runtime.DOCKER: - return - - reserved_ids = platform_config.docker.get_reserved_gpu_ids() - if reserved_ids is None: - return - if len(reserved_ids) > 0: - return - - if not spec_has_gpu_step(job): - return + """Fail fast when job requires GPU but platform Docker has no GPUs configured. - raise PlatformJobCompilationError(_GPU_UNAVAILABLE_MSG) + Delegates to :func:`nemo_platform_plugin.jobs.docker.validate_gpu_available_for_docker` + so soft-downgraded ``Runtime.NONE`` with a reachable Docker daemon still + enforces reserved-GPU checks (AIRCORE-971). + """ + _plugin_validate(job) diff --git a/packages/nmp_common/tests/jobs/test_docker.py b/packages/nmp_common/tests/jobs/test_docker.py index 60b54000b0..907dcd160d 100644 --- a/packages/nmp_common/tests/jobs/test_docker.py +++ b/packages/nmp_common/tests/jobs/test_docker.py @@ -6,6 +6,7 @@ from unittest.mock import MagicMock, patch import pytest +from nemo_platform_plugin.config import Runtime from nemo_platform_plugin.jobs.api_factory import ( ContainerSpec, CPUExecutionProviderSpec, @@ -13,7 +14,6 @@ PlatformJobSpec, PlatformJobStep, ) -from nmp.common.config import Runtime from nmp.common.jobs.docker import spec_has_gpu_step, validate_gpu_available_for_docker from nmp.common.jobs.exceptions import PlatformJobCompilationError @@ -50,16 +50,27 @@ def test_spec_has_gpu_step(): @pytest.mark.parametrize( - "runtime,reserved_gpu_ids,config_raises,expect_raise,message_contains", + "runtime,reserved_gpu_ids,config_raises,docker_available,expect_raise,message_contains", [ - (Runtime.DOCKER, [], False, True, ("no GPUs configured", "Docker")), - (Runtime.DOCKER, [0, 1], False, False, ()), - (Runtime.KUBERNETES, [], False, False, ()), - (None, None, True, False, ()), + (Runtime.DOCKER, [], False, True, True, ("no GPUs configured", "Docker")), + (Runtime.DOCKER, [0, 1], False, True, False, ()), + (Runtime.KUBERNETES, [], False, True, False, ()), + (Runtime.NONE, [], False, True, True, ("no GPUs configured", "Docker")), + (Runtime.NONE, [], False, False, False, ()), + (None, None, True, True, False, ()), + ], + ids=[ + "docker_no_gpus_raises", + "docker_has_gpus_passes", + "kubernetes_passes", + "none_with_docker_no_gpus_raises", + "none_without_docker_skips", + "config_fails_skips", ], - ids=["docker_no_gpus_raises", "docker_has_gpus_passes", "kubernetes_passes", "config_fails_skips"], ) -def test_validate_gpu_available_for_docker(runtime, reserved_gpu_ids, config_raises, expect_raise, message_contains): +def test_validate_gpu_available_for_docker( + runtime, reserved_gpu_ids, config_raises, docker_available, expect_raise, message_contains +): """GPU job validation: raise when Docker has no GPUs; pass or skip otherwise.""" gpu_executor = GPUExecutionProviderSpec( provider="gpu", profile="default", container=ContainerSpec(image="gpu_image") @@ -70,14 +81,26 @@ def test_validate_gpu_available_for_docker(runtime, reserved_gpu_ids, config_rai ] ) if config_raises: - get_config = patch("nmp.common.jobs.docker.get_platform_config", side_effect=ValueError("no config")) + get_config = patch( + "nemo_platform_plugin.jobs.docker.get_platform_config", + side_effect=ValueError("no config"), + ) else: mock_platform_config = MagicMock() mock_platform_config.runtime = runtime mock_platform_config.docker.get_reserved_gpu_ids.return_value = reserved_gpu_ids - get_config = patch("nmp.common.jobs.docker.get_platform_config", return_value=mock_platform_config) + get_config = patch( + "nemo_platform_plugin.jobs.docker.get_platform_config", + return_value=mock_platform_config, + ) - with get_config: + with ( + get_config, + patch( + "nemo_platform_plugin.jobs.docker.probe_docker", + return_value=MagicMock(available=docker_available), + ), + ): if expect_raise: with pytest.raises(PlatformJobCompilationError) as exc_info: validate_gpu_available_for_docker(gpu_job) diff --git a/packages/nmp_customization_common/src/nmp/customization_common/contributor/jobs.py b/packages/nmp_customization_common/src/nmp/customization_common/contributor/jobs.py index 83da9ead9b..1a187b6e5c 100644 --- a/packages/nmp_customization_common/src/nmp/customization_common/contributor/jobs.py +++ b/packages/nmp_customization_common/src/nmp/customization_common/contributor/jobs.py @@ -14,6 +14,7 @@ from typing import ClassVar, cast from nemo_platform import AsyncNeMoPlatform +from nemo_platform_plugin.capabilities import probe_docker from nemo_platform_plugin.config import NemoPlatformConfig, Runtime from nemo_platform_plugin.job import NemoJob from nemo_platform_plugin.jobs.exceptions import PlatformJobCompilationError @@ -28,15 +29,17 @@ def require_container_runtime(backend_label: str, *, num_nodes: int = 1) -> None - **Kubernetes** — the platform schedules GPU pods, including multi-node ``gpu_distributed`` jobs via Volcano. - - **Docker** — the platform's local Docker GPU executor (single host). - - Single-node jobs accept either runtime. **Multi-node jobs** (``num_nodes > - 1``) compile to a ``gpu_distributed`` executor that only the Volcano - (Kubernetes) backend can place — Docker has no multi-node/``gpu_distributed`` - backend — so they require ``platform.runtime: kubernetes``. Failing here - surfaces the misconfiguration at compile time instead of as an opaque - "no backend found" scheduling error (or, for ``runtime: none``, before the - Jobs API rejects the spec). + - **Docker** — the platform's local Docker GPU executor (single host), + detected via :func:`~nemo_platform_plugin.capabilities.probe_docker` + rather than treating ``Runtime.DOCKER`` / ``Runtime.NONE`` as the + capability signal (AIRCORE-971). + + Single-node jobs accept Kubernetes topology or a reachable Docker daemon. + **Multi-node jobs** (``num_nodes > 1``) compile to a ``gpu_distributed`` + executor that only the Volcano (Kubernetes) backend can place — Docker has + no multi-node/``gpu_distributed`` backend — so they require + ``platform.runtime: kubernetes``. Failing here surfaces the misconfiguration at + compile time instead of as an opaque "no backend found" scheduling error. """ platform_config = NemoPlatformConfig.get() runtime = platform_config.runtime @@ -52,19 +55,16 @@ def require_container_runtime(backend_label: str, *, num_nodes: int = 1) -> None if runtime == Runtime.KUBERNETES: return - if runtime == Runtime.DOCKER: - from nemo_platform_plugin.config import validate_docker_available - - if not validate_docker_available(): - raise PlatformJobCompilationError( - f"{backend_label} training requires a reachable Docker daemon (platform.runtime: docker).", - ) + # Capability probe — independent of Runtime.DOCKER vs soft-downgraded NONE. + # Use the process cache so compile agrees with the jobs registry boot probe. + # Mid-process "start Docker then retry compile" requires a platform restart. + result = probe_docker() + if result.available: return + detail = result.detail or "Docker daemon is unavailable" raise PlatformJobCompilationError( - f"{backend_label} training requires a container runtime: set platform.runtime to " - "'kubernetes' (schedules GPU pods) or 'docker' (local GPU executor). " - f"Current runtime: {runtime.value}.", + f"{backend_label} training requires a reachable Docker daemon (or platform.runtime: kubernetes). {detail}", ) diff --git a/packages/nmp_customization_common/tests/contributor/test_jobs.py b/packages/nmp_customization_common/tests/contributor/test_jobs.py index b2bce9f5ed..12ff2ea1ca 100644 --- a/packages/nmp_customization_common/tests/contributor/test_jobs.py +++ b/packages/nmp_customization_common/tests/contributor/test_jobs.py @@ -5,7 +5,7 @@ Covers :func:`require_container_runtime` (automodel / unsloth) and :func:`require_distributed_runtime` (rl). We patch the platform config and the -Docker-availability probe so the checks are exercised without a real runtime. +Docker capability probe so the checks are exercised without a real runtime. """ from __future__ import annotations @@ -13,6 +13,7 @@ from types import SimpleNamespace import pytest +from nemo_platform_plugin.capabilities import ProbeResult from nemo_platform_plugin.config import Runtime from nemo_platform_plugin.jobs.exceptions import PlatformJobCompilationError from nmp.customization_common.contributor import jobs as jobs_mod @@ -33,8 +34,12 @@ def _apply(runtime: Runtime, *, docker_available: bool = True) -> None: classmethod(lambda cls: SimpleNamespace(runtime=runtime)), ) monkeypatch.setattr( - "nemo_platform_plugin.config.validate_docker_available", - lambda: docker_available, + jobs_mod, + "probe_docker", + lambda **kwargs: ProbeResult( + available=docker_available, + detail=None if docker_available else "Docker daemon unreachable (test)", + ), ) return _apply @@ -63,9 +68,14 @@ def test_docker_multi_node_requires_kubernetes(self, _patch_runtime) -> None: with pytest.raises(PlatformJobCompilationError, match="multi-node training .* requires"): require_container_runtime("Automodel", num_nodes=2) - def test_none_runtime_raises(self, _patch_runtime) -> None: - _patch_runtime(Runtime.NONE) - with pytest.raises(PlatformJobCompilationError, match="requires a container runtime"): + def test_none_runtime_with_docker_ok(self, _patch_runtime) -> None: + """Soft-downgraded NONE still allows compile when Docker is reachable.""" + _patch_runtime(Runtime.NONE, docker_available=True) + require_container_runtime("Automodel") # no raise + + def test_none_runtime_without_docker_raises(self, _patch_runtime) -> None: + _patch_runtime(Runtime.NONE, docker_available=False) + with pytest.raises(PlatformJobCompilationError, match="reachable Docker daemon"): require_container_runtime("Automodel") def test_none_runtime_multi_node_raises_kubernetes_message(self, _patch_runtime) -> None: diff --git a/plugins/nemo-automodel/src/nemo_automodel_plugin/jobs/jobs.py b/plugins/nemo-automodel/src/nemo_automodel_plugin/jobs/jobs.py index 85fb6f11a0..e9d05c3a35 100644 --- a/plugins/nemo-automodel/src/nemo_automodel_plugin/jobs/jobs.py +++ b/plugins/nemo-automodel/src/nemo_automodel_plugin/jobs/jobs.py @@ -11,6 +11,7 @@ from __future__ import annotations +import asyncio from typing import ClassVar, cast from nemo_automodel_plugin.config import get_config @@ -51,12 +52,19 @@ async def compile( options: dict | None = None, ) -> PlatformJobSpec: del entity_client, options + if not isinstance(async_sdk, AsyncNeMoPlatform): + raise TypeError(f"async_sdk must be AsyncNeMoPlatform, got {type(async_sdk).__name__}") canonical = ( spec if isinstance(spec, AutomodelJobOutput) else AutomodelJobOutput.model_validate(spec.model_dump()) ) # Multi-node jobs compile to a gpu_distributed (Volcano) executor, which # only exists on Kubernetes; gate here so docker platforms fail fast. - require_container_runtime(cls.runtime_label, num_nodes=canonical.parallelism.num_nodes) + # Probe is sync (≤5s); keep it off the event loop. + await asyncio.to_thread( + require_container_runtime, + cls.runtime_label, + num_nodes=canonical.parallelism.num_nodes, + ) try: canonical.validate_for_training() except ValidationError as e: @@ -70,7 +78,7 @@ async def compile( platform_spec = await platform_job_config_compiler( canonical, workspace, - cast(AsyncNeMoPlatform, async_sdk), + async_sdk, job_name=job_name, profile=execution_profile, ) diff --git a/plugins/nemo-automodel/tests/test_jobs.py b/plugins/nemo-automodel/tests/test_jobs.py index 4fa89d98c5..d168aa55f3 100644 --- a/plugins/nemo-automodel/tests/test_jobs.py +++ b/plugins/nemo-automodel/tests/test_jobs.py @@ -18,6 +18,7 @@ import pytest from nemo_automodel_plugin.jobs.jobs import AutomodelJob from nemo_automodel_plugin.schema import AutomodelJobOutput +from nemo_platform import AsyncNeMoPlatform from nemo_platform_plugin.jobs.exceptions import PlatformJobCompilationError @@ -43,7 +44,7 @@ def _compile(canonical: AutomodelJobOutput) -> Any: spec=canonical, entity_client=object(), job_name=None, - async_sdk=object(), + async_sdk=object.__new__(AsyncNeMoPlatform), ), ) diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/base.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/base.py index f96e5fcb4d..94b3d6c672 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/base.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/base.py @@ -12,17 +12,20 @@ from nemo_deployments_plugin.types import DeploymentStatus, Endpoint, VolumeStatus from nemo_platform import AsyncNeMoPlatform +from nemo_platform_plugin.capabilities import CapabilityUnavailableError from pydantic import BaseModel, Field -class MissingBackendDependencyError(RuntimeError): +class MissingBackendDependencyError(CapabilityUnavailableError): """A backend cannot initialize because a required capability is unavailable. Raised when an optional packaging extra is missing (e.g. ``openshell``) or when a runtime substrate is unreachable (e.g. Docker daemon/socket). The executor registry catches this, skips that executor with a warning, and continues - starting the deployments service. Subclasses ``RuntimeError`` so existing - ``except RuntimeError`` paths keep working. + starting the deployments service. Subclasses + :class:`~nemo_platform_plugin.capabilities.CapabilityUnavailableError` (and + therefore ``RuntimeError``) so existing ``except RuntimeError`` paths keep + working and jobs/deployments share one unavailable-capability type. """ diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/backend.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/backend.py index 4d4462fdb9..641defb6db 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/backend.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/backend.py @@ -60,6 +60,7 @@ from nemo_deployments_plugin.entities import Container, Deployment, DeploymentConfig from nemo_deployments_plugin.secrets import SecretResolutionError, resolve_deployment_config_secrets from nemo_deployments_plugin.types import Endpoint, RestartPolicy +from nemo_platform_plugin.capabilities import docker_from_env_kwargs, probe_docker from nemo_platform_plugin.client.adapter import client_from_platform from nemo_platform_plugin.config import LOOPBACK_ADDRESSES from nemo_platform_plugin.entities.client import AsyncEntitiesClient @@ -114,6 +115,13 @@ def init(self) -> None: self._executor_config = DockerExecutorConfig.model_validate(self._config) self._entities = NemoEntitiesClient(client_from_platform(self._sdk, AsyncEntitiesClient)) self._gpu_pool = get_shared_gpu_pool() + docker_host = self._executor_config.docker_host + probe = probe_docker(docker_host=docker_host) + if not probe.available: + detail = probe.detail or "Docker daemon unreachable" + raise MissingBackendDependencyError( + f"Docker daemon is unavailable ({detail}). Docker-backed deployments will be disabled." + ) try: self._client = self._create_client() except (DockerException, RequestsConnectionError, RequestsTimeout, OSError) as exc: @@ -122,10 +130,14 @@ def init(self) -> None: ) from exc def _create_client(self) -> docker.DockerClient: - kwargs: dict[str, Any] = {"timeout": self._executor_config.docker_timeout} - if self._executor_config.docker_host: - kwargs["base_url"] = self._executor_config.docker_host - client = self._docker.from_env(**kwargs) + # docker-py 7.x rejects base_url= on from_env; override DOCKER_HOST instead + # so TLS env vars (DOCKER_TLS_VERIFY, cert paths) still apply. + client = self._docker.from_env( + **docker_from_env_kwargs( + timeout=self._executor_config.docker_timeout, + docker_host=self._executor_config.docker_host, + ) + ) client.api.timeout = self._executor_config.docker_timeout client.ping() return client diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/registry.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/registry.py index 4f74d7784c..ba9f8080ec 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/registry.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/registry.py @@ -96,9 +96,16 @@ def from_config( exc, ) if default_executor and default_executor not in executors: + default_spec = next((spec for spec in specs if spec.name == default_executor), None) + backend_hint = ( + f"backend '{default_spec.backend}'" + if default_spec is not None + else "its required backend capability" + ) raise ExecutorNotFoundError( f"default_executor '{default_executor}' is not registered " - "(unavailable backend or missing from executor config)." + "(the backend is unavailable or the executor is missing from configuration). " + f"Configure a registered default_executor, or restore {backend_hint}." ) except Exception: for backend in executors.values(): diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/conftest.py b/plugins/nemo-deployments/tests/unit/backends/docker/conftest.py index a3905e7409..5bded7e9b1 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/conftest.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/conftest.py @@ -10,6 +10,15 @@ import pytest from nemo_deployments_plugin.backends.docker.backend import DockerDeploymentBackend +from nemo_platform_plugin.capabilities import reset_capability_cache + + +@pytest.fixture(autouse=True) +def _reset_capability_cache() -> Iterator[None]: + """Prevent probe_docker process cache from poisoning later fixtures.""" + reset_capability_cache() + yield + reset_capability_cache() @pytest.fixture diff --git a/plugins/nemo-deployments/tests/unit/test_registry.py b/plugins/nemo-deployments/tests/unit/test_registry.py index 8262412e8c..b54b85534b 100644 --- a/plugins/nemo-deployments/tests/unit/test_registry.py +++ b/plugins/nemo-deployments/tests/unit/test_registry.py @@ -214,7 +214,10 @@ def test_registry_missing_dependency_default_executor_fails( # register must fail fast so misconfiguration is obvious at startup. sdk = AsyncNeMoPlatform(base_url="http://localhost:8080") classes = {**backend_classes, "sandbox": _MissingDepBackend} - with pytest.raises(ExecutorNotFoundError, match="default_executor 'sandbox-local' is not registered"): + with pytest.raises( + ExecutorNotFoundError, + match=r"default_executor 'sandbox-local' is not registered.*backend 'sandbox'", + ): ExecutorRegistry.from_config( sdk, [ExecutorSpec(name="sandbox-local", backend="sandbox", config={})], @@ -256,7 +259,10 @@ def init(self) -> None: raise MissingBackendDependencyError("Docker daemon is unavailable") classes = {**backend_classes, "docker": _UnavailableDocker} - with pytest.raises(ExecutorNotFoundError, match="default_executor 'local-docker' is not registered"): + with pytest.raises( + ExecutorNotFoundError, + match=r"default_executor 'local-docker' is not registered.*backend 'docker'", + ): ExecutorRegistry.from_config( sdk, [ diff --git a/plugins/nemo-unsloth/src/nemo_unsloth_plugin/jobs/jobs.py b/plugins/nemo-unsloth/src/nemo_unsloth_plugin/jobs/jobs.py index 196d651a5f..8d7cb35f16 100644 --- a/plugins/nemo-unsloth/src/nemo_unsloth_plugin/jobs/jobs.py +++ b/plugins/nemo-unsloth/src/nemo_unsloth_plugin/jobs/jobs.py @@ -13,6 +13,7 @@ from __future__ import annotations +import asyncio from typing import ClassVar, cast from nemo_platform import AsyncNeMoPlatform @@ -59,7 +60,8 @@ async def compile( ``unsloth_config.default_training_execution_profile``. """ del entity_client, options - require_container_runtime(cls.runtime_label) + # Probe is sync (≤5s); keep it off the event loop. + await asyncio.to_thread(require_container_runtime, cls.runtime_label) canonical = spec if isinstance(spec, UnslothJobOutput) else UnslothJobOutput.model_validate(spec.model_dump()) execution_profile = profile or unsloth_config.default_training_execution_profile diff --git a/sdk/python/nemo-platform/.nmpcontext/openapi.yaml b/sdk/python/nemo-platform/.nmpcontext/openapi.yaml index 5376947a7e..33a0d3c022 100644 --- a/sdk/python/nemo-platform/.nmpcontext/openapi.yaml +++ b/sdk/python/nemo-platform/.nmpcontext/openapi.yaml @@ -4907,7 +4907,16 @@ paths: tags: - Jobs summary: Get Execution Profiles - description: Get all currently configured execution profiles. + description: 'Get all currently configured execution profiles. + + + Returns the capability-filtered merge from jobs config. In local standalone + + the controller may prune the shared list further after registry boot; in + + split topologies the API advertises its own merge result (not controller + + process memory).' operationId: get_execution_profiles_apis_jobs_v2_execution_profiles_get responses: '200': diff --git a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/services/cli.py b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/services/cli.py index f3855c4c39..c63b6bf9db 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/services/cli.py +++ b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/services/cli.py @@ -14,6 +14,7 @@ import httpx import typer from nemo_platform.cli.core.help_formatter import create_typer_app +from nemo_platform.cli.docker_preflight import require_docker_for_default_local from nemo_platform.local.process import ( ForegroundInstanceError, InstanceAlreadyRunningError, @@ -236,16 +237,6 @@ def run_services( _ensure_port_available(host, port, scope, base_dir=base_dir) - try: - lock_fd = acquire_lock(scope, base_dir=base_dir) - except InstanceAlreadyRunningError: - _fail_already_running(scope, base_dir) - - # _NMP_LAUNCH_MODE is set by start_background() when this process was - # spawned via ``nemo services start``. Without it we default to - # "foreground", which protects interactive ``run`` sessions from being - # killed by ``stop``. - mode = "background" if os.environ.get("_NMP_LAUNCH_MODE") == "background" else "foreground" platform_config = PlatformAppConfig( services=_parse_csv_option(services), service_group=service_group, @@ -259,6 +250,18 @@ def run_services( keep_alive_timeout_seconds=keep_alive_timeout_seconds, state_root=base_dir, ) + require_docker_for_default_local(platform_config) + + try: + lock_fd = acquire_lock(scope, base_dir=base_dir) + except InstanceAlreadyRunningError: + _fail_already_running(scope, base_dir) + + # _NMP_LAUNCH_MODE is set by start_background() when this process was + # spawned via ``nemo services start``. Without it we default to + # "foreground", which protects interactive ``run`` sessions from being + # killed by ``stop``. + mode = "background" if os.environ.get("_NMP_LAUNCH_MODE") == "background" else "foreground" desc = InstanceDescriptor.from_config( platform_config, @@ -382,6 +385,7 @@ def start_services( keep_alive_timeout_seconds=keep_alive_timeout_seconds, state_root=base_dir, ) + require_docker_for_default_local(platform_config) typer.echo("Starting platform services...") proc = start_background(platform_config) @@ -565,11 +569,6 @@ def restart_services( ) raise typer.Exit(1) - typer.echo("Stopping platform services...") - # restart always produces a background instance, so force=True is - # appropriate even for foreground targets. - stop_instance(scope, base_dir=base_dir, force=True) - previous_config = prev.config if prev else None effective_services = _parse_csv_option(services) if services is not None else None if services is None and previous_config is not None: @@ -601,9 +600,6 @@ def restart_services( ) ) - _warn_bind_all(effective_host) - - _ensure_port_available(effective_host, effective_port, scope, base_dir=base_dir) platform_config = PlatformAppConfig( services=effective_services, service_group=effective_service_group, @@ -617,6 +613,17 @@ def restart_services( keep_alive_timeout_seconds=effective_keep_alive_timeout_seconds, state_root=base_dir, ) + # Preflight before stop so a missing Docker daemon does not tear down a healthy instance. + require_docker_for_default_local(platform_config) + + typer.echo("Stopping platform services...") + # restart always produces a background instance, so force=True is + # appropriate even for foreground targets. + stop_instance(scope, base_dir=base_dir, force=True) + + _warn_bind_all(effective_host) + + _ensure_port_available(effective_host, effective_port, scope, base_dir=base_dir) typer.echo("Starting platform services...") proc = start_background(platform_config) diff --git a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/setup.py b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/setup.py index 295f963a51..beda9ede74 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/setup.py +++ b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/setup.py @@ -27,8 +27,8 @@ import typer import yaml as _yaml from nemo_platform import NeMoPlatform +from nemo_platform_plugin.capabilities import probe_docker from nemo_platform_plugin.client.adapter import client_from_platform -from nemo_platform_plugin.config import validate_docker_available from nemo_platform_plugin.secrets.client import SecretsClient from nemo_platform_plugin.secrets.types import PlatformSecretCreateRequest, PlatformSecretUpdateRequest from nmp.common.config import nmp_user_data_dir @@ -43,6 +43,7 @@ from nemo_platform.cli.commands.skills.registry import get_installer, load_skills from nemo_platform.cli.core.context import CLIContext from nemo_platform.cli.core.errors import handle_errors +from nemo_platform.cli.docker_preflight import DOCKER_PREFLIGHT_MESSAGE, require_docker_for_default_local from nemo_platform.cli.telemetry import emit from nemo_platform.cli.telemetry.events import OnboardingStepEvent, TaskStatusEnum from nemo_platform.client.tls import client_verify_from_env @@ -1057,7 +1058,7 @@ def _should_hint_docker_unavailable(*, exit_code: int | None, log_path: Path | N """ if _services_log_suggests_docker_failure(log_path): return True - if exit_code is not None and not validate_docker_available(): + if exit_code is not None and not probe_docker(use_cache=False).available: return True return False @@ -1158,6 +1159,9 @@ def _maybe_start_services( console.print(" [cyan]PYO3_USE_ABI3_FORWARD_COMPATIBILITY=1 pip install 'nemo-platform\\[all]'[/cyan]") raise typer.Exit(1) + # Fail before stop/spawn when default local needs Docker (NVBug 6537617). + require_docker_for_default_local(console=console) + if already_running: console.print(" Restarting platform services...") _kill_existing_services(base_url) @@ -1186,10 +1190,7 @@ def _maybe_start_services( console.print(f"{CROSS} Platform did not become ready within {timeout}s") console.print(f" Check {log} for details.") if _should_hint_docker_unavailable(exit_code=exit_code, log_path=log): - console.print( - " Docker does not appear to be available. " - "Install and start Docker, or configure non-Docker executors, then retry." - ) + console.print(f" {DOCKER_PREFLIGHT_MESSAGE}") raise typer.Exit(1) console.print(f"{CHECK} Platform running at {base_url} (pid {proc.pid})\n") diff --git a/sdk/python/nemo-platform/src/nemo_platform/cli/docker_preflight.py b/sdk/python/nemo-platform/src/nemo_platform/cli/docker_preflight.py new file mode 100644 index 0000000000..c57ffc8c1c --- /dev/null +++ b/sdk/python/nemo-platform/src/nemo_platform/cli/docker_preflight.py @@ -0,0 +1,147 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Config-aware Docker preflight for local platform startup (NVBug 6537617). + +Fails fast before spawn/wait when the resolved run would start the deployments +service/controller with a docker-backed ``default_executor`` while the daemon is +unreachable. Does not treat Docker as a global platform dependency — kubernetes and +reduced selections that never hit that fail-close path are left alone. +""" + +from __future__ import annotations + +from dataclasses import dataclass, replace +from pathlib import Path +from typing import Any + +import typer +import yaml +from nemo_platform_plugin.capabilities import probe_docker +from nmp.platform_runner.config import PlatformAppConfig, default_config_path, resolve_run_configuration +from rich.console import Console + +DOCKER_PREFLIGHT_MESSAGE = ( + "Docker is required for this default local setup (deployments default_executor " + "uses the docker backend) but the Docker daemon is not available. " + "Install and start Docker, or use a kubernetes / non-docker config " + "(and omit the deployments service if you do not need it), then retry." +) + + +@dataclass(frozen=True) +class _DefaultDockerExecutorProbe: + """Whether the default deployments executor is docker, and its optional host override.""" + + is_docker: bool + docker_host: str | None = None + + +def _load_platform_yaml(config_path: str) -> dict[str, Any]: + """Load platform YAML without running Pydantic validators (no soft-downgrade).""" + path = Path(config_path) + if not path.is_file(): + return {} + with path.open(encoding="utf-8") as fh: + data = yaml.safe_load(fh) or {} + return data if isinstance(data, dict) else {} + + +def _intended_runtime_is_kubernetes(raw: dict[str, Any]) -> bool: + platform = raw.get("platform") + if not isinstance(platform, dict): + return False + runtime = platform.get("runtime") + return isinstance(runtime, str) and runtime.strip().lower() == "kubernetes" + + +def _default_docker_executor_probe(raw: dict[str, Any]) -> _DefaultDockerExecutorProbe: + """Resolve deployments.default_executor → docker backend + optional config.docker_host.""" + deployments = raw.get("deployments") + if not isinstance(deployments, dict): + return _DefaultDockerExecutorProbe(is_docker=False) + default_name = deployments.get("default_executor") + if not isinstance(default_name, str) or not default_name: + return _DefaultDockerExecutorProbe(is_docker=False) + executors = deployments.get("executors") + if not isinstance(executors, list): + return _DefaultDockerExecutorProbe(is_docker=False) + for spec in executors: + if not isinstance(spec, dict): + continue + if spec.get("name") != default_name: + continue + backend = spec.get("backend") + if not (isinstance(backend, str) and backend.strip().lower() == "docker"): + return _DefaultDockerExecutorProbe(is_docker=False) + docker_host: str | None = None + config = spec.get("config") + if isinstance(config, dict): + host = config.get("docker_host") + if isinstance(host, str) and host.strip(): + docker_host = host.strip() + return _DefaultDockerExecutorProbe(is_docker=True, docker_host=docker_host) + return _DefaultDockerExecutorProbe(is_docker=False) + + +def _resolve_default_local_docker_probe( + platform_config: PlatformAppConfig | None = None, + *, + config_path: str | None = None, +) -> _DefaultDockerExecutorProbe | None: + """Return the docker probe target when this run needs the default-local gate; else None.""" + app_config = platform_config or PlatformAppConfig() + if config_path is not None: + app_config = replace(app_config, config_path=config_path) + + try: + resolved = resolve_run_configuration(app_config) + except ValueError: + # Invalid selections fail elsewhere; do not block on Docker here. + return None + + starts_deployments = "deployments" in resolved.services or "deployments" in resolved.controllers + if not starts_deployments: + return None + + raw = _load_platform_yaml(resolved.config_path or default_config_path()) + if _intended_runtime_is_kubernetes(raw): + return None + + probe_target = _default_docker_executor_probe(raw) + if not probe_target.is_docker: + return None + return probe_target + + +def default_local_needs_docker( + platform_config: PlatformAppConfig | None = None, + *, + config_path: str | None = None, +) -> bool: + """Return whether this run would fail closed on a missing Docker daemon.""" + return _resolve_default_local_docker_probe(platform_config, config_path=config_path) is not None + + +def require_docker_for_default_local( + platform_config: PlatformAppConfig | None = None, + *, + config_path: str | None = None, + console: Console | None = None, +) -> None: + """Exit with a clear message when default-local Docker is required but missing.""" + probe_target = _resolve_default_local_docker_probe(platform_config, config_path=config_path) + if probe_target is None: + return + if probe_docker(docker_host=probe_target.docker_host, use_cache=False).available: + return + out = console or Console(stderr=True) + out.print(f"[red]✗[/red] {DOCKER_PREFLIGHT_MESSAGE}") + raise typer.Exit(1) + + +__all__ = [ + "DOCKER_PREFLIGHT_MESSAGE", + "default_local_needs_docker", + "require_docker_for_default_local", +] diff --git a/sdk/python/nemo-platform/src/nemo_platform/quickstart/preflight.py b/sdk/python/nemo-platform/src/nemo_platform/quickstart/preflight.py index f7f9290657..0cccad3bf2 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/quickstart/preflight.py +++ b/sdk/python/nemo-platform/src/nemo_platform/quickstart/preflight.py @@ -11,6 +11,8 @@ from dataclasses import dataclass from enum import Enum +from nemo_platform_plugin.capabilities import probe_docker + from .config import QuickstartConfig @@ -96,11 +98,9 @@ def run_all(self) -> list[PreflightResult]: def _check_docker_available(self) -> None: """Verify Docker daemon is running and accessible.""" - try: - import docker - - client = docker.from_env() - client.ping() + # Preflight can be re-run after the user starts Docker. + result = probe_docker(use_cache=False) + if result.available: self.results.append( PreflightResult( name="Docker Available", @@ -108,13 +108,13 @@ def _check_docker_available(self) -> None: message="Docker daemon is running", ) ) - except Exception as e: + else: self.results.append( PreflightResult( name="Docker Available", status=CheckStatus.FAIL, message="Docker daemon is not accessible", - details=str(e), + details=result.detail or "Docker daemon unreachable", ) ) diff --git a/sdk/python/nemo-platform/src/nemo_platform/quickstart/validators.py b/sdk/python/nemo-platform/src/nemo_platform/quickstart/validators.py index 966c8110f7..216f0a172e 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/quickstart/validators.py +++ b/sdk/python/nemo-platform/src/nemo_platform/quickstart/validators.py @@ -9,6 +9,8 @@ from dataclasses import dataclass from pathlib import Path +from nemo_platform_plugin.capabilities import probe_docker + from .config import QuickstartConfig @@ -25,19 +27,16 @@ def __bool__(self) -> bool: def validate_docker_available() -> ValidationResult: - """Check if Docker daemon is available. + """Check if Docker daemon is available via an uncached reachability probe (≤5s). Returns: ValidationResult indicating success or failure. """ - try: - import docker - - client = docker.from_env() - client.ping() + # CLI may re-run after the user starts Docker; do not pin a prior miss. + result = probe_docker(use_cache=False) + if result.available: return ValidationResult(True, "Docker is available") - except Exception as e: - return ValidationResult(False, f"Docker is not available: {e}") + return ValidationResult(False, result.detail or "Docker is not available") def validate_ngc_credentials(api_key: str) -> ValidationResult: diff --git a/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/commands/test_setup.py b/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/commands/test_setup.py index e6ffe98bc9..4accc1fdee 100644 --- a/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/commands/test_setup.py +++ b/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/commands/test_setup.py @@ -702,6 +702,8 @@ def maybe_start_preflight_mocks(): patch(f"{SETUP_MOD}._start_services_background") as mock_start, patch(f"{SETUP_MOD}.prompt_choice", return_value="yes"), patch(f"{SETUP_MOD}._prompt_data_dir", return_value="/tmp/data"), + # Docker preflight is covered separately; keep other start-path tests focused. + patch(f"{SETUP_MOD}.require_docker_for_default_local"), ): yield mock_start @@ -814,6 +816,7 @@ def test_restarts_when_running_and_start_services_true(self): patch(f"{SETUP_MOD}._prompt_data_dir", return_value="/tmp/test-data") as mock_db_prompt, patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._ensure_port_available_for_start", wraps=_ensure_port_available_for_start) as mock_port, + patch(f"{SETUP_MOD}.require_docker_for_default_local"), patch(f"{SETUP_MOD}._pause"), ): mock_start.return_value = MagicMock(pid=999) @@ -834,6 +837,7 @@ def test_restarts_exits_when_port_still_occupied_after_kill(self, capsys): patch(f"{SETUP_MOD}._start_services_background") as mock_start, patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=conflict), patch(f"{SETUP_MOD}._prompt_data_dir", return_value="/tmp/test-data"), + patch(f"{SETUP_MOD}.require_docker_for_default_local"), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), ): @@ -882,7 +886,7 @@ def test_early_exit_names_eaddrinuse_when_startup_log_has_bind_failure( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", return_value=False), patch(f"{SETUP_MOD}.log_path_for", return_value=log), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=True), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=True)), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), ): @@ -911,7 +915,7 @@ def test_early_exit_prints_docker_hint_when_daemon_unavailable(self, maybe_start with ( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", wait), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=False), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=False)), patch(f"{SETUP_MOD}.log_path_for", return_value=MagicMock(__str__=lambda self: "/tmp/services.log")), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), @@ -923,7 +927,20 @@ def test_early_exit_prints_docker_hint_when_daemon_unavailable(self, maybe_start captured = capsys.readouterr() assert "exited early (exit code 3)" in captured.err assert "Check /tmp/services.log for details." in captured.err - assert "Docker does not appear to be available" in captured.err + assert "Docker is required for this default local setup" in captured.err + + def test_docker_preflight_blocks_before_spawn(self, maybe_start_preflight_mocks, capsys): + """Default-local Docker gate exits before start_background (NVBug 6537617).""" + with ( + patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), + patch( + f"{SETUP_MOD}.require_docker_for_default_local", + side_effect=typer.Exit(1), + ), + pytest.raises(ClickExit), + ): + _maybe_start_services("http://localhost:8080", auto=False, start_services=True) + maybe_start_preflight_mocks.assert_not_called() def test_readiness_timeout_does_not_hint_docker_without_evidence( self, maybe_start_preflight_mocks, capsys, tmp_path @@ -935,7 +952,7 @@ def test_readiness_timeout_does_not_hint_docker_without_evidence( with ( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", return_value=False), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=False), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=False)), patch(f"{SETUP_MOD}.log_path_for", return_value=log), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), @@ -954,7 +971,7 @@ def test_docker_hint_from_services_log_markers(self, maybe_start_preflight_mocks with ( patch(f"{SETUP_MOD}.check_port_available_for_start", return_value=None), patch(f"{SETUP_MOD}._wait_for_platform", return_value=False), - patch(f"{SETUP_MOD}.validate_docker_available", return_value=True), + patch(f"{SETUP_MOD}.probe_docker", return_value=MagicMock(available=True)), patch(f"{SETUP_MOD}.log_path_for", return_value=log), patch(f"{SETUP_MOD}._pause"), pytest.raises(ClickExit), @@ -962,7 +979,7 @@ def test_docker_hint_from_services_log_markers(self, maybe_start_preflight_mocks maybe_start_preflight_mocks.return_value = alive _maybe_start_services("http://localhost:8080", auto=False, start_services=True) captured = capsys.readouterr() - assert "Docker does not appear to be available" in captured.err + assert "Docker is required for this default local setup" in captured.err def test_startup_port_conflict_prefers_live_port_probe(self, tmp_path): log = tmp_path / "services.log" @@ -1273,6 +1290,7 @@ def test_auto_mode_skips_prompt_and_uses_persisted(self, tmp_path, monkeypatch): patch(f"{SETUP_MOD}._start_services_background") as mock_start, patch(f"{SETUP_MOD}._wait_for_platform", return_value=True), patch(f"{SETUP_MOD}._prompt_data_dir") as mock_prompt, + patch(f"{SETUP_MOD}.require_docker_for_default_local"), patch(f"{SETUP_MOD}._pause"), ): mock_start.return_value = MagicMock(pid=999) diff --git a/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/test_docker_preflight.py b/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/test_docker_preflight.py new file mode 100644 index 0000000000..b221d6ea08 --- /dev/null +++ b/sdk/python/nemo-platform/tests/vendored/nemo_platform_ext/cli/test_docker_preflight.py @@ -0,0 +1,209 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Tests for config-aware Docker preflight (NVBug 6537617).""" + +from __future__ import annotations + +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest +import typer +from nemo_platform.cli.docker_preflight import ( + DOCKER_PREFLIGHT_MESSAGE, + default_local_needs_docker, + require_docker_for_default_local, +) +from nemo_platform_plugin.capabilities import ProbeResult +from nmp.platform_runner.config import PlatformAppConfig + + +@pytest.fixture +def stock_local_yaml(tmp_path: Path) -> Path: + path = tmp_path / "local.yaml" + path.write_text( + """ +platform: + runtime: docker +deployments: + executors: + - name: local-docker + backend: docker + config: {} + default_executor: local-docker +""", + encoding="utf-8", + ) + return path + + +@pytest.fixture +def k8s_yaml(tmp_path: Path) -> Path: + path = tmp_path / "k8s.yaml" + path.write_text( + """ +platform: + runtime: kubernetes +deployments: + executors: + - name: local-k8s + backend: kubernetes + config: {} + default_executor: local-k8s +""", + encoding="utf-8", + ) + return path + + +@pytest.fixture +def non_docker_default_yaml(tmp_path: Path) -> Path: + path = tmp_path / "subprocess.yaml" + path.write_text( + """ +platform: + runtime: docker +deployments: + executors: + - name: local-sandbox + backend: sandbox + config: {} + default_executor: local-sandbox +""", + encoding="utf-8", + ) + return path + + +def test_default_local_needs_docker_for_stock_full_platform(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml)) + with patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments", "entities"}, + controllers={"deployments"}, + config_path=str(stock_local_yaml), + ), + ): + assert default_local_needs_docker(cfg) is True + + +def test_default_local_skips_when_deployments_not_selected(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml), services=["entities"], service_group=None) + with patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"entities"}, + controllers=set(), + config_path=str(stock_local_yaml), + ), + ): + assert default_local_needs_docker(cfg) is False + + +def test_default_local_skips_kubernetes_runtime(k8s_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(k8s_yaml)) + with patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(k8s_yaml), + ), + ): + assert default_local_needs_docker(cfg) is False + + +def test_default_local_skips_non_docker_default_executor(non_docker_default_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(non_docker_default_yaml)) + with patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers=set(), + config_path=str(non_docker_default_yaml), + ), + ): + assert default_local_needs_docker(cfg) is False + + +def test_require_docker_exits_without_spawn_when_probe_false(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml)) + with ( + patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(stock_local_yaml), + ), + ), + patch( + "nemo_platform.cli.docker_preflight.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ) as probe, + pytest.raises(typer.Exit) as exc, + ): + require_docker_for_default_local(cfg) + assert exc.value.exit_code == 1 + probe.assert_called_once_with(docker_host=None, use_cache=False) + + +def test_require_docker_noop_when_probe_true(stock_local_yaml: Path) -> None: + cfg = PlatformAppConfig(config_path=str(stock_local_yaml)) + with ( + patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(stock_local_yaml), + ), + ), + patch( + "nemo_platform.cli.docker_preflight.probe_docker", + return_value=ProbeResult(available=True), + ), + ): + require_docker_for_default_local(cfg) + + +def test_require_docker_probes_executor_docker_host(tmp_path: Path) -> None: + path = tmp_path / "remote-docker.yaml" + path.write_text( + """ +platform: + runtime: docker +deployments: + executors: + - name: remote-docker + backend: docker + config: + docker_host: tcp://docker.example:2375 + default_executor: remote-docker +""", + encoding="utf-8", + ) + cfg = PlatformAppConfig(config_path=str(path)) + with ( + patch( + "nemo_platform.cli.docker_preflight.resolve_run_configuration", + return_value=MagicMock( + services={"deployments"}, + controllers={"deployments"}, + config_path=str(path), + ), + ), + patch( + "nemo_platform.cli.docker_preflight.probe_docker", + return_value=ProbeResult(available=True), + ) as probe, + ): + require_docker_for_default_local(cfg) + probe.assert_called_once_with(docker_host="tcp://docker.example:2375", use_cache=False) + + +def test_preflight_message_names_docker() -> None: + assert "Docker" in DOCKER_PREFLIGHT_MESSAGE + assert "default local" in DOCKER_PREFLIGHT_MESSAGE.lower() or "default_executor" in DOCKER_PREFLIGHT_MESSAGE diff --git a/services/automodel/src/nmp/automodel/compile.py b/services/automodel/src/nmp/automodel/compile.py index a5476fbef6..f5798bf316 100644 --- a/services/automodel/src/nmp/automodel/compile.py +++ b/services/automodel/src/nmp/automodel/compile.py @@ -5,6 +5,8 @@ from __future__ import annotations +from nemo_platform import AsyncNeMoPlatform +from nemo_platform_plugin.jobs.api_factory import PlatformJobSpec from nmp.automodel.adapter import automodel_spec_to_compiler_output from nmp.automodel.api.v2.jobs.schemas import CustomizationJobOutput from nmp.automodel.app.jobs.compiler import platform_job_config_compiler as _compile_canonical @@ -13,10 +15,10 @@ async def platform_job_config_compiler( job_spec: CustomizationJobOutput | object, workspace: str, - sdk: object, + sdk: AsyncNeMoPlatform, job_name: str | None = None, profile: str | None = None, -) -> object: +) -> PlatformJobSpec: """Compile Automodel job spec (plugin or legacy shape) to PlatformJobSpec.""" if not isinstance(job_spec, CustomizationJobOutput): job_spec = automodel_spec_to_compiler_output(job_spec) diff --git a/services/core/jobs/pyproject.toml b/services/core/jobs/pyproject.toml index 4cff0ee7be..7cd3461a9d 100644 --- a/services/core/jobs/pyproject.toml +++ b/services/core/jobs/pyproject.toml @@ -25,6 +25,7 @@ dependencies = [ "pyyaml>=6.0.2", "hvac>=2.3.0", "nmp-common", + "nemo-platform-plugin", "aiosqlite>=0.20.0", # Required for SQLite async support "duckdb<2.0.0,>=1.1.3", "pandas>=1.5.3", @@ -36,6 +37,7 @@ jobs-server = "nmp.core.jobs.main:run_standalone" [tool.uv.sources] nmp-common = { workspace = true } +nemo-platform-plugin = { workspace = true } diff --git a/services/core/jobs/src/nmp/core/jobs/api/v2/jobs/endpoints.py b/services/core/jobs/src/nmp/core/jobs/api/v2/jobs/endpoints.py index 1422c0eca6..2fad1befba 100644 --- a/services/core/jobs/src/nmp/core/jobs/api/v2/jobs/endpoints.py +++ b/services/core/jobs/src/nmp/core/jobs/api/v2/jobs/endpoints.py @@ -241,7 +241,13 @@ def configured_subprocess_translation_profiles() -> set[str]: # Execution Profiles Endpoint @router.get("/v2/execution-profiles") async def get_execution_profiles() -> list[ExecutionProfileT]: - """Get all currently configured execution profiles.""" + """Get all currently configured execution profiles. + + Returns the capability-filtered merge from jobs config. In local standalone + the controller may prune the shared list further after registry boot; in + split topologies the API advertises its own merge result (not controller + process memory). + """ return profiles diff --git a/services/core/jobs/src/nmp/core/jobs/config.py b/services/core/jobs/src/nmp/core/jobs/config.py index 3ab8698a0d..0c9df51d60 100644 --- a/services/core/jobs/src/nmp/core/jobs/config.py +++ b/services/core/jobs/src/nmp/core/jobs/config.py @@ -71,6 +71,10 @@ def validate_executors(self) -> Self: # Module-level singleton instances config = get_service_config(JobsServiceConfig) platform_runtime = get_platform_config().runtime +# Capability-filtered at merge time (probe_docker). In local standalone the +# controller may further prune this list in place after backend registration. +# Split topologies (API vs controller pods) advertise this merge result from the +# API process — do not gate on controller-only process state (AIRCORE-971). profiles = merge_executor_profiles( config.executors, get_default_executor_profiles_for_runtime( diff --git a/services/core/jobs/src/nmp/core/jobs/controllers/backends/config.py b/services/core/jobs/src/nmp/core/jobs/controllers/backends/config.py index 998d22eedd..4582317af9 100644 --- a/services/core/jobs/src/nmp/core/jobs/controllers/backends/config.py +++ b/services/core/jobs/src/nmp/core/jobs/controllers/backends/config.py @@ -3,7 +3,7 @@ import logging -from nemo_platform_plugin.config import validate_docker_available +from nemo_platform_plugin.capabilities import probe_docker from nmp.common.config import Runtime from nmp.core.jobs.app.profiles import ExecutionProfileT from nmp.core.jobs.controllers.backends.docker import DockerJobExecutionProfile, DockerJobExecutionProfileConfig @@ -148,7 +148,10 @@ def merge_executor_profiles( def _docker_is_available() -> bool: nonlocal docker_available if docker_available is None: - docker_available = validate_docker_available() + # Uncached: this runs at jobs config import and must not pin the + # process-wide probe cache before BackendRegistry.from_config's + # authoritative boot probe. + docker_available = probe_docker(use_cache=False).available return docker_available # Add default profiles first diff --git a/services/core/jobs/src/nmp/core/jobs/controllers/backends/docker.py b/services/core/jobs/src/nmp/core/jobs/controllers/backends/docker.py index bed8f82ae6..43dba2f13b 100644 --- a/services/core/jobs/src/nmp/core/jobs/controllers/backends/docker.py +++ b/services/core/jobs/src/nmp/core/jobs/controllers/backends/docker.py @@ -20,6 +20,7 @@ from docker.errors import APIError, ImageNotFound, NotFound from docker.models.containers import Container from docker.types import LogConfig, Mount +from nemo_platform_plugin.capabilities import CapabilityUnavailableError, probe_docker from nemo_platform_plugin.jobs.execution_profiles import ( DockerJobExecutionProfile as PluginDockerJobExecutionProfile, ) @@ -289,6 +290,11 @@ def init(self) -> None: self._container_start_admission = threading.BoundedSemaphore(DOCKER_CONTAINER_START_WORKERS) self._container_run_threadpool = ThreadPoolExecutor(max_workers=DOCKER_CONTAINER_START_WORKERS) self._workload_identity_refreshers: dict[str, SubjectTokenRefreshLoop] = {} + # Short probe first — avoid docker.from_env(timeout=180) hanging when the + # daemon is down. CapabilityUnavailableError is soft-skipped by the registry. + probe = probe_docker() + if not probe.available: + raise CapabilityUnavailableError(probe.detail or "Docker daemon is unavailable") self._client = docker.from_env(timeout=180) if NEMO_JOBS_IMAGE_REGISTRY: logger.info( diff --git a/services/core/jobs/src/nmp/core/jobs/controllers/backends/registry.py b/services/core/jobs/src/nmp/core/jobs/controllers/backends/registry.py index 52ae04460a..3a489d7cb9 100644 --- a/services/core/jobs/src/nmp/core/jobs/controllers/backends/registry.py +++ b/services/core/jobs/src/nmp/core/jobs/controllers/backends/registry.py @@ -7,7 +7,11 @@ from docker.errors import DockerException from nemo_platform import NeMoPlatform -from nemo_platform_plugin.config import validate_docker_available +from nemo_platform_plugin.capabilities import ( + CapabilityUnavailableError, + probe_docker, + reset_capability_cache, +) from nmp.core.jobs.app.profiles import ExecutionProfileT from nmp.core.jobs.app.schemas import BackendRef, ProfileRef, ProviderRef from nmp.core.jobs.controllers.backends.base import DEFAULT_PROFILE, DEFAULT_PROVIDER, JobBackend @@ -24,9 +28,16 @@ logger = logging.getLogger(__name__) -# Skip Docker backends only for daemon/connection failures after validate_docker_available() -# was True. ValidationError and other programming/config errors must still fail startup. -_DOCKER_BACKEND_INIT_SKIPPABLE_ERRORS = (DockerException, RequestsConnectionError, RequestsTimeout, OSError) +# Skip Docker backends only for daemon/connection failures after probe_docker() +# was True, or when init raises CapabilityUnavailableError. ValidationError and +# other programming/config errors must still fail startup. +_DOCKER_BACKEND_INIT_SKIPPABLE_ERRORS = ( + CapabilityUnavailableError, + DockerException, + RequestsConnectionError, + RequestsTimeout, + OSError, +) @dataclass(frozen=True) @@ -132,8 +143,17 @@ def from_config( backend = backends[backend_key] if executor.backend == "docker": + # Process-lifetime cache is intentional: executor registration is + # fixed at controller startup. Import-time merge uses + # probe_docker(use_cache=False) so it does not pin this cache; + # compile shares the cached boot verdict (restart required after + # starting Docker mid-process). + # Reset before the boot probe so an earlier transient miss (e.g. + # deployments init probing before the daemon is ready) cannot + # permanently skip Docker executors for this process. if docker_available is None: - docker_available = validate_docker_available() + reset_capability_cache() + docker_available = probe_docker().available if not docker_available: logger.warning( "Skipping job executor profile %s/%s using backend 'docker' because Docker is unavailable.", diff --git a/services/core/jobs/tests/conftest.py b/services/core/jobs/tests/conftest.py index 9791e085f1..213e2046d2 100644 --- a/services/core/jobs/tests/conftest.py +++ b/services/core/jobs/tests/conftest.py @@ -3,6 +3,7 @@ import datetime import tempfile +from collections.abc import Iterator from contextlib import ExitStack from pathlib import Path from typing import AsyncGenerator @@ -13,6 +14,7 @@ from fastapi import FastAPI from httpx import ASGITransport, AsyncClient from nemo_platform import AsyncNeMoPlatform +from nemo_platform_plugin.capabilities import reset_capability_cache from nemo_platform_plugin.jobs.api_factory import ContainerSpec as FactoryContainerSpec from nemo_platform_plugin.jobs.api_factory import CPUExecutionProviderSpec as FactoryCPUExecutionProviderSpec from nemo_platform_plugin.jobs.api_factory import PlatformJobEnvironmentVariableParam, job_route_factory @@ -69,6 +71,27 @@ # ============================================================================ +@pytest.fixture(autouse=True) +def _reset_capability_cache() -> Iterator[None]: + """Prevent probe_docker process cache from poisoning jobs unit tests.""" + reset_capability_cache() + yield + reset_capability_cache() + + +@pytest.fixture +def docker_probe_unavailable(): + """Patch probe_docker as unavailable on merge + registry import paths.""" + from nemo_platform_plugin.capabilities import ProbeResult + + result = ProbeResult(available=False, detail="down") + with ( + patch("nmp.core.jobs.controllers.backends.config.probe_docker", return_value=result), + patch("nmp.core.jobs.controllers.backends.registry.probe_docker", return_value=result), + ): + yield result + + def pytest_collection_modifyitems(config, items): """ Modify test items during collection. diff --git a/services/core/jobs/tests/test_config.py b/services/core/jobs/tests/test_config.py index b2272f58cb..d2c2e28dc3 100644 --- a/services/core/jobs/tests/test_config.py +++ b/services/core/jobs/tests/test_config.py @@ -6,6 +6,7 @@ from unittest.mock import MagicMock, patch import pytest +from nemo_platform_plugin.capabilities import ProbeResult from nmp.common.config import Configuration, Runtime from nmp.core.jobs.app.providers import ( ComputeResources, @@ -415,7 +416,10 @@ def __init__(self, nmp_sdk, execution_profile_config, profile_name): ] caplog.set_level(logging.WARNING) - with patch("nmp.core.jobs.controllers.backends.registry.validate_docker_available", return_value=False): + with patch( + "nmp.core.jobs.controllers.backends.registry.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ): registry = BackendRegistry.from_config( nmp_sdk=mock_nmp_client, profiles=profiles, @@ -431,6 +435,42 @@ def __init__(self, nmp_sdk, execution_profile_config, profile_name): assert "Skipping job executor profile cpu/default" in caplog.text +def test_backend_registry_boot_probe_clears_poisoned_cache(mock_nmp_client): + """Registry boot must not permanently skip Docker after an earlier transient miss.""" + + class DummyBackend: + def __init__(self, nmp_sdk, execution_profile_config, profile_name): + self.nmp_sdk = nmp_sdk + self.execution_profile_config = execution_profile_config + self.profile_name = profile_name + + profiles = [ + DockerJobExecutionProfile( + provider="cpu", + profile="default", + backend="docker", + config=DockerJobExecutionProfileConfig(), + ), + ] + + with ( + patch("nmp.core.jobs.controllers.backends.registry.reset_capability_cache") as reset_cache, + patch( + "nmp.core.jobs.controllers.backends.registry.probe_docker", + return_value=ProbeResult(available=True), + ) as probe, + ): + registry = BackendRegistry.from_config( + nmp_sdk=mock_nmp_client, + profiles=profiles, + backends={BackendKey("cpu", "docker"): DummyBackend}, + ) + + reset_cache.assert_called_once_with() + probe.assert_called_once_with() + assert registry.get_backend(provider="cpu", profile="default") is not None + + def test_backend_registry_registered_profile_keys_match_constructed_backends(mock_nmp_client): class DummyBackend: def __init__(self, nmp_sdk, execution_profile_config, profile_name): @@ -452,7 +492,10 @@ def __init__(self, nmp_sdk, execution_profile_config, profile_name): ), ] - with patch("nmp.core.jobs.controllers.backends.registry.validate_docker_available", return_value=False): + with patch( + "nmp.core.jobs.controllers.backends.registry.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ): registry = BackendRegistry.from_config( nmp_sdk=mock_nmp_client, profiles=profiles, @@ -491,7 +534,7 @@ def __init__(self, nmp_sdk, execution_profile_config, profile_name): ] caplog.set_level(logging.WARNING) - with patch("nmp.core.jobs.controllers.backends.registry.validate_docker_available", return_value=True): + with patch("nmp.core.jobs.controllers.backends.registry.probe_docker", return_value=ProbeResult(available=True)): registry = BackendRegistry.from_config( nmp_sdk=mock_nmp_client, profiles=profiles, @@ -520,7 +563,7 @@ def __init__(self, nmp_sdk, execution_profile_config, profile_name): ] with ( - patch("nmp.core.jobs.controllers.backends.registry.validate_docker_available", return_value=True), + patch("nmp.core.jobs.controllers.backends.registry.probe_docker", return_value=ProbeResult(available=True)), pytest.raises(ValueError, match="bad executor config"), ): BackendRegistry.from_config( @@ -550,7 +593,10 @@ def __init__(self, nmp_sdk, execution_profile_config, profile_name): ] caplog.set_level(logging.WARNING) - with patch("nmp.core.jobs.controllers.backends.registry.validate_docker_available", return_value=False): + with patch( + "nmp.core.jobs.controllers.backends.registry.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ): registry = BackendRegistry.from_config( nmp_sdk=mock_nmp_client, profiles=advertised, @@ -635,7 +681,10 @@ def test_merge_executor_profiles_skips_container_backends_for_none_runtime(caplo ] caplog.set_level(logging.WARNING) - with patch("nmp.core.jobs.controllers.backends.config.validate_docker_available", return_value=False) as validate: + with patch( + "nmp.core.jobs.controllers.backends.config.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ) as validate: merged = merge_executor_profiles(custom, defaults, runtime=Runtime.NONE) assert [(p.provider, p.profile, p.backend) for p in merged] == [ @@ -665,7 +714,10 @@ def test_merge_executor_profiles_skips_docker_when_unavailable_under_docker_runt ] caplog.set_level(logging.WARNING) - with patch("nmp.core.jobs.controllers.backends.config.validate_docker_available", return_value=False) as validate: + with patch( + "nmp.core.jobs.controllers.backends.config.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ) as validate: merged = merge_executor_profiles(custom, defaults, runtime=Runtime.DOCKER) assert [(p.provider, p.profile, p.backend) for p in merged] == [ @@ -687,7 +739,9 @@ def test_merge_executor_profiles_keeps_docker_executor_for_none_runtime_when_doc ) ] - with patch("nmp.core.jobs.controllers.backends.config.validate_docker_available", return_value=True) as validate: + with patch( + "nmp.core.jobs.controllers.backends.config.probe_docker", return_value=ProbeResult(available=True) + ) as validate: merged = merge_executor_profiles(custom, defaults, runtime=Runtime.NONE) assert [(p.provider, p.profile, p.backend) for p in merged] == [ @@ -857,7 +911,10 @@ def test_merge_executor_profiles_keeps_subprocess_override_for_none_runtime(capl ), ] - with patch("nmp.core.jobs.controllers.backends.config.validate_docker_available", return_value=False): + with patch( + "nmp.core.jobs.controllers.backends.config.probe_docker", + return_value=ProbeResult(available=False, detail="down"), + ): merged = merge_executor_profiles(custom_executors, default_executors, runtime=Runtime.NONE) assert [(p.provider, p.profile, p.backend) for p in merged] == [("subprocess", "default", "subprocess")] diff --git a/services/core/jobs/tests/test_jobs_api.py b/services/core/jobs/tests/test_jobs_api.py index e8de61fe8a..4d93f8d2d9 100644 --- a/services/core/jobs/tests/test_jobs_api.py +++ b/services/core/jobs/tests/test_jobs_api.py @@ -432,7 +432,7 @@ async def test_create_job_gpu_fail_fast_when_docker_no_gpus(test_client: AsyncCl mock_platform_config.runtime = Runtime.DOCKER mock_platform_config.docker.get_reserved_gpu_ids.return_value = [] - with patch("nmp.common.jobs.docker.get_platform_config", return_value=mock_platform_config): + with patch("nemo_platform_plugin.jobs.docker.get_platform_config", return_value=mock_platform_config): response = await test_client.post("/apis/jobs/v2/workspaces/default/jobs", json=req.model_dump()) assert response.status_code == 422 diff --git a/services/core/jobs/tests/test_jobs_client.py b/services/core/jobs/tests/test_jobs_client.py index f5bd5750c1..00661dfa64 100644 --- a/services/core/jobs/tests/test_jobs_client.py +++ b/services/core/jobs/tests/test_jobs_client.py @@ -67,6 +67,20 @@ async def test_get_execution_profiles_parses_response(jobs_client: AsyncJobsClie assert isinstance(profiles, list) +@pytest.mark.asyncio +async def test_get_execution_profiles_available_without_controller_ready_gate( + test_client: AsyncClient, +): + """API must advertise merge-filtered profiles without a controller ready flag. + + Split topologies (API pod without controllers) previously 503'd forever when + readiness lived in controller-only process memory (AIRCORE-971). + """ + raw = await test_client.get("/apis/jobs/v2/execution-profiles") + assert raw.status_code == 200 + assert isinstance(raw.json(), list) + + async def _create_hello_world_job(test_client: AsyncClient, name: str = "e2e-client-job") -> None: """Create a job via the hello-world factory route (service-specific body).""" resp = await test_client.post( diff --git a/uv.lock b/uv.lock index ad0fde3fd5..88fcd53790 100644 --- a/uv.lock +++ b/uv.lock @@ -5104,6 +5104,7 @@ jobs-service = [ { name = "fastapi", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "hvac", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "kubernetes", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, + { name = "nemo-platform-plugin", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "nmp-common", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "pandas", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "pydantic", extra = ["email"], marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, @@ -5702,6 +5703,7 @@ requires-dist = [ { name = "nemo-platform-plugin", marker = "extra == 'all'", editable = "packages/nemo_platform_plugin" }, { name = "nemo-platform-plugin", marker = "extra == 'core-service'", editable = "packages/nemo_platform_plugin" }, { name = "nemo-platform-plugin", marker = "extra == 'inference-gateway-service'", editable = "packages/nemo_platform_plugin" }, + { name = "nemo-platform-plugin", marker = "extra == 'jobs-service'", editable = "packages/nemo_platform_plugin" }, { name = "nemo-platform-plugin", marker = "extra == 'models-service'", editable = "packages/nemo_platform_plugin" }, { name = "nemo-platform-plugin", marker = "extra == 'nemo-agents-plugin'", editable = "packages/nemo_platform_plugin" }, { name = "nemo-platform-plugin", marker = "extra == 'nemo-anonymizer-plugin'", editable = "packages/nemo_platform_plugin" }, @@ -7722,6 +7724,7 @@ dependencies = [ { name = "fastapi", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "hvac", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "kubernetes", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, + { name = "nemo-platform-plugin", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "nmp-common", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "pandas", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, { name = "pydantic", extra = ["email"], marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux')" }, @@ -7750,6 +7753,7 @@ requires-dist = [ { name = "fastapi", specifier = ">=0.115.8" }, { name = "hvac", specifier = ">=2.3.0" }, { name = "kubernetes", specifier = ">=30.1.0" }, + { name = "nemo-platform-plugin", editable = "packages/nemo_platform_plugin" }, { name = "nmp-common", editable = "packages/nmp_common" }, { name = "pandas", specifier = ">=1.5.3" }, { name = "pydantic", specifier = ">=2.10.6" },