From b3b63e79d6ad45fad34d7530baf23ce9b2b9afd0 Mon Sep 17 00:00:00 2001 From: Matt Kornfield Date: Mon, 27 Jul 2026 18:39:28 +0000 Subject: [PATCH] chore: harden container labeling and e2e tests Signed-off-by: Matt Kornfield --- e2e/services_pool.py | 36 ++- e2e/test_anonymizer_plugin.py | 64 ++++- e2e/test_nemo_agents.py | 11 +- .../tests/unit/test_e2e_harness.py | 19 ++ .../nemo_agents_plugin/api/v2/deployments.py | 41 ++- .../tests/unit/test_deployments_api.py | 78 ++++- .../backends/docker/backend.py | 45 ++- .../backends/docker/config.py | 9 + .../backends/docker/volumes.py | 5 +- .../backends/k8s/jobs.py | 3 + .../backends/labels.py | 25 +- .../backends/docker/test_backend_mocked.py | 269 +++++++++++++++++- .../backends/docker/test_executor_config.py | 2 + .../unit/backends/docker/test_idempotency.py | 3 + .../tests/unit/backends/docker/test_labels.py | 31 ++ 15 files changed, 601 insertions(+), 40 deletions(-) diff --git a/e2e/services_pool.py b/e2e/services_pool.py index c27093c99a..e7b5f1e128 100644 --- a/e2e/services_pool.py +++ b/e2e/services_pool.py @@ -290,7 +290,12 @@ def _register_config_state( def _materialize_config_path(self, state: ModuleConfigState) -> ModuleConfigState: data_dir = e2e_services_data_dir(self._get_log_dir(), state.key.config_hash) - rendered_config_data = _render_e2e_config_for_backend(state.config_data, data_dir, state.harness_config) + rendered_config_data = _render_e2e_config_for_backend( + state.config_data, + data_dir, + state.harness_config, + state.key.config_hash, + ) rendered_config = yaml.safe_dump(rendered_config_data, default_flow_style=False, sort_keys=True) config_path = self._get_generated_config_dir() / f"platform-{state.key.config_hash}.yaml" if not config_path.exists(): @@ -546,7 +551,12 @@ def e2e_services_data_dir(log_dir: Path, config_hash: str) -> Path: return log_dir / f"data-{config_hash}" -def with_e2e_instance_paths(config_data: dict[str, Any], data_dir: Path) -> dict[str, Any]: +def with_e2e_instance_paths( + config_data: dict[str, Any], + data_dir: Path, + *, + resource_scope: str | None = None, +) -> dict[str, Any]: """Return config data with per-instance filesystem paths rooted under ``data_dir``.""" rendered = deepcopy(config_data) subprocess_working_dir = str(data_dir / "subprocess-jobs") @@ -577,9 +587,27 @@ def with_e2e_instance_paths(config_data: dict[str, Any], data_dir: Path) -> dict if isinstance(default_storage_config, dict) and default_storage_config.get("type") == "local": default_storage_config["path"] = files_root + if resource_scope: + _stamp_e2e_docker_deployment_scope(rendered, resource_scope) + return rendered +def _stamp_e2e_docker_deployment_scope(config_data: dict[str, Any], resource_scope: str) -> None: + deployments = config_data.get("deployments") + if not isinstance(deployments, dict): + return + executors = deployments.get("executors") + if not isinstance(executors, list): + return + for executor in executors: + if not isinstance(executor, dict) or executor.get("backend") != "docker": + continue + executor_config = executor.setdefault("config", {}) + if isinstance(executor_config, dict): + executor_config.setdefault("resource_scope", resource_scope) + + # The authz e2e suite (``e2e/authz_oidc``) editable-installs intentionally-broken # fixture plugins (named ``harness-*``) into the shared venv to exercise the authz # fail-modes. Their ``nemo.services`` entry points persist for the rest of the pytest @@ -610,11 +638,11 @@ def _e2e_backend(harness_config: E2EHarnessConfig) -> Literal["subprocess", "doc def _render_e2e_config_for_backend( - config_data: dict[str, Any], data_dir: Path, harness_config: E2EHarnessConfig + config_data: dict[str, Any], data_dir: Path, harness_config: E2EHarnessConfig, config_hash: str ) -> dict[str, Any]: if _e2e_backend(harness_config) in {"docker", "docker_compose"}: return deepcopy(config_data) - return with_e2e_instance_paths(config_data, data_dir) + return with_e2e_instance_paths(config_data, data_dir, resource_scope=f"e2e-{config_hash}") class DockerBackendOverrides(TypedDict, total=False): diff --git a/e2e/test_anonymizer_plugin.py b/e2e/test_anonymizer_plugin.py index 1f4d8de75f..2a58030710 100644 --- a/e2e/test_anonymizer_plugin.py +++ b/e2e/test_anonymizer_plugin.py @@ -14,6 +14,7 @@ from collections.abc import Iterator from contextlib import suppress from pathlib import Path +from typing import cast import data_designer.config as dd import httpx @@ -28,7 +29,11 @@ AnonymizerConfigValidationError, AnonymizerPreviewError, ) -from nemo_anonymizer_plugin.sdk.job_resources import TERMINAL_INCOMPLETE_STATUSES, AnonymizerJobResource +from nemo_anonymizer_plugin.sdk.job_resources import ( + MAX_CONSECUTIVE_POLL_ERRORS, + TERMINAL_INCOMPLETE_STATUSES, + AnonymizerJobResource, +) from nemo_anonymizer_plugin.sdk.resources import AnonymizerPreviewResult from nemo_platform import NeMoPlatform from nemo_platform_plugin.files.client import FilesClient @@ -113,13 +118,20 @@ def _workspace_client(sdk: NeMoPlatform, workspace: str) -> NeMoPlatform: ) +def _require_workspace(workspace: str | None) -> str: + assert workspace is not None + return workspace + + def _anonymizer_url(sdk: NeMoPlatform, workspace: str, path: str) -> str: return f"{str(sdk.base_url).rstrip('/')}/apis/anonymizer/v2/workspaces/{workspace}/{path.lstrip('/')}" -def _raw_anonymizer_post(sdk: NeMoPlatform, workspace: str, path: str, payload: dict[str, object]) -> httpx.Response: +def _raw_anonymizer_post( + sdk: NeMoPlatform, workspace: str | None, path: str, payload: dict[str, object] +) -> httpx.Response: return sdk._client.post( - _anonymizer_url(sdk, workspace, path), + _anonymizer_url(sdk, _require_workspace(workspace), path), json=payload, headers=_string_headers(sdk), timeout=sdk.timeout, @@ -166,12 +178,12 @@ def _rewrite_config() -> AnonymizerConfig: ) -def _fileset_ref(workspace: str, fileset: str, path: str) -> str: - return f"{workspace}/{fileset}#{path}" +def _fileset_ref(workspace: str | None, fileset: str, path: str) -> str: + return f"{_require_workspace(workspace)}/{fileset}#{path}" -def _fileset_uri_ref(workspace: str, fileset: str, path: str) -> str: - return f"fileset://{workspace}/{fileset}#{path}" +def _fileset_uri_ref(workspace: str | None, fileset: str, path: str) -> str: + return f"fileset://{_require_workspace(workspace)}/{fileset}#{path}" def _input_spec(source: str, *, text_column: str = TEXT_COLUMN) -> AnonymizerInputSpec: @@ -281,14 +293,32 @@ def _job_name(job: AnonymizerJobResource) -> str: def _wait_for_anonymizer_job(job: AnonymizerJobResource, *, timeout_seconds: float) -> None: deadline = time.monotonic() + timeout_seconds - status = job.get_job_status() + status = None + consecutive_poll_errors = 0 + last_poll_error: Exception | None = None while status not in {"completed", *TERMINAL_INCOMPLETE_STATUSES}: if time.monotonic() >= deadline: - logs = job.get_logs() - tail = logs[-5:] if logs else [] - raise TimeoutError(f"Anonymizer job {_job_name(job)} timed out with status {status!r}; logs={tail!r}") - time.sleep(ANONYMIZER_POLL_INTERVAL_SECONDS) - status = job.get_job_status() + try: + logs = job.get_logs() + tail = logs[-5:] if logs else [] + except Exception as exc: + tail = [f""] + raise TimeoutError( + f"Anonymizer job {_job_name(job)} timed out with status {status!r}; " + f"last_poll_error={last_poll_error!r}; logs={tail!r}" + ) + try: + status = job.get_job_status() + except Exception as exc: + consecutive_poll_errors += 1 + last_poll_error = exc + if consecutive_poll_errors >= MAX_CONSECUTIVE_POLL_ERRORS: + raise + else: + consecutive_poll_errors = 0 + last_poll_error = None + if status not in {"completed", *TERMINAL_INCOMPLETE_STATUSES}: + time.sleep(ANONYMIZER_POLL_INTERVAL_SECONDS) assert status == "completed" @@ -433,7 +463,11 @@ def test_mock_provider_chat_completion_works_through_minikube_ingress( }, ) - assert SUBSTITUTE_NAME in response["choices"][0]["message"]["content"] + choices = cast(list[dict[str, object]], response["choices"]) + message = cast(dict[str, object], choices[0]["message"]) + content = message["content"] + assert isinstance(content, str) + assert SUBSTITUTE_NAME in content def test_file_upload_round_trips_through_minikube_ingress( @@ -558,7 +592,7 @@ def test_preview_missing_text_column_is_rejected( def test_preview_invalid_strategy_payload_is_rejected(anonymizer_sdk: NeMoPlatform, anonymizer_fileset: str) -> None: - payload = { + payload: dict[str, object] = { "config": {"replace": {"kind": "explode"}, "emit_telemetry": False}, "data": { "source": _fileset_ref(anonymizer_sdk.workspace, anonymizer_fileset, CSV_REMOTE_PATH), diff --git a/e2e/test_nemo_agents.py b/e2e/test_nemo_agents.py index 110175cac1..3464c443f9 100644 --- a/e2e/test_nemo_agents.py +++ b/e2e/test_nemo_agents.py @@ -134,8 +134,15 @@ def _delete_deployment_if_exists(sdk: NeMoPlatform, *, workspace: str, name: str try: sdk.agents.deployments.delete(name, workspace=workspace) except httpx.HTTPStatusError as exc: - if exc.response.status_code != 404: - raise + if exc.response.status_code == 404: + return + if exc.response.status_code in {409, 500}: + try: + sdk.agents.deployments.get(name, workspace=workspace) + except httpx.HTTPStatusError as get_exc: + if get_exc.response.status_code == 404: + return + raise def _get_deployment_log_text(sdk: NeMoPlatform, *, workspace: str, name: str) -> str: diff --git a/packages/nmp_testing/tests/unit/test_e2e_harness.py b/packages/nmp_testing/tests/unit/test_e2e_harness.py index 95e28ff52d..9bb3f7eff3 100644 --- a/packages/nmp_testing/tests/unit/test_e2e_harness.py +++ b/packages/nmp_testing/tests/unit/test_e2e_harness.py @@ -64,3 +64,22 @@ def test_with_e2e_instance_paths_namespaces_local_filesystem_paths(tmp_path): jobs_config = cast(dict[str, Any], config_data["jobs"]) executors = cast(list[dict[str, Any]], jobs_config["executors"]) assert executors[0]["config"]["working_directory"] == ".tmp/e2e/subprocess-jobs" + + +def test_with_e2e_instance_paths_scopes_docker_deployments_executor(tmp_path): + data_dir = tmp_path / "data-abc123def456" + config_data: dict[str, Any] = { + "deployments": { + "executors": [ + {"name": "local-docker", "backend": "docker", "config": {"pull_images": False}}, + {"name": "local-k8s", "backend": "k8s", "config": {}}, + ], + }, + } + + rendered = services_pool.with_e2e_instance_paths(config_data, data_dir, resource_scope="e2e-abc123def456") + + deployments = cast(dict[str, Any], rendered["deployments"]) + executors = cast(list[dict[str, Any]], deployments["executors"]) + assert executors[0]["config"]["resource_scope"] == "e2e-abc123def456" + assert "resource_scope" not in executors[1]["config"] diff --git a/plugins/nemo-agents/src/nemo_agents_plugin/api/v2/deployments.py b/plugins/nemo-agents/src/nemo_agents_plugin/api/v2/deployments.py index a0cbba077f..03efa7ba67 100644 --- a/plugins/nemo-agents/src/nemo_agents_plugin/api/v2/deployments.py +++ b/plugins/nemo-agents/src/nemo_agents_plugin/api/v2/deployments.py @@ -49,6 +49,7 @@ router = APIRouter() _deployment_filter_dep = make_filter_obj_dep(DeploymentFilter) +_DELETE_MARK_ATTEMPTS = 3 @router.post("/deployments", response_model=AgentDeployment, status_code=201, tags=["Agent Deployments"]) @@ -206,9 +207,42 @@ async def delete_deployment( Marks the deployment as ``deleting``. The controller terminates the subprocess and removes the entity on the next reconcile cycle. """ + for attempt in range(_DELETE_MARK_ATTEMPTS): + try: + await _mark_deployment_deleting_once( + entity_client, + workspace=workspace, + name=name, + retrying=attempt > 0, + ) + return + except HTTPException: + raise + except NemoEntityConflictError as exc: + if attempt + 1 >= _DELETE_MARK_ATTEMPTS: + raise HTTPException( + status_code=409, + detail=f"Deployment '{name}' is being modified concurrently.", + ) from exc + logger.info("Retrying delete for deployment '%s' after concurrent update", name) + except Exception as exc: + logger.exception("Failed to mark deployment '%s' as deleting", name) + raise HTTPException(status_code=500, detail="Failed to update deployment.") from exc + + +async def _mark_deployment_deleting_once( + entity_client: NemoEntitiesClient, + *, + workspace: str, + name: str, + retrying: bool, +) -> None: try: dep = await entity_client.get(AgentDeployment, name=name, workspace=workspace) except NemoEntityNotFoundError as exc: + if retrying: + logger.info("Deployment '%s' already deleted during delete retry", name) + return raise HTTPException( status_code=404, detail=f"Deployment '{name}' not found in workspace '{workspace}'.", @@ -217,11 +251,12 @@ async def delete_deployment( logger.exception("Failed to look up deployment '%s' before delete", name) raise HTTPException(status_code=500, detail="Failed to look up deployment.") from exc + if dep.status == "deleting": + return + dep.status = "deleting" try: await entity_client.update(dep) except NemoEntityNotFoundError: logger.info("Deployment '%s' already deleted before status update", name) - except Exception as exc: - logger.exception("Failed to mark deployment '%s' as deleting", name) - raise HTTPException(status_code=500, detail="Failed to update deployment.") from exc + return diff --git a/plugins/nemo-agents/tests/unit/test_deployments_api.py b/plugins/nemo-agents/tests/unit/test_deployments_api.py index 1bb1839b32..96d2e0fe83 100644 --- a/plugins/nemo-agents/tests/unit/test_deployments_api.py +++ b/plugins/nemo-agents/tests/unit/test_deployments_api.py @@ -13,7 +13,8 @@ from fastapi.testclient import TestClient from nemo_agents_plugin.api.v2 import deployments as deployments_router_module from nemo_agents_plugin.api.v2.dependencies import get_entity_client -from nemo_agents_plugin.entities import NEMO_AGENTS_SPEC_CONFIG_FORMAT, Agent, AgentDeployment +from nemo_agents_plugin.entities import NEMO_AGENTS_SPEC_CONFIG_FORMAT, Agent, AgentDeployment, DeploymentStatus +from nemo_platform_plugin.entity_client import NemoEntityConflictError, NemoEntityNotFoundError NOW = datetime.now(timezone.utc) @@ -56,6 +57,19 @@ def _make_agent( return agent +def _make_deployment( + *, + name: str = "fabric-dep", + workspace: str = "default", + agent: str = "fabric-agent", + status: DeploymentStatus = "pending", +) -> AgentDeployment: + deployment = AgentDeployment(name=name, workspace=workspace, agent=agent, status=status) + deployment._id = f"deployment-{name}-id" + deployment._created_at = NOW + return deployment + + def _test_client(mock_entity_client: AsyncMock) -> TestClient: app = FastAPI() app.include_router( @@ -106,3 +120,65 @@ def test_create_rejects_invalid_platform_agent_config(self) -> None: assert resp.status_code == 400 assert "Invalid agent config" in resp.json()["detail"] mock_entity_client.create.assert_not_called() + + +class TestDeleteDeployment: + def test_delete_marks_deployment_deleting(self) -> None: + mock_entity_client = AsyncMock() + mock_entity_client.get = AsyncMock(return_value=_make_deployment(status="starting")) + mock_entity_client.update = AsyncMock(return_value=None) + client = _test_client(mock_entity_client) + + resp = client.delete("/apis/agents/v2/workspaces/default/deployments/fabric-dep") + + assert resp.status_code == 204 + updated: AgentDeployment = mock_entity_client.update.call_args[0][0] + assert updated.status == "deleting" + + def test_delete_retries_concurrent_update_conflict(self) -> None: + mock_entity_client = AsyncMock() + mock_entity_client.get = AsyncMock( + side_effect=[ + _make_deployment(status="pending"), + _make_deployment(status="starting"), + ] + ) + mock_entity_client.update = AsyncMock(side_effect=[NemoEntityConflictError("conflict"), None]) + client = _test_client(mock_entity_client) + + resp = client.delete("/apis/agents/v2/workspaces/default/deployments/fabric-dep") + + assert resp.status_code == 204 + assert mock_entity_client.get.await_count == 2 + assert mock_entity_client.update.await_count == 2 + + def test_delete_returns_success_when_entity_disappears_during_retry(self) -> None: + mock_entity_client = AsyncMock() + mock_entity_client.get = AsyncMock( + side_effect=[ + _make_deployment(status="pending"), + NemoEntityNotFoundError("gone"), + ] + ) + mock_entity_client.update = AsyncMock(side_effect=NemoEntityConflictError("conflict")) + client = _test_client(mock_entity_client) + + resp = client.delete("/apis/agents/v2/workspaces/default/deployments/fabric-dep") + + assert resp.status_code == 204 + + def test_delete_returns_409_when_conflicts_exhausted(self) -> None: + mock_entity_client = AsyncMock() + mock_entity_client.get = AsyncMock( + side_effect=[ + _make_deployment(status="pending") + for _ in range(deployments_router_module._DELETE_MARK_ATTEMPTS) # noqa: SLF001 + ] + ) + mock_entity_client.update = AsyncMock(side_effect=NemoEntityConflictError("conflict")) + client = _test_client(mock_entity_client) + + resp = client.delete("/apis/agents/v2/workspaces/default/deployments/fabric-dep") + + assert resp.status_code == 409 + assert mock_entity_client.update.await_count == deployments_router_module._DELETE_MARK_ATTEMPTS # noqa: SLF001 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 43c410457d..94c2aeb60c 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 @@ -47,6 +47,7 @@ DEPLOYMENT_NAME_LABEL, DEPLOYMENT_WORKSPACE_LABEL, MANAGED_BY_KEY, + RESOURCE_SCOPE_LABEL, RESTART_POLICY_LABEL, companion_container_name, container_name, @@ -245,6 +246,7 @@ async def create_deployment( config.restart_policy, config_name=config_name, backoff_limit=config.backoff_limit, + resource_scope=self._executor_config.resource_scope, ), } @@ -391,7 +393,8 @@ async def _run_init_container( # Best-effort clean any stale init container from a prior attempt. try: stale = await asyncio.to_thread(self._client.containers.get, init_name) - await asyncio.to_thread(stale.remove, force=True) + if self._container_matches_deployment_group(stale, workspace, name): + await asyncio.to_thread(stale.remove, force=True) except self._docker_errors.NotFound: pass except Exception: @@ -448,6 +451,9 @@ async def read_status(self, *, workspace: str, name: str) -> BackendStatusUpdate c_name = container_name(workspace, name) try: container = await asyncio.to_thread(self._client.containers.get, c_name) + if not self._container_matches_deployment_group(container, workspace, name): + restart_policy = await self._resolve_restart_policy(workspace, name) + return missing_container_status(restart_policy, container_name=c_name) await asyncio.to_thread(container.reload) except self._docker_errors.NotFound: restart_policy = await self._resolve_restart_policy(workspace, name) @@ -586,6 +592,7 @@ def _list() -> list[Any]: f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}", f"{DEPLOYMENT_WORKSPACE_LABEL}={workspace}", f"{DEPLOYMENT_NAME_LABEL}={name}", + *self._resource_scope_filter_label(), ] }, ) @@ -599,6 +606,8 @@ def _list() -> list[Any]: running_roles: set[str] = set() for container in containers: + if not self._container_matches_deployment_group(container, workspace, name): + continue role = (container.labels or {}).get(CONTAINER_ROLE_LABEL, "") if role not in expected_roles: continue @@ -628,16 +637,20 @@ def _delete() -> None: f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}", f"{DEPLOYMENT_WORKSPACE_LABEL}={workspace}", f"{DEPLOYMENT_NAME_LABEL}={name}", + *self._resource_scope_filter_label(), ] }, ): - group[container.name] = container + if self._container_matches_deployment_group(container, workspace, name): + group[container.name] = container except Exception: logger.warning("Failed to list group containers for %s; falling back to primary", c_name, exc_info=True) # Ensure the primary is included even if the label list query missed it. if c_name not in group: try: - group[c_name] = self._client.containers.get(c_name) + primary = self._client.containers.get(c_name) + if self._container_matches_deployment_group(primary, workspace, name): + group[c_name] = primary except self._docker_errors.NotFound: pass for container in group.values(): @@ -662,7 +675,7 @@ async def list_managed_deployment_names(self) -> list[str]: containers = await asyncio.to_thread( self._client.containers.list, all=True, - filters=managed_by_filter(), + filters=managed_by_filter(resource_scope=self._executor_config.resource_scope), ) except Exception: logger.warning("Failed to list managed containers", exc_info=True) @@ -673,17 +686,36 @@ async def list_managed_deployment_names(self) -> list[str]: container_labels = container.labels or {} if container_labels.get(MANAGED_BY_KEY) != MANAGED_BY_LABEL: continue + if not self._labels_match_resource_scope(container_labels): + continue ws = container_labels.get(DEPLOYMENT_WORKSPACE_LABEL) dep_name = container_labels.get(DEPLOYMENT_NAME_LABEL) if ws and dep_name: seen.add(f"{ws}/{dep_name}") return sorted(seen) + def _resource_scope_filter_label(self) -> list[str]: + return [f"{RESOURCE_SCOPE_LABEL}={self._executor_config.resource_scope}"] + + def _labels_match_resource_scope(self, labels: dict[str, Any]) -> bool: + return labels.get(RESOURCE_SCOPE_LABEL) == self._executor_config.resource_scope + + def _container_matches_deployment_group(self, container: DockerContainer, workspace: str, name: str) -> bool: + labels = container.labels or {} + return ( + labels.get(DEPLOYMENT_WORKSPACE_LABEL) == workspace + and labels.get(DEPLOYMENT_NAME_LABEL) == name + and labels.get(MANAGED_BY_KEY) == MANAGED_BY_LABEL + and self._labels_match_resource_scope(labels) + ) + async def get_logs(self, *, workspace: str, name: str, tail: int = 100) -> LogResult: c_name = container_name(workspace, name) def _logs() -> bytes: container = self._client.containers.get(c_name) + if not self._container_matches_deployment_group(container, workspace, name): + raise self._docker_errors.NotFound(f"Container {c_name} not found") return container.logs(tail=tail, timestamps=True) try: @@ -729,6 +761,7 @@ async def create_volume( driver=driver, init_chmod=init_chmod, init_image=init_image, + resource_scope=self._executor_config.resource_scope, ) async def read_volume_status( @@ -760,10 +793,8 @@ def _container_matches_deployment( ) -> bool: labels = container.labels or {} return ( - labels.get(DEPLOYMENT_WORKSPACE_LABEL) == workspace - and labels.get(DEPLOYMENT_NAME_LABEL) == name + self._container_matches_deployment_group(container, workspace, name) and labels.get(CONFIG_NAME_LABEL) == config_name - and labels.get(MANAGED_BY_KEY) == MANAGED_BY_LABEL ) async def _pull_image(self, image: str, *, ngc_api_key: str | None) -> str | None: 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 c066b07e59..f27dfca0ae 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 @@ -5,6 +5,7 @@ from __future__ import annotations +from nemo_deployments_plugin.backends.labels import DEFAULT_RESOURCE_SCOPE from pydantic import BaseModel, Field, model_validator @@ -18,6 +19,14 @@ class DockerExecutorConfig(BaseModel): description="Docker client timeout in seconds for pull/create/status operations (default: 10 minutes).", ) pull_images: bool = Field(default=True, description="Pull container images before run when missing locally.") + resource_scope: str = Field( + default=DEFAULT_RESOURCE_SCOPE, + min_length=1, + description=( + "Owner scope label for Docker-managed resources. Orphan cleanup only lists resources with the " + "same scope, allowing multiple platform instances to share one Docker daemon." + ), + ) port_range_start: int = Field( default=9000, ge=1, diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/volumes.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/volumes.py index df849df7ef..117c315d1d 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/volumes.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/docker/volumes.py @@ -10,7 +10,7 @@ from typing import Any from nemo_deployments_plugin.backends.base import VolumeStatusUpdate -from nemo_deployments_plugin.backends.labels import docker_volume_name, volume_identity_labels +from nemo_deployments_plugin.backends.labels import DEFAULT_RESOURCE_SCOPE, docker_volume_name, volume_identity_labels import docker @@ -25,9 +25,10 @@ async def create_volume( driver: str = "local", init_chmod: str | None = None, init_image: str | None = None, + resource_scope: str = DEFAULT_RESOURCE_SCOPE, ) -> VolumeStatusUpdate: vol_name = docker_volume_name(workspace, name) - labels = volume_identity_labels(workspace, name) + labels = volume_identity_labels(workspace, name, resource_scope=resource_scope) def _create() -> bool: """Create the volume if missing. Returns True if newly created.""" diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/k8s/jobs.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/k8s/jobs.py index f03e4e87a5..bf8cdb05ef 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/k8s/jobs.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/k8s/jobs.py @@ -30,9 +30,11 @@ ) from nemo_deployments_plugin.backends.labels import ( CONFIG_NAME_LABEL, + DEFAULT_RESOURCE_SCOPE, DEPLOYMENT_NAME_LABEL, DEPLOYMENT_WORKSPACE_LABEL, MANAGED_BY_KEY, + RESOURCE_SCOPE_LABEL, deployment_identity_labels, k8s_deployment_configmap_name, k8s_deployment_resource_name, @@ -74,6 +76,7 @@ def deployment_scope_labels(workspace: str, name: str) -> dict[str, str]: MANAGED_BY_KEY: MANAGED_BY_LABEL, DEPLOYMENT_WORKSPACE_LABEL: workspace, DEPLOYMENT_NAME_LABEL: name, + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, } diff --git a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/labels.py b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/labels.py index 99ace6ebe1..38157f07ba 100644 --- a/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/labels.py +++ b/plugins/nemo-deployments/src/nemo_deployments_plugin/backends/labels.py @@ -14,6 +14,8 @@ resource kind without ambiguous selectors, even though the workspace value is the same string. """ +from typing import Any + from nemo_deployments_plugin.constants import MANAGED_BY_LABEL from nemo_platform_plugin.k8s_naming import k8s_safe_name, workspace_name_identity @@ -25,6 +27,8 @@ VOLUME_WORKSPACE_LABEL = "nemo.nvidia.com/volume-workspace" VOLUME_NAME_LABEL = "nemo.nvidia.com/volume-name" BACKOFF_LIMIT_LABEL = "nemo.nvidia.com/backoff-limit" +RESOURCE_SCOPE_LABEL = "nemo.nvidia.com/resource-scope" +DEFAULT_RESOURCE_SCOPE = "default" # Role of a docker container within a multi-container deployment group # (e.g. "server" or a sidecar container name). Used to discover companion # containers (LoRA adapters sidecar) that share the primary server container. @@ -95,6 +99,7 @@ def deployment_identity_labels( *, config_name: str, backoff_limit: int = 6, + resource_scope: str = DEFAULT_RESOURCE_SCOPE, ) -> dict[str, str]: """Return identity labels attached to deployment backend resources.""" return { @@ -104,23 +109,33 @@ def deployment_identity_labels( RESTART_POLICY_LABEL: restart_policy, CONFIG_NAME_LABEL: config_name, BACKOFF_LIMIT_LABEL: str(backoff_limit), + RESOURCE_SCOPE_LABEL: resource_scope, } -def volume_identity_labels(workspace: str, name: str) -> dict[str, str]: +def volume_identity_labels( + workspace: str, name: str, *, resource_scope: str = DEFAULT_RESOURCE_SCOPE +) -> dict[str, str]: """Return identity labels attached to volume backend resources.""" return { MANAGED_BY_KEY: MANAGED_BY_LABEL, VOLUME_WORKSPACE_LABEL: workspace, VOLUME_NAME_LABEL: name, + RESOURCE_SCOPE_LABEL: resource_scope, } -def managed_by_filter() -> dict[str, str]: +def managed_by_filter(*, resource_scope: str | None = None) -> dict[str, Any]: """Return a Docker SDK filter dict for plugin-managed resources.""" - return {"label": f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}"} + labels = [f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}"] + if resource_scope is not None: + labels.append(f"{RESOURCE_SCOPE_LABEL}={resource_scope}") + return {"label": labels if len(labels) > 1 else labels[0]} -def managed_by_label_selector() -> str: +def managed_by_label_selector(*, resource_scope: str | None = None) -> str: """Kubernetes label selector for plugin-managed resources.""" - return f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}" + selector = f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}" + if resource_scope is not None: + selector = f"{selector},{RESOURCE_SCOPE_LABEL}={resource_scope}" + return selector 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 9be1747f83..6b79ca41cc 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 @@ -5,7 +5,7 @@ from __future__ import annotations -from unittest.mock import AsyncMock, MagicMock +from unittest.mock import AsyncMock, MagicMock, patch import pytest from backends.docker.docker_helpers import container_attrs, lora_config, sample_config @@ -14,8 +14,11 @@ from nemo_deployments_plugin.backends.labels import ( CONFIG_NAME_LABEL, CONTAINER_ROLE_LABEL, + DEFAULT_RESOURCE_SCOPE, DEPLOYMENT_NAME_LABEL, DEPLOYMENT_WORKSPACE_LABEL, + MANAGED_BY_KEY, + RESOURCE_SCOPE_LABEL, RESTART_POLICY_LABEL, companion_container_name, container_name, @@ -243,8 +246,20 @@ async def test_delete_removes_whole_group( """delete_deployment stops+removes every container in the group.""" server = MagicMock() server.name = container_name("default", "srv") + server.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, + } sidecar = MagicMock() sidecar.name = companion_container_name("default", "srv", "lora-adapters") + sidecar.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, + } mock_docker_client.containers.list.return_value = [server, sidecar] update = await docker_backend.delete_deployment("default", "srv") @@ -254,6 +269,108 @@ async def test_delete_removes_whole_group( sidecar.remove.assert_called_once() +@pytest.mark.asyncio +async def test_delete_scoped_deployment_does_not_remove_foreign_primary( + 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": 60, "pull_images": False, "resource_scope": "e2e-abc123"}, + ) + + foreign = MagicMock() + foreign.name = container_name("default", "srv") + foreign.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: "other-scope", + } + mock_docker_client.containers.list.return_value = [] + mock_docker_client.containers.get.return_value = foreign + + update = await backend.delete_deployment("default", "srv") + + assert update.status == "SUCCEEDED" + foreign.stop.assert_not_called() + foreign.remove.assert_not_called() + + +@pytest.mark.asyncio +async def test_delete_default_scoped_deployment_does_not_remove_foreign_primary( + docker_backend: DockerDeploymentBackend, + mock_docker_client: MagicMock, +) -> None: + scoped = MagicMock() + scoped.name = container_name("default", "srv") + scoped.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: "e2e-abc123", + } + mock_docker_client.containers.list.return_value = [scoped] + mock_docker_client.containers.get.return_value = scoped + + update = await docker_backend.delete_deployment("default", "srv") + + assert update.status == "SUCCEEDED" + scoped.stop.assert_not_called() + scoped.remove.assert_not_called() + + +@pytest.mark.asyncio +async def test_create_lora_group_does_not_remove_foreign_stale_init_container( + 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": 60, "pull_images": False, "resource_scope": "e2e-abc123"}, + ) + + foreign_stale = MagicMock() + foreign_stale.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: "other-scope", + } + init_container = MagicMock() + init_container.wait.return_value = {"StatusCode": 0} + server_container = MagicMock(id="server123") + sidecar_container = MagicMock(id="sidecar123") + mock_entities.get.return_value = lora_config() + mock_docker_client.containers.get.side_effect = [NotFound("missing"), foreign_stale] + mock_docker_client.containers.run.side_effect = [init_container, server_container, sidecar_container] + + update = await backend.create_deployment( + workspace="default", + name="srv", + config_name="cfg1", + labels={"managed-by": MANAGED_BY_LABEL}, + backend_config={}, + ) + + assert update.status == "STARTING" + foreign_stale.remove.assert_not_called() + + @pytest.mark.asyncio async def test_create_volume_runs_init_chmod_container( docker_backend: DockerDeploymentBackend, @@ -275,6 +392,8 @@ async def test_create_volume_runs_init_chmod_container( ) assert update.status == "BOUND" + labels = mock_docker_client.volumes.create.call_args.kwargs["labels"] + assert labels[RESOURCE_SCOPE_LABEL] == DEFAULT_RESOURCE_SCOPE mock_docker_client.containers.run.assert_called_once() args, run_kwargs = mock_docker_client.containers.run.call_args assert args[0] == "docker.io/library/busybox" @@ -340,8 +459,11 @@ async def test_read_status_ready_when_running_without_probe( container.status = "running" container.labels = { "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", RESTART_POLICY_LABEL: "Always", CONFIG_NAME_LABEL: "cfg1", + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, } container.ports = {} container.attrs = container_attrs() @@ -359,8 +481,11 @@ def _running_server_container() -> MagicMock: container.status = "running" container.labels = { "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", RESTART_POLICY_LABEL: "Always", CONFIG_NAME_LABEL: "cfg1", + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, } container.ports = {} container.attrs = container_attrs() @@ -375,6 +500,7 @@ def _sidecar_container(role: str, status: str) -> MagicMock: DEPLOYMENT_WORKSPACE_LABEL: "default", DEPLOYMENT_NAME_LABEL: "srv", CONTAINER_ROLE_LABEL: role, + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, } return sidecar @@ -461,6 +587,48 @@ async def get_side_effect(entity_type, name, workspace=None): assert update.status == "LOST" +@pytest.mark.asyncio +async def test_read_status_treats_foreign_container_as_missing( + 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": 60, "pull_images": False, "resource_scope": "e2e-abc123"}, + ) + + foreign = MagicMock() + foreign.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: "other-scope", + } + mock_docker_client.containers.get.return_value = foreign + + deployment_entity = MagicMock() + deployment_entity.deployment_config = "cfg1" + + async def get_side_effect(entity_type, name, workspace=None): + if entity_type is Deployment: + return deployment_entity + return sample_config(restart_policy="Always") + + mock_entities.get.side_effect = get_side_effect + + update = await backend.read_status(workspace="default", name="srv") + + assert update.status == "LOST" + foreign.reload.assert_not_called() + + @pytest.mark.asyncio async def test_read_status_unknown_on_transient_docker_error( docker_backend: DockerDeploymentBackend, @@ -496,6 +664,7 @@ async def test_list_managed_deployment_names( "managed-by": MANAGED_BY_LABEL, DEPLOYMENT_WORKSPACE_LABEL: "default", DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, } mock_docker_client.containers.list.return_value = [container] @@ -504,6 +673,104 @@ async def test_list_managed_deployment_names( assert names == ["default/srv"] +@pytest.mark.asyncio +async def test_default_list_managed_deployment_names_ignores_foreign_scoped_resources( + docker_backend: DockerDeploymentBackend, + mock_docker_client: MagicMock, +) -> None: + default_scoped = MagicMock() + default_scoped.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, + } + scoped = MagicMock() + scoped.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "foreign", + RESOURCE_SCOPE_LABEL: "e2e-abc123", + } + mock_docker_client.containers.list.return_value = [default_scoped, scoped] + + names = await docker_backend.list_managed_deployment_names() + + assert names == ["default/srv"] + + +@pytest.mark.asyncio +async def test_list_managed_deployment_names_scopes_docker_query( + 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": 60, "pull_images": False, "resource_scope": "e2e-abc123"}, + ) + + container = MagicMock() + container.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: "e2e-abc123", + } + mock_docker_client.containers.list.return_value = [container] + + names = await backend.list_managed_deployment_names() + + assert names == ["default/srv"] + mock_docker_client.containers.list.assert_called_once_with( + all=True, + filters={ + "label": [ + f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}", + f"{RESOURCE_SCOPE_LABEL}=e2e-abc123", + ] + }, + ) + + +@pytest.mark.asyncio +async def test_get_logs_treats_foreign_container_as_missing( + 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": 60, "pull_images": False, "resource_scope": "e2e-abc123"}, + ) + + foreign = MagicMock() + foreign.labels = { + "managed-by": MANAGED_BY_LABEL, + DEPLOYMENT_WORKSPACE_LABEL: "default", + DEPLOYMENT_NAME_LABEL: "srv", + RESOURCE_SCOPE_LABEL: "other-scope", + } + mock_docker_client.containers.get.return_value = foreign + + logs = await backend.get_logs(workspace="default", name="srv") + + assert logs.lines == [f"Container {container_name('default', 'srv')} not found"] + foreign.logs.assert_not_called() + + @pytest.mark.asyncio async def test_create_volume_bound( 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 53268dd8a1..395f6dd825 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 @@ -3,6 +3,7 @@ import pytest from nemo_deployments_plugin.backends.docker.config import DockerExecutorConfig +from nemo_deployments_plugin.backends.labels import DEFAULT_RESOURCE_SCOPE from pydantic import ValidationError @@ -10,6 +11,7 @@ def test_docker_executor_config_defaults() -> None: cfg = DockerExecutorConfig() assert cfg.port_range_start == 9000 assert cfg.port_range_end == 9100 + assert cfg.resource_scope == DEFAULT_RESOURCE_SCOPE 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 59fdbdbf69..175f10c021 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/test_idempotency.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/test_idempotency.py @@ -13,8 +13,10 @@ from nemo_deployments_plugin.backends.docker.backend import DockerDeploymentBackend from nemo_deployments_plugin.backends.labels import ( CONFIG_NAME_LABEL, + DEFAULT_RESOURCE_SCOPE, DEPLOYMENT_NAME_LABEL, DEPLOYMENT_WORKSPACE_LABEL, + RESOURCE_SCOPE_LABEL, RESTART_POLICY_LABEL, ) from nemo_deployments_plugin.constants import MANAGED_BY_LABEL @@ -28,6 +30,7 @@ def _matching_labels(*, name: str, restart_policy: str = "Always", config_name: DEPLOYMENT_NAME_LABEL: name, RESTART_POLICY_LABEL: restart_policy, CONFIG_NAME_LABEL: config_name, + RESOURCE_SCOPE_LABEL: DEFAULT_RESOURCE_SCOPE, } diff --git a/plugins/nemo-deployments/tests/unit/backends/docker/test_labels.py b/plugins/nemo-deployments/tests/unit/backends/docker/test_labels.py index 9466919130..cf19162168 100644 --- a/plugins/nemo-deployments/tests/unit/backends/docker/test_labels.py +++ b/plugins/nemo-deployments/tests/unit/backends/docker/test_labels.py @@ -9,12 +9,16 @@ from nemo_deployments_plugin.backends.labels import ( CONFIG_NAME_LABEL, + DEFAULT_RESOURCE_SCOPE, DEPLOYMENT_NAME_LABEL, DEPLOYMENT_WORKSPACE_LABEL, MANAGED_BY_KEY, + RESOURCE_SCOPE_LABEL, container_name, deployment_identity_labels, docker_volume_name, + managed_by_filter, + managed_by_label_selector, ) from nemo_deployments_plugin.constants import MANAGED_BY_LABEL from nemo_platform_plugin.k8s_naming import k8s_safe_name @@ -80,3 +84,30 @@ def test_deployment_identity_labels() -> None: assert labels[DEPLOYMENT_WORKSPACE_LABEL] == "default" assert labels[DEPLOYMENT_NAME_LABEL] == "srv" assert labels[CONFIG_NAME_LABEL] == "cfg1" + assert labels[RESOURCE_SCOPE_LABEL] == DEFAULT_RESOURCE_SCOPE + + +def test_deployment_identity_labels_can_override_resource_scope() -> None: + labels = deployment_identity_labels( + "default", + "srv", + "Always", + config_name="cfg1", + resource_scope="e2e-abc123", + ) + + assert labels[RESOURCE_SCOPE_LABEL] == "e2e-abc123" + + +def test_managed_by_filters_include_optional_resource_scope() -> None: + assert managed_by_filter() == {"label": f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}"} + assert managed_by_filter(resource_scope="e2e-abc123") == { + "label": [ + f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}", + f"{RESOURCE_SCOPE_LABEL}=e2e-abc123", + ] + } + assert managed_by_label_selector() == f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL}" + assert managed_by_label_selector(resource_scope="e2e-abc123") == ( + f"{MANAGED_BY_KEY}={MANAGED_BY_LABEL},{RESOURCE_SCOPE_LABEL}=e2e-abc123" + )