Skip to content

Commit 6bd8e4d

Browse files
committed
refactor(models): address PR review feedback
Composition over inheritance for the k8s reconcilers (review #2/#3): * Extract StatusProjector (pod-status projection, crash-loop/pending-timeout error builders, host URL) and ResourceDeleter (idempotent 404-tolerant delete) as standalone collaborators. * Reconciler is now a pure interface (the 5 verbs); NimOperatorReconciler and K8sReconciler compose the projector + deleter instead of inheriting them. The backend builds both collaborators in init() and injects them. Thread the reconcile context through the backend interface (review #19): * create/update/get_model_deployment_status now take a single ctx: ModelContext instead of (deployment, config, model_entity); applied across the ServiceBackend ABC and the docker / none / k8s backends, the deployment reconciler call sites, and the test mocks. delete stays (workspace, name). Fixes + nits: * Harden NIMService status read against a null status/state (review #15): (nim_status.get("state") or "").lower() can no longer raise. * Convert nim_operator logging to structured extra={} (review #13); avoid the reserved LogRecord 'name' key (use resource_name / deployment_name). * Flatten the Files-service create/update branches into a guard-clause helper (review #14). * compile_puller_job: rename args -> container_args (review #17). * Reconciler nits: import the vllm_k8s_compiler module under its full name (review #7), reflow the P3 (a)/(b) comment (review #8), quote values in the model-source error (review #9), drop the _ = image_pull_secrets dance (review #12), name the event-message cap MAX_EVENT_MESSAGE_CHARS (review #6), and document the _select_reconciler None contract (review #16). Signed-off-by: Ben McCown <[email protected]>
1 parent cceba77 commit 6bd8e4d

17 files changed

Lines changed: 1116 additions & 911 deletions

File tree

services/core/inference-gateway/tests/integration/conftest.py

Lines changed: 6 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -87,31 +87,19 @@ def init(self) -> None:
8787
"""No-op init for mock backend."""
8888
pass
8989

90-
async def create_model_deployment(
91-
self,
92-
deployment: Any,
93-
config: Any,
94-
model_entity: Any = None,
95-
) -> DeploymentStatusUpdate:
90+
async def create_model_deployment(self, ctx: Any) -> DeploymentStatusUpdate:
9691
"""Record call and return configured response."""
97-
self.create_calls.append((deployment, config, model_entity))
92+
self.create_calls.append((ctx.model_deployment, ctx.model_deployment_config, ctx.model_entity))
9893
return self.create_response
9994

100-
async def update_model_deployment(
101-
self,
102-
deployment: Any,
103-
config: Any,
104-
model_entity: Any = None,
105-
) -> DeploymentStatusUpdate:
95+
async def update_model_deployment(self, ctx: Any) -> DeploymentStatusUpdate:
10696
"""Record call and return configured response."""
107-
self.update_calls.append((deployment, config, model_entity))
97+
self.update_calls.append((ctx.model_deployment, ctx.model_deployment_config, ctx.model_entity))
10898
return self.create_response
10999

110-
async def get_model_deployment_status(
111-
self, deployment: Any, config: Any = None, model_entity: Any = None
112-
) -> DeploymentStatusUpdate:
100+
async def get_model_deployment_status(self, ctx: Any) -> DeploymentStatusUpdate:
113101
"""Record call and return configured response."""
114-
self.status_calls.append(deployment)
102+
self.status_calls.append(ctx.model_deployment)
115103
return self.status_response
116104

117105
async def delete_model_deployment(self, deployment: Any) -> DeploymentStatusUpdate:

services/core/models/src/nmp/core/models/controllers/backends/backends.py

Lines changed: 14 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -4,13 +4,11 @@
44
"""Base backend interface for Models Controller service."""
55

66
from abc import ABC, abstractmethod
7-
from typing import Any, Dict, Optional
7+
from typing import Any, Dict
88

99
from nemo_platform import AsyncNeMoPlatform
1010
from nemo_platform.types.inference import ModelDeploymentStatus
11-
from nemo_platform.types.inference.model_deployment import ModelDeployment
12-
from nemo_platform.types.inference.model_deployment_config import ModelDeploymentConfig
13-
from nemo_platform.types.models.model_entity import ModelEntity
11+
from nmp.core.models.controllers.context import ModelContext
1412
from pydantic import BaseModel
1513

1614

@@ -67,15 +65,12 @@ def shutdown(self) -> None:
6765
...
6866

6967
@abstractmethod
70-
async def create_model_deployment(
71-
self, deployment: ModelDeployment, config: ModelDeploymentConfig, model_entity: Optional[ModelEntity] = None
72-
) -> DeploymentStatusUpdate:
68+
async def create_model_deployment(self, ctx: ModelContext) -> DeploymentStatusUpdate:
7369
"""Create a new model deployment.
7470
7571
Args:
76-
deployment: The ModelDeployment object to create
77-
config: The ModelDeploymentConfig for this deployment
78-
model_entity: Optional Model entity from Entity Store (contains peft, artifact, etc.)
72+
ctx: The reconciliation context bundling the ModelDeployment, its
73+
ModelDeploymentConfig, and the optional Model entity.
7974
8075
Returns:
8176
DeploymentStatusUpdate with the current status after creation attempt
@@ -86,15 +81,13 @@ async def create_model_deployment(
8681
...
8782

8883
@abstractmethod
89-
async def update_model_deployment(
90-
self, deployment: ModelDeployment, config: ModelDeploymentConfig, model_entity: Optional[ModelEntity] = None
91-
) -> DeploymentStatusUpdate:
84+
async def update_model_deployment(self, ctx: ModelContext) -> DeploymentStatusUpdate:
9285
"""Update an existing model deployment.
9386
9487
Args:
95-
deployment: The ModelDeployment object with updated configuration
96-
config: The ModelDeploymentConfig for this deployment (may be a new version)
97-
model_entity: Optional Model entity from Entity Store (contains peft, artifact, etc.)
88+
ctx: The reconciliation context bundling the ModelDeployment, its
89+
(possibly new-version) ModelDeploymentConfig, and the optional
90+
Model entity.
9891
9992
Returns:
10093
DeploymentStatusUpdate with the current status after update attempt
@@ -105,20 +98,14 @@ async def update_model_deployment(
10598
...
10699

107100
@abstractmethod
108-
async def get_model_deployment_status(
109-
self,
110-
deployment: ModelDeployment,
111-
config: Optional[ModelDeploymentConfig] = None,
112-
model_entity: Optional[ModelEntity] = None,
113-
) -> DeploymentStatusUpdate:
101+
async def get_model_deployment_status(self, ctx: ModelContext) -> DeploymentStatusUpdate:
114102
"""Get the current status of a model deployment.
115103
116104
Args:
117-
deployment: The ModelDeployment object to check
118-
config: The ModelDeploymentConfig for this deployment. Some backends
119-
need it to advance creation (e.g. the k8s vLLM path emits the
120-
serving Deployment once the weight-puller Job completes).
121-
model_entity: Optional Model entity from Entity Store.
105+
ctx: The reconciliation context bundling the ModelDeployment, its
106+
ModelDeploymentConfig, and the optional Model entity. Some backends
107+
need the config to advance creation (e.g. the k8s vLLM path emits
108+
the serving Deployment once the weight-puller Job completes).
122109
123110
Returns:
124111
DeploymentStatusUpdate with the current deployment status

services/core/models/src/nmp/core/models/controllers/backends/docker/backend.py

Lines changed: 11 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -12,14 +12,12 @@
1212
import asyncio
1313
import os
1414
from logging import getLogger
15-
from typing import Any, Optional
15+
from typing import Any
1616

1717
import httpx
1818
from docker.errors import APIError, NotFound
1919
from nemo_platform import NotFoundError
2020
from nemo_platform.types.inference.model_deployment import ModelDeployment
21-
from nemo_platform.types.inference.model_deployment_config import ModelDeploymentConfig
22-
from nemo_platform.types.models.model_entity import ModelEntity
2321
from nmp.common.config import get_platform_config
2422
from nmp.common.docker.gpu_pool import DockerGPUPool
2523
from nmp.common.resources import SharedResourceManager
@@ -37,6 +35,7 @@
3735
NGC_IMAGE_REGISTRY_USER_NAME,
3836
DockerDeploymentCreationReconciler,
3937
)
38+
from nmp.core.models.controllers.context import ModelContext
4039
from requests.exceptions import ConnectionError as RequestsConnectionError
4140
from requests.exceptions import ReadTimeout
4241
from urllib3.exceptions import ReadTimeoutError as Urllib3ReadTimeoutError
@@ -184,17 +183,15 @@ async def _ensure_ngc_login(self, ngc_api_key: str | None) -> None:
184183
# ServiceBackend CRUD interface
185184
# ==================================================================
186185

187-
async def create_model_deployment(
188-
self,
189-
deployment: ModelDeployment,
190-
config: ModelDeploymentConfig,
191-
model_entity: Optional[ModelEntity] = None,
192-
) -> DeploymentStatusUpdate:
186+
async def create_model_deployment(self, ctx: ModelContext) -> DeploymentStatusUpdate:
193187
"""Create a new model deployment as a Docker container.
194188
195189
Resolves NGC credentials and delegates the multi-stage creation
196190
pipeline to :class:`DockerDeploymentCreationReconciler`.
197191
"""
192+
deployment = ctx.model_deployment
193+
config = ctx.model_deployment_config
194+
model_entity = ctx.model_entity
198195
resolved_ngc_key = await self._resolve_ngc_api_key()
199196
await self._ensure_ngc_login(resolved_ngc_key)
200197

@@ -205,30 +202,22 @@ async def create_model_deployment(
205202
resolved_ngc_key,
206203
)
207204

208-
async def update_model_deployment(
209-
self,
210-
deployment: ModelDeployment,
211-
config: ModelDeploymentConfig,
212-
model_entity: Optional[ModelEntity] = None,
213-
) -> DeploymentStatusUpdate:
205+
async def update_model_deployment(self, ctx: ModelContext) -> DeploymentStatusUpdate:
214206
"""Update a model deployment by recreating the container."""
207+
deployment = ctx.model_deployment
215208
logger.info(f"Updating Docker deployment: {deployment.workspace}/{deployment.name}")
216209
delete_result = await self.delete_model_deployment(deployment.workspace, deployment.name)
217210
if delete_result.status == "ERROR":
218211
return delete_result
219-
return await self.create_model_deployment(deployment, config, model_entity)
212+
return await self.create_model_deployment(ctx)
220213

221-
async def get_model_deployment_status(
222-
self,
223-
deployment: ModelDeployment,
224-
config: Optional[ModelDeploymentConfig] = None,
225-
model_entity: Optional[ModelEntity] = None,
226-
) -> DeploymentStatusUpdate:
214+
async def get_model_deployment_status(self, ctx: ModelContext) -> DeploymentStatusUpdate:
227215
"""Get the status of a Docker model deployment.
228216
229217
While the deployment is still progressing through the creation
230218
pipeline this delegates to the reconciler's ``advance`` method.
231219
"""
220+
deployment = ctx.model_deployment
232221
if self._reconciler.is_deploying(deployment.workspace, deployment.name):
233222
deployment_key = self._reconciler.get_deployment_key(deployment.workspace, deployment.name)
234223
return await self._reconciler.advance(deployment_key)

services/core/models/src/nmp/core/models/controllers/backends/k8s_nim_operator/backend.py

Lines changed: 51 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@
2323
from kubernetes import config as k8s_config
2424
from kubernetes.dynamic import DynamicClient
2525
from nemo_platform.types.inference.model_deployment import ModelDeployment
26-
from nemo_platform.types.inference.model_deployment_config import ModelDeploymentConfig
2726
from nemo_platform.types.models.model_entity import ModelEntity
2827
from nmp.common.config import get_platform_config
2928
from nmp.core.models.app import (
@@ -40,11 +39,14 @@
4039
)
4140
from nmp.core.models.controllers.backends.engine import ENGINE_GENERIC, ENGINE_VLLM, config_engine
4241
from nmp.core.models.controllers.backends.k8s_nim_operator.config import K8sNimOperatorConfig
43-
from nmp.core.models.controllers.backends.k8s_nim_operator.reconcilers.base import ResolvedDeployment
42+
from nmp.core.models.controllers.backends.k8s_nim_operator.reconcilers.base import Reconciler, ResolvedDeployment
4443
from nmp.core.models.controllers.backends.k8s_nim_operator.reconcilers.k8s import K8sReconciler
4544
from nmp.core.models.controllers.backends.k8s_nim_operator.reconcilers.nim_operator import (
4645
NimOperatorReconciler,
4746
)
47+
from nmp.core.models.controllers.backends.k8s_nim_operator.reconcilers.resource_deleter import ResourceDeleter
48+
from nmp.core.models.controllers.backends.k8s_nim_operator.reconcilers.status_projector import StatusProjector
49+
from nmp.core.models.controllers.context import ModelContext
4850

4951
logger = getLogger(__name__)
5052

@@ -62,6 +64,8 @@ def __init__(self, nmp_sdk, config, huggingface_model_puller: str):
6264
self._k8s_namespace: str | None = None
6365
self._backend_config: K8sNimOperatorConfig | None = None
6466
self._huggingface_model_puller = huggingface_model_puller
67+
self._status_projector: StatusProjector | None = None
68+
self._resource_deleter: ResourceDeleter | None = None
6569
self._nim_reconciler: NimOperatorReconciler | None = None
6670
self._k8s_reconciler: K8sReconciler | None = None
6771
super().__init__(nmp_sdk, config)
@@ -88,18 +92,30 @@ def init(self) -> None:
8892
self._k8s_namespace = self._get_current_namespace()
8993
logger.info(f"Models controller will deploy models to namespace: {self._k8s_namespace}")
9094

91-
self._nim_reconciler = NimOperatorReconciler(
95+
# Shared collaborators composed into both reconcilers (and used directly
96+
# by the PENDING-timeout policy below).
97+
self._status_projector = StatusProjector(
9298
k8s_client_=self._k8s_client,
99+
backend_config=self._backend_config,
100+
k8s_namespace=self._k8s_namespace,
101+
)
102+
self._resource_deleter = ResourceDeleter(k8s_namespace=self._k8s_namespace)
103+
104+
self._nim_reconciler = NimOperatorReconciler(
93105
dynamic_client=self._dynamic_client,
94106
backend_config=self._backend_config,
95107
k8s_namespace=self._k8s_namespace,
96108
huggingface_model_puller=self._huggingface_model_puller,
109+
status=self._status_projector,
110+
deleter=self._resource_deleter,
97111
)
98112
self._k8s_reconciler = K8sReconciler(
99113
k8s_client_=self._k8s_client,
100114
backend_config=self._backend_config,
101115
k8s_namespace=self._k8s_namespace,
102116
huggingface_model_puller=self._huggingface_model_puller,
117+
status=self._status_projector,
118+
deleter=self._resource_deleter,
103119
)
104120

105121
def shutdown(self) -> None:
@@ -189,13 +205,11 @@ def _remote_files_hf_url(self) -> str:
189205
files_url = platform_config.service_discovery.get("files") or platform_config.base_url
190206
return urljoin(files_url.rstrip("/") + "/", "apis/files/v2/hf")
191207

192-
def _resolve(
193-
self,
194-
deployment: ModelDeployment,
195-
config: ModelDeploymentConfig,
196-
model_entity: Optional[ModelEntity],
197-
) -> ResolvedDeployment:
208+
def _resolve(self, ctx: ModelContext) -> ResolvedDeployment:
198209
"""Resolve everything a reconciler needs from the API object + SDK state."""
210+
deployment = ctx.model_deployment
211+
config = ctx.model_deployment_config
212+
model_entity = ctx.model_entity
199213
view = deployment_config_view(config)
200214
model_namespace, model_name, model_revision = self._resolve_model_source(model_entity, view)
201215
weights_type = get_model_weights_type(
@@ -218,8 +232,14 @@ def _resolve(
218232
huggingface_model_puller=self._huggingface_model_puller,
219233
)
220234

221-
def _select_reconciler(self, engine: str):
222-
"""Select the reconciler for an engine, rejecting the unsupported one."""
235+
def _select_reconciler(self, engine: str) -> Optional[Reconciler]:
236+
"""Select the reconciler for an engine.
237+
238+
Returns the vLLM reconciler for ``vllm``, the NIM-operator reconciler for
239+
any other engine (the default), and ``None`` for ``generic`` -- which the
240+
callers treat as the "unsupported engine" rejection (see
241+
:meth:`_unsupported_engine`).
242+
"""
223243
if engine == ENGINE_VLLM:
224244
return self._k8s_reconciler
225245
if engine == ENGINE_GENERIC:
@@ -239,48 +259,42 @@ def _unsupported_engine(engine: str) -> DeploymentStatusUpdate:
239259
# ServiceBackend interface (resolve + select + delegate)
240260
# ------------------------------------------------------------------
241261

242-
async def create_model_deployment(
243-
self, deployment: ModelDeployment, config: ModelDeploymentConfig, model_entity: Optional[ModelEntity] = None
244-
) -> DeploymentStatusUpdate:
245-
"""Create a new model deployment (dispatches on ``config.engine``)."""
246-
engine = config_engine(config)
262+
async def create_model_deployment(self, ctx: ModelContext) -> DeploymentStatusUpdate:
263+
"""Create a new model deployment (dispatches on the config's engine)."""
264+
engine = config_engine(ctx.model_deployment_config)
247265
reconciler = self._select_reconciler(engine)
248266
if reconciler is None:
249267
return self._unsupported_engine(engine)
250-
resolved = self._resolve(deployment, config, model_entity)
268+
resolved = self._resolve(ctx)
251269
return await reconciler.create(resolved)
252270

253-
async def update_model_deployment(
254-
self, deployment: ModelDeployment, config: ModelDeploymentConfig, model_entity: Optional[ModelEntity] = None
255-
) -> DeploymentStatusUpdate:
256-
"""Update an existing model deployment (dispatches on ``config.engine``)."""
257-
engine = config_engine(config)
271+
async def update_model_deployment(self, ctx: ModelContext) -> DeploymentStatusUpdate:
272+
"""Update an existing model deployment (dispatches on the config's engine)."""
273+
engine = config_engine(ctx.model_deployment_config)
258274
reconciler = self._select_reconciler(engine)
259275
if reconciler is None:
260276
return self._unsupported_engine(engine)
261-
resolved = self._resolve(deployment, config, model_entity)
277+
resolved = self._resolve(ctx)
262278
return await reconciler.update(resolved)
263279

264-
async def get_model_deployment_status(
265-
self,
266-
deployment: ModelDeployment,
267-
config: Optional[ModelDeploymentConfig] = None,
268-
model_entity: Optional[ModelEntity] = None,
269-
) -> DeploymentStatusUpdate:
280+
async def get_model_deployment_status(self, ctx: ModelContext) -> DeploymentStatusUpdate:
270281
"""Get the current status of a model deployment.
271282
272-
The engine is taken from ``config`` (same selection as create/update), so a
273-
config is required. When ``config`` is ``None`` (e.g. the controller failed
274-
to fetch it this cycle) the backend cannot determine the deployment's state
275-
and returns ``UNKNOWN``; the controller retries on the next poll (which
276-
normally has a config) and escalates to ERROR after its retry budget.
283+
The engine is taken from the config (same selection as create/update), so a
284+
config is required. When ``ctx.model_deployment_config`` is ``None`` (e.g.
285+
the controller failed to fetch it this cycle) the backend cannot determine
286+
the deployment's state and returns ``UNKNOWN``; the controller retries on
287+
the next poll (which normally has a config) and escalates to ERROR after
288+
its retry budget.
277289
278290
In addition to the reconciler's status, this method enforces the PENDING
279291
timeout policy: if the deployment has been alive longer than
280292
``pending_timeout_seconds`` and is still PENDING, transition to ERROR with
281293
diagnostic information. (Crash-loop detection is handled inside the
282294
reconciler's pod drill-down.)
283295
"""
296+
deployment = ctx.model_deployment
297+
config = ctx.model_deployment_config
284298
logger.debug(
285299
f"Checking deployment status: {deployment.workspace}/{deployment.name} "
286300
f"(version: {deployment.entity_version})"
@@ -305,15 +319,15 @@ async def get_model_deployment_status(
305319
return self._unsupported_engine(engine)
306320
# A reconciler MAY advance creation in get_status; it needs the
307321
# resolved config to compile the serving spec.
308-
resolved = self._resolve(deployment, config, model_entity)
322+
resolved = self._resolve(ctx)
309323
result = await reconciler.get_status(resolved)
310324

311325
if result.status == "PENDING":
312326
elapsed = deployment_elapsed_seconds(deployment)
313327

314328
if elapsed >= self._backend_config.pending_timeout_seconds:
315-
pod_name = self._nim_reconciler._find_pod_name(resource_name)
316-
return self._nim_reconciler._build_pending_timeout_error(resource_name, elapsed, pod_name)
329+
pod_name = self._status_projector.find_pod_name(resource_name)
330+
return self._status_projector.build_pending_timeout_error(resource_name, elapsed, pod_name)
317331

318332
# Use a stable message (no elapsed/timeout) so we don't create a new history entry every poll
319333

0 commit comments

Comments
 (0)