diff --git a/e2e/configs/local-docker-agents.yaml b/e2e/configs/local-docker-agents.yaml index f14663bb49..953716c47e 100644 --- a/e2e/configs/local-docker-agents.yaml +++ b/e2e/configs/local-docker-agents.yaml @@ -49,8 +49,8 @@ deployments: # local Docker daemon (pulled by the test's CI job), so disable the # per-run pull and use the image as-is. pull_images: false - port_range_start: 9000 - port_range_end: 9100 + port_range_start: 49152 + port_range_end: 49251 models: controller: diff --git a/packages/nmp_platform/config/local.yaml b/packages/nmp_platform/config/local.yaml index fa2dbc301d..cd723a93d3 100644 --- a/packages/nmp_platform/config/local.yaml +++ b/packages/nmp_platform/config/local.yaml @@ -128,8 +128,8 @@ deployments: backend: docker config: pull_images: false - port_range_start: 9000 - port_range_end: 9100 + port_range_start: 49152 + port_range_end: 49251 default_executor: local-docker models: diff --git a/packages/nmp_platform_runner/src/nmp/platform_runner/config/local.yaml b/packages/nmp_platform_runner/src/nmp/platform_runner/config/local.yaml index 2636247775..1e22469627 100644 --- a/packages/nmp_platform_runner/src/nmp/platform_runner/config/local.yaml +++ b/packages/nmp_platform_runner/src/nmp/platform_runner/config/local.yaml @@ -78,8 +78,8 @@ deployments: backend: docker config: pull_images: false - port_range_start: 9000 - port_range_end: 9100 + port_range_start: 49152 + port_range_end: 49251 default_executor: local-docker inference_gateway: {} diff --git a/plugins/nemo-deployments/README.md b/plugins/nemo-deployments/README.md index 8af6d497a1..4d74e1e1cd 100644 --- a/plugins/nemo-deployments/README.md +++ b/plugins/nemo-deployments/README.md @@ -30,8 +30,8 @@ Host port allocation for published container ports is configured per named docke executor (not on `DeploymentConfig.backend_config.docker`). Set the inclusive `port_range_start` / `port_range_end` bounds on the executor `config` block in platform YAML. The allocator scans every host port from `port_range_start` -through `port_range_end`, including both endpoints (for example, 9000–9100 -allows 101 ports): +through `port_range_end`, including both endpoints (for example, 49152–49251 +allows 100 ports): ```yaml deployments: @@ -39,11 +39,17 @@ deployments: - name: local-docker backend: docker config: - port_range_start: 9000 - port_range_end: 9100 # inclusive + port_range_start: 49152 + port_range_end: 49251 # inclusive default_executor: local-docker ``` +A port is considered taken when any container on the daemon publishes it — not just +platform-managed ones — or when the host socket is already bound. Docker's own +reservations are not observable until a publish is attempted, so a port claimed +concurrently is retried against a different port rather than failing the deployment. +Pick a range nothing else on the host publishes on. + Entity-level `backend_config.docker` accepts only deployment-specific overrides such as `network`. 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 b997537877..81b803490b 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 @@ -75,6 +75,11 @@ _ONE_SHOT_RESTART_POLICIES = frozenset({"Never", "OnFailure"}) _EXITED_CONTAINER_STATES = frozenset({"exited", "dead"}) +# Docker rejects a publish with this message when its own allocator already holds +# the port. Its reservations are invisible to a host-socket probe, so a container +# created concurrently can claim a port between allocation and run. +_PORT_CONFLICT_MARKER = "port is already allocated" +_PORT_CONFLICT_ATTEMPTS = 3 NGC_IMAGE_REGISTRY = os.getenv("NGC_IMAGE_REGISTRY", "nvcr.io") NGC_IMAGE_REGISTRY_USER_NAME = os.getenv("NGC_IMAGE_REGISTRY_USER_NAME", "$oauthtoken") @@ -198,21 +203,11 @@ async def create_deployment( except GPUAllocationError as exc: return BackendStatusUpdate(status="FAILED", status_message=str(exc)) - host_ports: dict[int, int] = {} - for port_spec in container_spec.ports: - host_port = await find_available_port( - self._client, - self._executor_config.port_range_start, - self._executor_config.port_range_end, - exclude_ports=set(host_ports.values()), - ) - if host_port is None: - if gpu_ids and gpu_pool is not None: - gpu_pool.release_gpu(dep_key) - return BackendStatusUpdate( - status="FAILED", status_message="No host ports available in configured range" - ) - host_ports[port_spec.container_port] = host_port + host_ports = await self._allocate_host_ports(container_spec) + if host_ports is None: + if gpu_ids and gpu_pool is not None: + gpu_pool.release_gpu(dep_key) + return BackendStatusUpdate(status="FAILED", status_message="No host ports available in configured range") # Pull all images in the group (init + primary + sidecars) up front. if self._executor_config.pull_images: @@ -265,23 +260,20 @@ async def create_deployment( return init_status # 2) Primary (server) container: publishes ports, owns GPUs. - server_run_kwargs = self._build_run_kwargs( + server_container, host_ports, run_error = await self._run_server_container( workspace=workspace, config=config, - container=container_spec, + container_spec=container_spec, name=c_name, labels={**base_labels, CONTAINER_ROLE_LABEL: CONTAINER_ROLE_SERVER}, host_ports=host_ports, gpu_ids=gpu_ids, network=docker_cfg.network, ) - try: - await asyncio.to_thread(self._client.containers.run, **server_run_kwargs) - except Exception as exc: + if server_container is None: if gpu_ids and gpu_pool is not None: gpu_pool.release_gpu(dep_key) - logger.exception("Failed to start container %s", c_name) - return BackendStatusUpdate(status="FAILED", status_message=f"Failed to start container: {exc}") + return BackendStatusUpdate(status="FAILED", status_message=run_error) # 3) Sidecars share the primary's network namespace + volumes; no ports/GPU. # @@ -318,12 +310,111 @@ async def create_deployment( ) endpoints = self._build_endpoints(container_spec, host_ports) + if config.restart_policy in _ONE_SHOT_RESTART_POLICIES: + return await self._observe_one_shot_primary_after_create( + workspace=workspace, + name=name, + config=config, + container=server_container, + dep_key=dep_key, + endpoints=endpoints, + ) return BackendStatusUpdate( status="STARTING", status_message=f"Container {c_name} created", endpoints=endpoints, ) + async def _allocate_host_ports( + self, + container_spec: Container, + *, + exclude_ports: set[int] | None = None, + ) -> dict[int, int] | None: + """Map each published container port to a free host port. + + Returns None when the configured range holds no more free ports. + """ + excluded = set(exclude_ports or ()) + host_ports: dict[int, int] = {} + for port_spec in container_spec.ports: + host_port = await find_available_port( + self._client, + self._executor_config.port_range_start, + self._executor_config.port_range_end, + exclude_ports=excluded | set(host_ports.values()), + ) + if host_port is None: + return None + host_ports[port_spec.container_port] = host_port + return host_ports + + async def _remove_container_by_name(self, name: str) -> None: + """Best-effort removal of a container this call just created.""" + try: + container = await asyncio.to_thread(self._client.containers.get, name) + await asyncio.to_thread(container.remove, force=True) + except self._docker_errors.NotFound: + pass + except Exception: + logger.warning("Failed to remove container %s", name, exc_info=True) + + async def _run_server_container( + self, + *, + workspace: str, + config: DeploymentConfig, + container_spec: Container, + name: str, + labels: dict[str, str], + host_ports: dict[int, int], + gpu_ids: list[int], + network: str | None, + ) -> tuple[DockerContainer | None, dict[int, int], str]: + """Start the primary container, reallocating host ports on Docker port conflicts. + + Returns the started container (None on failure), the host ports it actually + published, and an error message when it could not be started. + """ + rejected_ports: set[int] = set() + attempt = 0 + while True: + attempt += 1 + run_kwargs = self._build_run_kwargs( + workspace=workspace, + config=config, + container=container_spec, + name=name, + labels=labels, + host_ports=host_ports, + gpu_ids=gpu_ids, + network=network, + ) + try: + container = await asyncio.to_thread(self._client.containers.run, **run_kwargs) + return container, host_ports, "" + except Exception as exc: + last_attempt = attempt == _PORT_CONFLICT_ATTEMPTS + if not host_ports or last_attempt or _PORT_CONFLICT_MARKER not in str(exc): + logger.exception("Failed to start container %s", name) + return None, host_ports, f"Failed to start container: {exc}" + + rejected_ports |= set(host_ports.values()) + logger.warning( + "Host port conflict starting %s (attempt %d/%d); reallocating outside %s", + name, + attempt, + _PORT_CONFLICT_ATTEMPTS, + sorted(rejected_ports), + ) + # containers.run() creates then starts, so a failed start leaves the + # created container holding the name and blocking the retry. + await self._remove_container_by_name(name) + reallocated = await self._allocate_host_ports(container_spec, exclude_ports=rejected_ports) + if reallocated is None: + return None, host_ports, "No host ports available in configured range" + host_ports = reallocated + def _build_run_kwargs( self, *, @@ -423,7 +514,7 @@ async def _run_init_container( def _run_and_wait() -> int: container = self._client.containers.run(**run_kwargs) result = container.wait(timeout=self._executor_config.docker_timeout) - exit_code = int(result.get("StatusCode", 1)) if isinstance(result, dict) else int(result) + exit_code = self._exit_code_from_wait_result(result) try: container.remove(force=True) except Exception: @@ -526,47 +617,156 @@ async def read_status(self, *, workspace: str, name: str) -> BackendStatusUpdate if state in ("exited", "dead"): exit_code = int(container.attrs.get("State", {}).get("ExitCode", 1)) - if exit_code == 0 and restart_policy in ("Never", "OnFailure"): - if self._gpu_pool is not None: - self._gpu_pool.release_gpu(dep_key) + restart_count = int(container.attrs.get("RestartCount", 0)) + return self._status_from_exited_container( + exit_code=exit_code, + restart_policy=restart_policy, + labels=labels, + dep_key=dep_key, + restart_count=restart_count, + endpoints=endpoints, + ) + + if state == "removing": + return BackendStatusUpdate(status="DELETING", status_message=f"Container removing (ID: {container_id})") + + return BackendStatusUpdate(status="STARTING", status_message=f"Container state: {state}") + + @staticmethod + def _exit_code_from_wait_result(result: dict[str, Any] | int) -> int: + """Normalize a docker `container.wait()` result to an exit code.""" + return int(result.get("StatusCode", 1)) if isinstance(result, dict) else int(result) + + def _status_from_exited_container( + self, + *, + exit_code: int, + restart_policy: RestartPolicy, + labels: dict[str, str], + dep_key: str, + restart_count: int = 0, + endpoints: list[Endpoint] | None = None, + ) -> BackendStatusUpdate: + """Map a stopped container's exit code to a deployment status update.""" + resolved_endpoints = endpoints or [] + if exit_code == 0 and restart_policy in ("Never", "OnFailure"): + if self._gpu_pool is not None: + self._gpu_pool.release_gpu(dep_key) + return BackendStatusUpdate( + status="SUCCEEDED", + status_message="Container exited successfully (code 0)", + exit_code=exit_code, + endpoints=resolved_endpoints, + ) + if restart_policy == "Always": + return BackendStatusUpdate( + status="STARTING", + status_message=f"Container exited (code {exit_code}); restart policy will recreate it", + exit_code=exit_code, + endpoints=resolved_endpoints, + ) + if restart_policy == "OnFailure": + backoff_limit = int(labels.get(BACKOFF_LIMIT_LABEL, "6")) + if restart_count < backoff_limit: return BackendStatusUpdate( - status="SUCCEEDED", - status_message="Container exited successfully (code 0)", + status="STARTING", + status_message=(f"Container exited (code {exit_code}); retry {restart_count}/{backoff_limit}"), exit_code=exit_code, - endpoints=endpoints, + endpoints=resolved_endpoints, ) - if restart_policy == "Always": + if self._gpu_pool is not None: + self._gpu_pool.release_gpu(dep_key) + status = map_exited_status(exit_code, restart_policy) + return BackendStatusUpdate( + status=status, + status_message=f"Container exited with code {exit_code}", + exit_code=exit_code, + endpoints=resolved_endpoints, + ) + + async def _observe_one_shot_primary_after_create( + self, + *, + workspace: str, + name: str, + config: DeploymentConfig, + container: DockerContainer, + dep_key: str, + endpoints: list[Endpoint], + ) -> BackendStatusUpdate: + """Observe a one-shot primary container immediately after create.""" + restart_policy = config.restart_policy + labels = container.labels or {} + + if restart_policy == "Never": + observe_timeout = self._executor_config.oneshot_observe_timeout_seconds + + def _wait_for_exit() -> int: + return self._exit_code_from_wait_result(container.wait(timeout=observe_timeout)) + + try: + exit_code = await asyncio.to_thread(_wait_for_exit) + except ( + ReadTimeout, + Urllib3ReadTimeoutError, + RequestsConnectionError, + self._docker_errors.APIError, + ): return BackendStatusUpdate( status="STARTING", - status_message=f"Container exited (code {exit_code}); restart policy will recreate it", - exit_code=exit_code, + status_message=( + f"Container {container.name} still running or status unavailable " + f"after observe wait ({observe_timeout}s)" + ), endpoints=endpoints, ) - if restart_policy == "OnFailure": - restart_count = int(container.attrs.get("RestartCount", 0)) - backoff_limit = int(labels.get(BACKOFF_LIMIT_LABEL, "6")) - if restart_count < backoff_limit: - return BackendStatusUpdate( - status="STARTING", - status_message=f"Container exited (code {exit_code}); retry {restart_count}/{backoff_limit}", - exit_code=exit_code, - endpoints=endpoints, - ) - if self._gpu_pool is not None: - self._gpu_pool.release_gpu(dep_key) - status = map_exited_status(exit_code, restart_policy) - message = f"Container exited with code {exit_code}" - return BackendStatusUpdate( - status=status, - status_message=message, + except Exception as exc: + logger.exception("Failed waiting for one-shot container %s to exit", container.name) + await self.delete_deployment(workspace, name) + return BackendStatusUpdate( + status="FAILED", + status_message=f"Failed waiting for container to exit: {exc}", + ) + return self._status_from_exited_container( exit_code=exit_code, + restart_policy=restart_policy, + labels=labels, + dep_key=dep_key, endpoints=endpoints, ) - if state == "removing": - return BackendStatusUpdate(status="DELETING", status_message=f"Container removing (ID: {container_id})") + try: + await asyncio.to_thread(container.reload) + except ( + self._docker_errors.APIError, + ReadTimeout, + Urllib3ReadTimeoutError, + RequestsConnectionError, + ) as exc: + return BackendStatusUpdate( + status="STARTING", + status_message=f"Container created but status check failed: {exc}", + endpoints=endpoints, + ) + if container.status in _EXITED_CONTAINER_STATES: + attrs = container.attrs or {} + exit_code = int(attrs.get("State", {}).get("ExitCode", 1)) + restart_count = int(attrs.get("RestartCount", 0)) + return self._status_from_exited_container( + exit_code=exit_code, + restart_policy=restart_policy, + labels=labels, + dep_key=dep_key, + restart_count=restart_count, + endpoints=endpoints, + ) - return BackendStatusUpdate(status="STARTING", status_message=f"Container state: {state}") + container_name_label = container.name or "primary" + return BackendStatusUpdate( + status="STARTING", + status_message=f"Container {container_name_label} created", + endpoints=endpoints, + ) async def _sidecars_healthy(self, workspace: str, name: str, config: DeploymentConfig | None) -> tuple[bool, str]: """Return (all_healthy, reason) for a deployment's expected sidecar containers. diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/config.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/config.py index f27dfca0ae..20e25b5d2d 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/config.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/config.py @@ -18,6 +18,15 @@ class DockerExecutorConfig(BaseModel): ge=1, description="Docker client timeout in seconds for pull/create/status operations (default: 10 minutes).", ) + oneshot_observe_timeout_seconds: int = Field( + default=5, + ge=1, + description=( + "Max seconds to wait for a Never one-shot container to exit during create. " + "Should stay near the deployments controller reconcile interval (default 5s). " + "Longer jobs return STARTING and finish via read_status polling." + ), + ) pull_images: bool = Field(default=True, description="Pull container images before run when missing locally.") resource_scope: str = Field( default=DEFAULT_RESOURCE_SCOPE, @@ -28,13 +37,13 @@ class DockerExecutorConfig(BaseModel): ), ) port_range_start: int = Field( - default=9000, + default=49152, ge=1, le=65535, description="First host port to consider when publishing container ports for this executor.", ) port_range_end: int = Field( - default=9100, + default=49251, ge=1, le=65535, description="Last host port (inclusive) to consider when publishing container ports for this executor.", diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/ports.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/ports.py index d9516a1f4f..57ec702b76 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/ports.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/ports.py @@ -11,8 +11,6 @@ import socket from typing import TYPE_CHECKING -from nemo_deployments_plugin.backends.labels import managed_by_filter - import docker if TYPE_CHECKING: @@ -30,8 +28,10 @@ def is_port_free(port: int) -> bool: if is_remote_docker_host(): return True try: + # Do not set SO_REUSEADDR: it can make this bind succeed while a Docker + # wildcard publisher already holds the port. Binding loopback is enough + # to detect that conflict without exposing a socket on external interfaces. with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: - sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) sock.bind(("127.0.0.1", port)) return True except OSError: @@ -62,11 +62,10 @@ async def find_available_port( exclude_ports: set[int] | None = None, ) -> int | None: try: - containers = await asyncio.to_thread( - client.containers.list, - all=True, - filters=managed_by_filter(), - ) + # Every container on the daemon competes for host ports, not just the ones + # this platform manages. ``resource_scope`` scopes ownership and cleanup; + # port safety has to consider foreign containers too. + containers = await asyncio.to_thread(client.containers.list, all=True) except Exception: logger.exception("Failed to list containers for port allocation") return None diff --git a/plugins/nemo-deployments/tests/integration/backends/docker/test_docker_backend.py b/plugins/nemo-deployments/tests/integration/backends/docker/test_docker_backend.py index 236adb4849..a2938160fa 100644 --- a/plugins/nemo-deployments/tests/integration/backends/docker/test_docker_backend.py +++ b/plugins/nemo-deployments/tests/integration/backends/docker/test_docker_backend.py @@ -6,6 +6,8 @@ from __future__ import annotations import asyncio +import time +from typing import Any from unittest.mock import AsyncMock, MagicMock, patch import pytest @@ -30,29 +32,70 @@ ] -@pytest.fixture -def docker_backend() -> DockerDeploymentBackend: +ALPINE_IMAGE = "alpine:3.20" + +# Keep these tests in a dedicated range that nothing else in CI claims, separate +# from both service ports and the product's dynamic/private default range. +TEST_PORT_RANGE_START = 21000 +TEST_PORT_RANGE_END = 21100 + + +def _build_docker_backend(**config_overrides: Any) -> DockerDeploymentBackend: mock_entities = AsyncMock() mock_sdk = MagicMock() + executor_config: dict[str, Any] = { + "pull_images": True, + "port_range_start": TEST_PORT_RANGE_START, + "port_range_end": TEST_PORT_RANGE_END, + **config_overrides, + } with ( patch("nemo_deployments_plugin.backends.docker.backend.AsyncEntitiesResource"), patch("nemo_deployments_plugin.backends.docker.backend.NemoEntitiesClient", return_value=mock_entities), patch("nemo_deployments_plugin.backends.docker.backend.get_shared_gpu_pool", return_value=None), ): - backend = DockerDeploymentBackend(mock_sdk, {"pull_images": True}) + backend = DockerDeploymentBackend(mock_sdk, executor_config) backend._entities = mock_entities return backend +@pytest.fixture +def docker_backend() -> DockerDeploymentBackend: + return _build_docker_backend() + + def _never_config() -> DeploymentConfig: return DeploymentConfig( name="echo-cfg", workspace="itest", restart_policy="Never", # ty: ignore[unknown-argument] - containers=[Container(name="main", image="alpine:3.20", command=["echo"], args=["hello"])], + containers=[Container(name="main", image=ALPINE_IMAGE, command=["echo"], args=["hello"])], + ) + + +def _never_sleep_config(*, sleep_seconds: int) -> DeploymentConfig: + return DeploymentConfig( + name="sleep-cfg", + workspace="itest", + restart_policy="Never", # ty: ignore[unknown-argument] + containers=[ + Container( + name="main", + image=ALPINE_IMAGE, + command=["sleep"], + args=[str(sleep_seconds)], + ) + ], ) +def _docker_backend_with_observe_timeout( + *, + oneshot_observe_timeout_seconds: int, +) -> DockerDeploymentBackend: + return _build_docker_backend(oneshot_observe_timeout_seconds=oneshot_observe_timeout_seconds) + + def _always_http_config() -> DeploymentConfig: return DeploymentConfig( name="http-cfg", @@ -89,7 +132,7 @@ async def test_volume_lifecycle(docker_backend: DockerDeploymentBackend) -> None @pytest.mark.asyncio async def test_never_deployment_succeeds(docker_backend: DockerDeploymentBackend) -> None: config = _never_config() - docker_backend._entities.get.return_value = config # type: ignore[attr-defined] + docker_backend._entities.get.return_value = config # ty: ignore[unresolved-attribute] c_name = container_name("itest", "echo-job") client = docker.from_env() @@ -101,18 +144,61 @@ async def test_never_deployment_succeeds(docker_backend: DockerDeploymentBackend labels={"managed-by": MANAGED_BY_LABEL}, backend_config={}, ) - assert created.status == "STARTING" + assert created.status == "SUCCEEDED" + assert created.exit_code == 0 + + status = await docker_backend.read_status(workspace="itest", name="echo-job") + assert status.status == "SUCCEEDED" + assert status.exit_code == 0 + finally: + await docker_backend.delete_deployment("itest", "echo-job") + force_remove_container(client, c_name) + + +@pytest.mark.asyncio +async def test_never_deployment_outlives_observe_wait_then_succeeds() -> None: + """Long Never jobs return STARTING on create and finish via read_status polling.""" + job_sleep_seconds = 5 + observe_timeout_seconds = 1 + docker_backend = _docker_backend_with_observe_timeout( + oneshot_observe_timeout_seconds=observe_timeout_seconds, + ) + config = _never_sleep_config(sleep_seconds=job_sleep_seconds) + docker_backend._entities.get.return_value = config # ty: ignore[unresolved-attribute] + c_name = container_name("itest", "sleep-job") + client = docker.from_env() + + try: + # Warm the image cache so the timed window below measures the observe wait + # rather than an uncached image pull. + await asyncio.to_thread(client.images.pull, ALPINE_IMAGE) - for _ in range(30): - status = await docker_backend.read_status(workspace="itest", name="echo-job") - if status.status in ("SUCCEEDED", "FAILED"): + started = time.monotonic() + created = await docker_backend.create_deployment( + workspace="itest", + name="sleep-job", + config_name="sleep-cfg", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + create_elapsed = time.monotonic() - started + + assert created.status == "STARTING" + assert "after observe wait" in (created.status_message or "") + assert create_elapsed < observe_timeout_seconds + 2.0 + + deadline = time.monotonic() + 15.0 + status = created + while time.monotonic() < deadline: + status = await docker_backend.read_status(workspace="itest", name="sleep-job") + if status.status in {"SUCCEEDED", "FAILED"}: break await asyncio.sleep(0.5) assert status.status == "SUCCEEDED" assert status.exit_code == 0 finally: - await docker_backend.delete_deployment("itest", "echo-job") + await docker_backend.delete_deployment("itest", "sleep-job") force_remove_container(client, c_name) @@ -126,7 +212,7 @@ async def get_side_effect(entity_type, name, workspace=None): return deployment return config - docker_backend._entities.get.side_effect = get_side_effect # type: ignore[attr-defined] + docker_backend._entities.get.side_effect = get_side_effect # ty: ignore[unresolved-attribute] c_name = container_name("itest", "lost-srv") client = docker.from_env() diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/docker_helpers.py b/plugins/nemo-deployments/tests/unit/backends/docker/docker_helpers.py index 93d4916d1c..5f3c887516 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/docker_helpers.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/docker_helpers.py @@ -27,6 +27,22 @@ def sample_config(*, restart_policy: RestartPolicy = "Always") -> DeploymentConf ) +def published_port_config(*, restart_policy: RestartPolicy = "Always") -> DeploymentConfig: + """Single-container config that publishes a host port.""" + return DeploymentConfig( + name="cfg1", + workspace="default", + containers=[ + Container( + name="main", + image="alpine:latest", + ports=[ContainerPort(name="http", containerPort=8000)], + ) + ], + restart_policy=restart_policy, # ty: ignore[unknown-argument] + ) + + def lora_config(*, restart_policy: RestartPolicy = "Always") -> DeploymentConfig: """A LoRA-shaped multi-container config: init + server + adapters sidecar. diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/test_backend_mocked.py b/plugins/nemo-deployments/tests/unit/backends/docker/test_backend_mocked.py index f419afcaa1..6b502504f2 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/test_backend_mocked.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/test_backend_mocked.py @@ -8,10 +8,22 @@ from unittest.mock import AsyncMock, MagicMock, patch import pytest -from backends.docker.docker_helpers import container_attrs, lora_config, sample_config +from backends.docker.docker_helpers import ( + container_attrs, + lora_config, + published_port_config, + sample_config, +) from docker.errors import APIError, NotFound -from nemo_deployments_plugin.backends.docker.backend import DockerDeploymentBackend +from nemo_deployments_plugin.backends.base import BackendStatusUpdate +from nemo_deployments_plugin.backends.docker import ports as ports_mod +from nemo_deployments_plugin.backends.docker.backend import ( + _PORT_CONFLICT_ATTEMPTS, + _PORT_CONFLICT_MARKER, + DockerDeploymentBackend, +) from nemo_deployments_plugin.backends.labels import ( + BACKOFF_LIMIT_LABEL, CONFIG_NAME_LABEL, CONTAINER_ROLE_LABEL, DEFAULT_RESOURCE_SCOPE, @@ -27,6 +39,8 @@ from nemo_deployments_plugin.constants import MANAGED_BY_LABEL from nemo_deployments_plugin.entities import Deployment from nemo_deployments_plugin.types import RestartPolicy +from requests.exceptions import ConnectionError as RequestsConnectionError +from requests.exceptions import ReadTimeout @pytest.mark.asyncio @@ -81,6 +95,105 @@ async def test_create_deployment_maps_command_to_entrypoint( assert run_kwargs["command"] == ["hello"] +def _port_conflict_error(port: int) -> APIError: + return APIError( + f"driver failed programming external connectivity: Bind for 0.0.0.0:{port} failed: {_PORT_CONFLICT_MARKER}" + ) + + +@pytest.fixture +def free_host_ports(monkeypatch: pytest.MonkeyPatch, mock_docker_client: MagicMock) -> None: + """Make host port allocation deterministic: nothing published, every probe free.""" + mock_docker_client.containers.list.return_value = [] + monkeypatch.setattr(ports_mod, "is_port_free", lambda port: True) + + +def _published_host_ports(run_mock: MagicMock) -> list[int]: + return [kwargs["ports"]["8000/tcp"] for _, kwargs in run_mock.call_args_list] + + +@pytest.mark.asyncio +async def test_create_deployment_reallocates_port_after_docker_conflict( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, + free_host_ports: None, +) -> None: + """Docker's port reservations are invisible to the probe, so a publish can still lose a race.""" + first_port = docker_backend._executor_config.port_range_start + leftover = MagicMock() + mock_entities.get.return_value = published_port_config() + mock_docker_client.containers.get.side_effect = [NotFound("missing"), leftover] + mock_docker_client.containers.run.side_effect = [ + _port_conflict_error(first_port), + MagicMock(id="abc123"), + ] + + update = await docker_backend.create_deployment( + workspace="default", + name="srv", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + assert _published_host_ports(mock_docker_client.containers.run) == [first_port, first_port + 1] + # run() creates then starts, so the container that failed to start still holds the name. + leftover.remove.assert_called_once_with(force=True) + + +@pytest.mark.asyncio +async def test_create_deployment_fails_after_repeated_port_conflicts( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, + free_host_ports: None, +) -> None: + first_port = docker_backend._executor_config.port_range_start + mock_entities.get.return_value = published_port_config() + mock_docker_client.containers.get.side_effect = NotFound("missing") + mock_docker_client.containers.run.side_effect = [ + _port_conflict_error(first_port + offset) for offset in range(_PORT_CONFLICT_ATTEMPTS) + ] + + update = await docker_backend.create_deployment( + workspace="default", + name="srv", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "FAILED" + assert _PORT_CONFLICT_MARKER in (update.status_message or "") + published = _published_host_ports(mock_docker_client.containers.run) + assert published == [first_port + offset for offset in range(_PORT_CONFLICT_ATTEMPTS)] + + +@pytest.mark.asyncio +async def test_create_deployment_does_not_retry_unrelated_start_failure( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, + free_host_ports: None, +) -> None: + mock_entities.get.return_value = published_port_config() + mock_docker_client.containers.get.side_effect = NotFound("missing") + mock_docker_client.containers.run.side_effect = APIError("no such image") + + update = await docker_backend.create_deployment( + workspace="default", + name="srv", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "FAILED" + mock_docker_client.containers.run.assert_called_once() + + @pytest.mark.asyncio async def test_create_lora_group_runs_init_server_and_sidecar( docker_backend: DockerDeploymentBackend, @@ -711,6 +824,296 @@ async def test_list_managed_deployment_names( assert names == ["default/srv"] +def _one_shot_server_container( + *, + restart_policy: str, + status: str = "exited", + exit_code: int = 0, + restart_count: int = 0, + backoff_limit: str = "6", +) -> MagicMock: + container = MagicMock() + container.name = container_name("default", "job") + container.wait.return_value = {"StatusCode": exit_code} + container.labels = { + RESTART_POLICY_LABEL: restart_policy, + BACKOFF_LIMIT_LABEL: backoff_limit, + } + container.status = status + container.attrs = { + **container_attrs(exit_code=exit_code), + "RestartCount": restart_count, + } + return container + + +@pytest.mark.asyncio +async def test_create_never_job_returns_succeeded_when_container_exits_immediately( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="Never") + mock_docker_client.containers.get.side_effect = NotFound("missing") + mock_docker_client.containers.run.return_value = _one_shot_server_container( + restart_policy="Never", + exit_code=0, + ) + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "SUCCEEDED" + assert update.exit_code == 0 + mock_docker_client.containers.run.return_value.wait.assert_called_once_with(timeout=5) + + +@pytest.mark.asyncio +async def test_create_never_job_uses_configured_oneshot_observe_timeout( + mock_sdk: MagicMock, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + with ( + patch("nemo_deployments_plugin.backends.docker.backend.AsyncEntitiesResource"), + patch("nemo_deployments_plugin.backends.docker.backend.NemoEntitiesClient", return_value=mock_entities), + patch("nemo_deployments_plugin.backends.docker.backend.get_shared_gpu_pool", return_value=None), + patch("docker.from_env", return_value=mock_docker_client), + ): + backend = DockerDeploymentBackend( + mock_sdk, + {"docker_timeout": 600, "oneshot_observe_timeout_seconds": 7, "pull_images": False}, + ) + backend._client = mock_docker_client + + mock_entities.get.return_value = sample_config(restart_policy="Never") + mock_docker_client.containers.get.side_effect = NotFound("missing") + mock_docker_client.containers.run.return_value = _one_shot_server_container( + restart_policy="Never", + exit_code=0, + ) + + update = await backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "SUCCEEDED" + mock_docker_client.containers.run.return_value.wait.assert_called_once_with(timeout=7) + + +@pytest.mark.asyncio +async def test_create_never_job_returns_failed_on_non_zero_exit( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="Never") + mock_docker_client.containers.get.side_effect = NotFound("missing") + mock_docker_client.containers.run.return_value = _one_shot_server_container( + restart_policy="Never", + exit_code=42, + ) + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "FAILED" + assert update.exit_code == 42 + + +@pytest.mark.asyncio +async def test_create_on_failure_returns_succeeded_when_already_exited_zero( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="OnFailure") + mock_docker_client.containers.get.side_effect = NotFound("missing") + server = _one_shot_server_container(restart_policy="OnFailure", exit_code=0) + mock_docker_client.containers.run.return_value = server + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "SUCCEEDED" + assert update.exit_code == 0 + server.wait.assert_not_called() + + +@pytest.mark.asyncio +async def test_create_on_failure_returns_starting_when_failed_under_backoff( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="OnFailure") + mock_docker_client.containers.get.side_effect = NotFound("missing") + server = _one_shot_server_container( + restart_policy="OnFailure", + exit_code=1, + restart_count=2, + backoff_limit="6", + ) + mock_docker_client.containers.run.return_value = server + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + assert update.exit_code == 1 + assert "retry 2/6" in update.status_message + + +@pytest.mark.asyncio +async def test_create_on_failure_returns_starting_when_still_running( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="OnFailure") + mock_docker_client.containers.get.side_effect = NotFound("missing") + server = _one_shot_server_container(restart_policy="OnFailure", status="running", exit_code=0) + mock_docker_client.containers.run.return_value = server + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + assert "created" in update.status_message.lower() + server.wait.assert_not_called() + + +@pytest.mark.asyncio +async def test_create_never_job_returns_starting_when_wait_times_out( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="Never") + mock_docker_client.containers.get.side_effect = NotFound("missing") + server = _one_shot_server_container(restart_policy="Never", exit_code=0) + server.wait.side_effect = ReadTimeout("timed out") + mock_docker_client.containers.run.return_value = server + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + assert "after observe wait (5s)" in update.status_message + server.wait.assert_called_once_with(timeout=5) + + +@pytest.mark.asyncio +async def test_create_never_job_returns_starting_when_wait_connection_error( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="Never") + mock_docker_client.containers.get.side_effect = NotFound("missing") + server = _one_shot_server_container(restart_policy="Never", exit_code=0) + server.wait.side_effect = RequestsConnectionError("connection reset") + mock_docker_client.containers.run.return_value = server + + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + assert "after observe wait (5s)" in update.status_message + server.wait.assert_called_once_with(timeout=5) + + +@pytest.mark.asyncio +async def test_create_never_job_cleans_up_on_wait_error( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="Never") + mock_docker_client.containers.get.side_effect = NotFound("missing") + server = _one_shot_server_container(restart_policy="Never", exit_code=0) + server.wait.side_effect = RuntimeError("boom") + mock_docker_client.containers.run.return_value = server + + with patch.object( + docker_backend, + "delete_deployment", + new_callable=AsyncMock, + return_value=BackendStatusUpdate(status="SUCCEEDED", status_message="deleted"), + ) as mock_delete: + update = await docker_backend.create_deployment( + workspace="default", + name="job", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "FAILED" + mock_delete.assert_awaited_once_with("default", "job") + + +@pytest.mark.asyncio +async def test_create_always_still_returns_starting( + docker_backend: DockerDeploymentBackend, + mock_entities: AsyncMock, + mock_docker_client: MagicMock, +) -> None: + mock_entities.get.return_value = sample_config(restart_policy="Always") + mock_docker_client.containers.get.side_effect = NotFound("missing") + mock_docker_client.containers.run.return_value = MagicMock(id="abc123") + + update = await docker_backend.create_deployment( + workspace="default", + name="srv", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + mock_docker_client.containers.run.return_value.wait.assert_not_called() + + @pytest.mark.asyncio async def test_default_list_managed_deployment_names_ignores_foreign_scoped_resources( docker_backend: DockerDeploymentBackend, diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/test_executor_config.py b/plugins/nemo-deployments/tests/unit/backends/docker/test_executor_config.py index 395f6dd825..2fa9f416bc 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/test_executor_config.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/test_executor_config.py @@ -9,9 +9,10 @@ def test_docker_executor_config_defaults() -> None: cfg = DockerExecutorConfig() - assert cfg.port_range_start == 9000 - assert cfg.port_range_end == 9100 + assert cfg.port_range_start == 49152 + assert cfg.port_range_end == 49251 assert cfg.resource_scope == DEFAULT_RESOURCE_SCOPE + assert cfg.oneshot_observe_timeout_seconds == 5 def test_docker_executor_config_rejects_inverted_port_range() -> None: diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/test_idempotency.py b/plugins/nemo-deployments/tests/unit/backends/docker/test_idempotency.py index 175f10c021..de1269f03b 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/test_idempotency.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/test_idempotency.py @@ -78,7 +78,14 @@ async def test_create_exited_one_shot_container_is_removed_and_recreated( existing.attrs = container_attrs(status="exited", exit_code=0) mock_docker_client.containers.get.return_value = existing mock_entities.get.return_value = sample_config(restart_policy=restart_policy) - mock_docker_client.containers.run.return_value = MagicMock(id="fresh456") + fresh = MagicMock(id="fresh456") + fresh.labels = _matching_labels(name="srv-puller", restart_policy=restart_policy) + fresh.attrs = container_attrs(exit_code=0) + if restart_policy == "Never": + fresh.wait.return_value = {"StatusCode": 0} + else: + fresh.status = "running" + mock_docker_client.containers.run.return_value = fresh update = await docker_backend.create_deployment( workspace="default", @@ -88,7 +95,13 @@ async def test_create_exited_one_shot_container_is_removed_and_recreated( backend_config={}, ) - assert update.status == "STARTING" + if restart_policy == "Never": + assert update.status == "SUCCEEDED" + assert update.exit_code == 0 + fresh.wait.assert_called_once() + else: + assert update.status == "STARTING" + fresh.wait.assert_not_called() existing.remove.assert_called_once_with(force=True) mock_docker_client.containers.run.assert_called_once() diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/test_ports.py b/plugins/nemo-deployments/tests/unit/backends/docker/test_ports.py index fdf72f26cb..5bf508e64c 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/test_ports.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/test_ports.py @@ -24,47 +24,50 @@ def test_is_port_free_skips_check_for_remote_host(monkeypatch: pytest.MonkeyPatc assert is_port_free(1) is True -def test_is_port_free_returns_false_when_bind_fails(monkeypatch: pytest.MonkeyPatch) -> None: - from nemo_deployments_plugin.backends.docker import ports as ports_mod +class FakeSock: + """Records what ``is_port_free`` does to its probe socket.""" - class FakeSock: - def __enter__(self) -> FakeSock: - return self + def __init__(self, *, bind_error: OSError | None = None) -> None: + self.bind_error = bind_error + self.bound: list[tuple[str, int]] = [] + self.sockopts: list[tuple[object, ...]] = [] - def __exit__(self, *args: object) -> None: - return None + def __enter__(self) -> FakeSock: + return self - def setsockopt(self, *args: object, **kwargs: object) -> None: - return None + def __exit__(self, *args: object) -> None: + return None - def bind(self, addr: tuple[str, int]) -> None: - raise OSError("Address already in use") + def setsockopt(self, *args: object) -> None: + self.sockopts.append(args) - monkeypatch.setattr(ports_mod, "is_remote_docker_host", lambda: False) - monkeypatch.setattr(ports_mod.socket, "socket", lambda *args, **kwargs: FakeSock()) - assert is_port_free(9000) is False + def bind(self, addr: tuple[str, int]) -> None: + if self.bind_error is not None: + raise self.bind_error + self.bound.append(addr) -def test_is_port_free_returns_true_when_bind_succeeds(monkeypatch: pytest.MonkeyPatch) -> None: +def _patch_local_probe(monkeypatch: pytest.MonkeyPatch, sock: FakeSock) -> None: from nemo_deployments_plugin.backends.docker import ports as ports_mod - class FakeSock: - def __enter__(self) -> FakeSock: - return self + monkeypatch.setattr(ports_mod, "is_remote_docker_host", lambda: False) + monkeypatch.setattr(ports_mod.socket, "socket", lambda *args, **kwargs: sock) - def __exit__(self, *args: object) -> None: - return None - def setsockopt(self, *args: object, **kwargs: object) -> None: - return None +def test_is_port_free_returns_false_when_bind_fails(monkeypatch: pytest.MonkeyPatch) -> None: + _patch_local_probe(monkeypatch, FakeSock(bind_error=OSError("Address already in use"))) + assert is_port_free(9000) is False - def bind(self, addr: tuple[str, int]) -> None: - assert addr == ("127.0.0.1", 9000) - return None - monkeypatch.setattr(ports_mod, "is_remote_docker_host", lambda: False) - monkeypatch.setattr(ports_mod.socket, "socket", lambda *args, **kwargs: FakeSock()) +def test_is_port_free_probes_loopback_without_reuseaddr(monkeypatch: pytest.MonkeyPatch) -> None: + # A Docker wildcard publisher blocks this loopback bind. Avoiding SO_REUSEADDR + # is what prevents the probe from incorrectly reporting the port free. + sock = FakeSock() + _patch_local_probe(monkeypatch, sock) + assert is_port_free(9000) is True + assert sock.bound == [("127.0.0.1", 9000)] + assert sock.sockopts == [] @pytest.mark.asyncio @@ -81,6 +84,27 @@ async def test_find_available_port_skips_used(mock_docker_client: MagicMock, mon assert port == 9002 +@pytest.mark.asyncio +async def test_find_available_port_skips_unmanaged_containers( + mock_docker_client: MagicMock, monkeypatch: pytest.MonkeyPatch +) -> None: + """Foreign containers (e.g. a test-fixture ClickHouse) compete for host ports too.""" + from nemo_deployments_plugin.backends.docker import ports as ports_mod + from nemo_deployments_plugin.backends.docker.ports import find_available_port + + foreign = MagicMock() + foreign.labels = {} + foreign.ports = {"9000/tcp": [{"HostPort": "9000"}]} + mock_docker_client.containers.list.return_value = [foreign] + # The daemon's own reservation is invisible to a host-socket probe. + monkeypatch.setattr(ports_mod, "is_port_free", lambda port: True) + + port = await find_available_port(mock_docker_client, 9000, 9002) + + assert port == 9001 + assert mock_docker_client.containers.list.call_args.kwargs == {"all": True} + + @pytest.mark.asyncio async def test_find_available_port_excludes_pending_assignments(mock_docker_client: MagicMock) -> None: from nemo_deployments_plugin.backends.docker.ports import find_available_port