diff --git a/config/litellm.yaml b/config/litellm.yaml new file mode 100644 index 0000000..20a04e5 --- /dev/null +++ b/config/litellm.yaml @@ -0,0 +1,25 @@ +# LiteLLM proxy configuration: the one place an approved external provider is +# named. The gateway reaches providers only through this proxy, so provider +# credentials never enter the gateway and adding a provider never changes it. +# +# Every alias under model_list must match a non-local model card in +# config/registry.yaml; a test enforces it. Keys are read from the environment, +# which the cluster secret manager populates. Nothing secret is committed here. +model_list: + - model_name: approved-external-fallback + litellm_params: + model: anthropic/claude-sonnet-5-5 + api_key: os.environ/EXTERNAL_PROVIDER_API_KEY + +general_settings: + master_key: os.environ/LITELLM_MASTER_KEY + +litellm_settings: + # Only public requests may reach this proxy, but prompts are still kept out + # of its logs: the gateway's redacted traces are the record of a request. + turn_off_message_logging: true + # The gateway owns retries and fallback, so a failure is reported once + # rather than retried here and again upstream. + num_retries: 0 + request_timeout: 60 + drop_params: true diff --git a/deploy/kubernetes/external-provider.yaml b/deploy/kubernetes/external-provider.yaml new file mode 100644 index 0000000..f83ce57 --- /dev/null +++ b/deploy/kubernetes/external-provider.yaml @@ -0,0 +1,156 @@ +# LiteLLM proxy for the approved external fallback. It is the only workload in +# the namespace allowed to reach the internet, and only the gateway may call it. +apiVersion: v1 +kind: ConfigMap +metadata: + name: litellm-config + namespace: llm-routing +data: + # Kept identical to config/litellm.yaml; a test fails if the two drift. + config.yaml: | + model_list: + - model_name: approved-external-fallback + litellm_params: + model: anthropic/claude-sonnet-5-5 + api_key: os.environ/EXTERNAL_PROVIDER_API_KEY + general_settings: + master_key: os.environ/LITELLM_MASTER_KEY + litellm_settings: + turn_off_message_logging: true + num_retries: 0 + request_timeout: 60 + drop_params: true +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + name: litellm-proxy + namespace: llm-routing + labels: + app.kubernetes.io/name: litellm-proxy + app.kubernetes.io/part-of: local-llm-router +spec: + replicas: 2 + selector: + matchLabels: + app.kubernetes.io/name: litellm-proxy + template: + metadata: + labels: + app.kubernetes.io/name: litellm-proxy + app.kubernetes.io/part-of: local-llm-router + spec: + securityContext: + runAsNonRoot: true + runAsUser: 10001 + seccompProfile: + type: RuntimeDefault + containers: + - name: litellm + image: ghcr.io/berriai/litellm@sha256:REPLACE_ME + imagePullPolicy: IfNotPresent + args: ["--config", "/etc/litellm/config.yaml", "--port", "4000"] + ports: + - name: http + containerPort: 4000 + env: + - name: LITELLM_MASTER_KEY + valueFrom: + secretKeyRef: + name: llm-gateway-credentials + key: external-proxy-key + - name: EXTERNAL_PROVIDER_API_KEY + valueFrom: + secretKeyRef: + name: llm-gateway-credentials + key: external-provider-key + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: true + capabilities: + drop: [ALL] + resources: + requests: + cpu: 250m + memory: 512Mi + limits: + cpu: "1" + memory: 1Gi + livenessProbe: + httpGet: + path: /health/liveliness + port: http + initialDelaySeconds: 15 + periodSeconds: 30 + readinessProbe: + httpGet: + path: /health/readiness + port: http + initialDelaySeconds: 15 + periodSeconds: 10 + volumeMounts: + - name: config + mountPath: /etc/litellm + readOnly: true + - name: tmp + mountPath: /tmp + volumes: + - name: config + configMap: + name: litellm-config + - name: tmp + emptyDir: {} + terminationGracePeriodSeconds: 60 +--- +apiVersion: v1 +kind: Service +metadata: + name: litellm-proxy + namespace: llm-routing + labels: + app.kubernetes.io/name: litellm-proxy +spec: + selector: + app.kubernetes.io/name: litellm-proxy + ports: + - name: http + port: 4000 + targetPort: http +--- +apiVersion: networking.k8s.io/v1 +kind: NetworkPolicy +metadata: + name: litellm-proxy + namespace: llm-routing +spec: + podSelector: + matchLabels: + app.kubernetes.io/name: litellm-proxy + policyTypes: [Ingress, Egress] + ingress: + - from: + - podSelector: + matchLabels: + app.kubernetes.io/name: llm-gateway + ports: + - protocol: TCP + port: 4000 + egress: + # Name resolution for the provider endpoint. + - to: + - namespaceSelector: + matchLabels: + kubernetes.io/metadata.name: kube-system + ports: + - protocol: UDP + port: 53 + - protocol: TCP + port: 53 + # HTTPS to the provider, and nothing inside the cluster or private ranges. + - to: + - ipBlock: + cidr: 0.0.0.0/0 + except: [10.0.0.0/8, 172.16.0.0/12, 192.168.0.0/16] + ports: + - protocol: TCP + port: 443 diff --git a/deploy/kubernetes/gateway.yaml b/deploy/kubernetes/gateway.yaml index d5f389e..a4e3106 100644 --- a/deploy/kubernetes/gateway.yaml +++ b/deploy/kubernetes/gateway.yaml @@ -42,6 +42,15 @@ spec: secretKeyRef: name: llm-gateway-credentials key: api-keys + # External fallback stays off until an operator enables it; the + # proxy address and key are in place so enabling it is one setting. + - name: ROUTER_EXTERNAL_BASE_URL + value: http://litellm-proxy.llm-routing.svc.cluster.local:4000 + - name: ROUTER_EXTERNAL_API_KEY + valueFrom: + secretKeyRef: + name: llm-gateway-credentials + key: external-proxy-key securityContext: allowPrivilegeEscalation: false readOnlyRootFilesystem: true diff --git a/deploy/kubernetes/kustomization.yaml b/deploy/kubernetes/kustomization.yaml index 5bac207..8e17a43 100644 --- a/deploy/kubernetes/kustomization.yaml +++ b/deploy/kubernetes/kustomization.yaml @@ -8,4 +8,5 @@ resources: - vllm-serve.yaml - state.yaml - network-policy.yaml + - external-provider.yaml - observability.yaml diff --git a/deploy/kubernetes/network-policy.yaml b/deploy/kubernetes/network-policy.yaml index e2b06b8..d8889b3 100644 --- a/deploy/kubernetes/network-policy.yaml +++ b/deploy/kubernetes/network-policy.yaml @@ -38,6 +38,14 @@ spec: ports: - protocol: TCP port: 6379 + # The gateway never reaches a provider itself, only the proxy that does. + - to: + - podSelector: + matchLabels: + app.kubernetes.io/name: litellm-proxy + ports: + - protocol: TCP + port: 4000 --- apiVersion: networking.k8s.io/v1 kind: NetworkPolicy diff --git a/deploy/kubernetes/state.yaml b/deploy/kubernetes/state.yaml index 3e2a260..f50d01d 100644 --- a/deploy/kubernetes/state.yaml +++ b/deploy/kubernetes/state.yaml @@ -86,3 +86,7 @@ spec: remoteRef: key: llm-routing/gateway property: external_provider_key + - secretKey: external-proxy-key + remoteRef: + key: llm-routing/gateway + property: external_proxy_key diff --git a/src/llm_router/app.py b/src/llm_router/app.py index 95aa55f..7f3093c 100644 --- a/src/llm_router/app.py +++ b/src/llm_router/app.py @@ -24,6 +24,8 @@ BackendOutOfMemoryError, BackendResult, BackendUnavailableError, + DispatchingBackend, + ExternalDispatchRefusedError, InferenceBackend, MockInferenceBackend, VLLMBackend, @@ -51,6 +53,7 @@ ChatCompletionRequest, ChatCompletionResponse, ChatMessage, + PrivacyClass, RouteDecision, Usage, ) @@ -163,7 +166,32 @@ def create_app( failure_threshold=runtime_settings.engine_failure_threshold, cooldown_seconds=runtime_settings.engine_cooldown_seconds, ) - inference_backend: InferenceBackend = ResilientBackend(raw_backend, circuit) + external_client = ( + httpx.AsyncClient() if backend is None and runtime_settings.external_base_url else None + ) + # An injected or mock backend answers external routes too, which keeps + # tests and local development free of a provider. + raw_external: InferenceBackend = ( + VLLMBackend( + base_url=runtime_settings.external_base_url, + client=external_client, + request_timeout_seconds=runtime_settings.backend_timeout_seconds, + api_key=runtime_settings.external_api_key, + health_path="/health/liveliness", + ) + if external_client is not None + else raw_backend + ) + # Each target has its own circuit: a lost GPU node must not close the + # route to the provider, nor a provider outage the route to the engine. + external_circuit = CircuitBreaker( + failure_threshold=runtime_settings.engine_failure_threshold, + cooldown_seconds=runtime_settings.engine_cooldown_seconds, + ) + inference_backend: InferenceBackend = DispatchingBackend( + local=ResilientBackend(raw_backend, circuit), + external=ResilientBackend(raw_external, external_circuit), + ) telemetry = metrics if metrics is not None else Metrics() engine_telemetry = engine_stats or ( EngineStatsCollector(base_url=runtime_settings.vllm_base_url, client=engine_client) @@ -211,6 +239,8 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: await admission.drain(runtime_settings.shutdown_grace_seconds) if engine_client is not None: await engine_client.aclose() + if external_client is not None: + await external_client.aclose() app = FastAPI( title="Local LLM Inference Router", @@ -261,6 +291,16 @@ async def admission_handler(_: Request, error: AdmissionRejectedError) -> JSONRe content={"error": {"message": str(error), "type": "overloaded"}}, ) + @app.exception_handler(ExternalDispatchRefusedError) + async def refused_handler(_: Request, error: ExternalDispatchRefusedError) -> JSONResponse: + # Reaching here means routing chose an external model for data that + # must stay local. The request is refused, and counted so it is seen. + telemetry.record_rejection("external_dispatch_refused") + return JSONResponse( + status_code=500, + content={"error": {"message": str(error), "type": "policy_violation"}}, + ) + @app.exception_handler(BackendOutOfMemoryError) async def out_of_memory_handler(_: Request, error: BackendOutOfMemoryError) -> JSONResponse: # Not a retry-as-is condition: the same request at the same size will @@ -323,6 +363,7 @@ async def prometheus_metrics() -> Response: # gateway stays the single scrape target for the whole serving path and # no background poller runs when nobody is collecting. telemetry.record_circuit_state(circuit.state, engine=engine_label) + telemetry.record_circuit_state(external_circuit.state, engine="external") if engine_telemetry is not None: stats = await engine_telemetry.sample() if stats is not None: @@ -555,6 +596,38 @@ def _record_canary( if verdict is not None and verdict.action == "rollback": telemetry.record_canary_rollback(decision.canary_subject) + def _fallback_route( + payload: ChatCompletionRequest, + failed: RouteDecision, + tenant: TenantRecord | None, + privacy_raised_from: PrivacyClass | None, + ) -> RouteDecision | None: + """Find an approved external route for a request the local engine failed. + + Only a local failure falls back, and only to a model the request was + already entitled to: the same privacy, tenant, operator, and opt-in + rules apply as on the first attempt. + """ + + if not failed.profile.local: + return None + try: + return router.select( + payload, + task=failed.task, + permitted_models=( + catalog.permitted_models_for(tenant) if catalog is not None else None + ), + quality_floor=tenant.quality_floor if tenant is not None else 0.0, + tenant_allows_external=( + tenant.allow_external_fallback is not False if tenant is not None else True + ), + privacy_raised_from=privacy_raised_from, + fallback_from=failed.profile.id, + ) + except NoEligibleModelError: + return None + async def _consume_quota(subject: str, tenant: TenantRecord | None) -> None: # A tenant may carry its own limit; absent one the platform default # applies, which is why None means "defer" rather than "unlimited". @@ -827,10 +900,23 @@ async def _complete( ) try: - result = await inference_backend.generate(payload, decision) - except BackendUnavailableError: - _record_canary(decision, ok=False, started=started) - raise + try: + result = await inference_backend.generate(payload, decision) + except BackendUnavailableError as error: + _record_canary(decision, ok=False, started=started) + fallback = _fallback_route(payload, decision, tenant, privacy_raised_from) + if fallback is None: + raise + # Declared, attributed, and counted: never a silent switch. + telemetry.record_fallback( + from_model=decision.profile.id, + to_model=fallback.profile.id, + cause=type(error).__name__, + ) + telemetry.record_route(fallback, privacy=payload.routing.privacy.value) + decision = fallback + span.set_route(decision) + result = await inference_backend.generate(payload, decision) finally: telemetry.inflight_requests.dec() admission.release() diff --git a/src/llm_router/backends.py b/src/llm_router/backends.py index 48f7641..6f05f49 100644 --- a/src/llm_router/backends.py +++ b/src/llm_router/backends.py @@ -12,7 +12,7 @@ import httpx -from llm_router.models import ChatCompletionRequest, RouteDecision +from llm_router.models import ChatCompletionRequest, PrivacyClass, RouteDecision class BackendUnavailableError(RuntimeError): @@ -100,6 +100,14 @@ class VLLMBackend: base_url: str client: httpx.AsyncClient request_timeout_seconds: float = 60.0 + # The same OpenAI-compatible client reaches the LiteLLM proxy that fronts + # external providers, which authenticates and has its own health path. + api_key: str = "" + health_path: str = "/health" + + @property + def _headers(self) -> dict[str, str]: + return {"Authorization": f"Bearer {self.api_key}"} if self.api_key else {} def _payload( self, request: ChatCompletionRequest, decision: RouteDecision, *, stream: bool @@ -119,6 +127,7 @@ async def generate( response = await self.client.post( f"{self.base_url}/v1/chat/completions", json=self._payload(request, decision, stream=False), + headers=self._headers, timeout=self.request_timeout_seconds, ) except httpx.HTTPError as error: @@ -151,6 +160,7 @@ async def stream( "POST", f"{self.base_url}/v1/chat/completions", json=payload, + headers=self._headers, timeout=self.request_timeout_seconds, ) as response: if response.status_code >= 400: @@ -165,12 +175,59 @@ async def stream( async def healthy(self) -> bool: try: - response = await self.client.get(f"{self.base_url}/health", timeout=2.0) + response = await self.client.get( + f"{self.base_url}{self.health_path}", headers=self._headers, timeout=2.0 + ) except httpx.HTTPError: return False return response.status_code < 400 +class ExternalDispatchRefusedError(RuntimeError): + """Raised when a non-public request reaches the external dispatch boundary.""" + + +@dataclass +class DispatchingBackend: + """Sends each decision to the local engine or the external provider proxy. + + Routing already refuses to pick an external model for private or + restricted data. The same rule is enforced again here, at the last point + before a prompt could leave the private environment, so a routing defect + cannot become a disclosure. + """ + + local: InferenceBackend + external: InferenceBackend | None = None + + def _target(self, request: ChatCompletionRequest, decision: RouteDecision) -> InferenceBackend: + if decision.profile.local: + return self.local + if request.routing.privacy is not PrivacyClass.PUBLIC: + raise ExternalDispatchRefusedError( + f"refused to send a {request.routing.privacy.value} request to an external model" + ) + if self.external is None: + raise BackendUnavailableError("no external provider is configured") + return self.external + + async def generate( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> BackendResult: + return await self._target(request, decision).generate(request, decision) + + async def stream( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> AsyncIterator[str]: + async for delta in self._target(request, decision).stream(request, decision): + yield delta + + async def healthy(self) -> bool: + """Readiness follows the local engine; the external provider is optional.""" + + return await self.local.healthy() + + def _parse_stream_line(line: str) -> str: """Extract the content delta from one server-sent-event line.""" diff --git a/src/llm_router/config.py b/src/llm_router/config.py index 6dd34e3..de83171 100644 --- a/src/llm_router/config.py +++ b/src/llm_router/config.py @@ -21,6 +21,10 @@ class Settings(BaseSettings): admission_timeout_seconds: float = Field(default=0.25, gt=0) quota_requests_per_minute: int = Field(default=120, ge=1) external_fallback_enabled: bool = False + # The LiteLLM proxy that normalizes approved external providers. Provider + # credentials live in the proxy; this key only authenticates to it. + external_base_url: str = "" + external_api_key: str = "" backend: Literal["mock", "vllm"] = "mock" vllm_base_url: str = "http://127.0.0.1:8001" backend_timeout_seconds: float = Field(default=60.0, gt=0) @@ -54,6 +58,15 @@ def reject_development_key_in_shared_environments(self) -> "Settings": raise ValueError("ROUTER_API_KEYS must be set outside development and test") return self + @model_validator(mode="after") + def require_a_provider_before_enabling_external_fallback(self) -> "Settings": + if self.backend == "vllm" and self.external_fallback_enabled and not self.external_base_url: + raise ValueError( + "ROUTER_EXTERNAL_BASE_URL must be set when external fallback is enabled; " + "without it an external route has nowhere to go" + ) + return self + @property def accepted_api_keys(self) -> frozenset[str]: return frozenset(self.tenant_by_key) diff --git a/src/llm_router/observability.py b/src/llm_router/observability.py index 3c55f67..60d73f1 100644 --- a/src/llm_router/observability.py +++ b/src/llm_router/observability.py @@ -129,6 +129,12 @@ def __init__(self, registry: CollectorRegistry | None = None) -> None: buckets=LATENCY_BUCKETS, registry=self.registry, ) + self.fallbacks_total = Counter( + "router_fallbacks_total", + "Requests re-routed after the model first chosen failed.", + ["from_model", "to_model", "cause"], + registry=self.registry, + ) self.canary_requests_total = Counter( "router_canary_requests_total", "Requests on a canaried route by subject, arm, and outcome.", @@ -255,6 +261,9 @@ def record_structured_output(self, decision: RouteDecision, *, valid: bool) -> N abs(decision.profile.quality - observed) ) + def record_fallback(self, *, from_model: str, to_model: str, cause: str) -> None: + self.fallbacks_total.labels(from_model=from_model, to_model=to_model, cause=cause).inc() + def record_canary(self, subject: str, *, arm: str, ok: bool) -> None: self.canary_requests_total.labels( subject=subject, arm=arm, outcome="success" if ok else "error" diff --git a/src/llm_router/routing.py b/src/llm_router/routing.py index f3ab257..7fe252f 100644 --- a/src/llm_router/routing.py +++ b/src/llm_router/routing.py @@ -132,6 +132,7 @@ def select( tenant_allows_external: bool = True, privacy_raised_from: PrivacyClass | None = None, canary_key: str = "", + fallback_from: str | None = None, ) -> RouteDecision: prediction = ( self.classifier.predict(request.prompt) if self.classifier is not None else None @@ -171,6 +172,9 @@ def select( # A tenant entitlement is a hard restriction, like privacy: it is # applied before scoring and can never be outscored. and (permitted_models is None or profile.id in permitted_models) + # A fallback leaves the engine that just failed: every local + # model shares it, so only a model served elsewhere can help. + and (fallback_from is None or not profile.local) ] if request.model != "auto": @@ -247,6 +251,8 @@ def score(profile: ModelProfile) -> float: f"; applied adapter {adapter.id} for domain {adapter.domain} " f"(measured quality delta {adapter.benchmark.quality_delta:+.3f})" ) + if fallback_from is not None: + reason += f"; fell back from {fallback_from} after the local engine failed" if canary_arm == "canary": reason += f"; canary arm of {canary_subject}" elif canary_arm == "stable": diff --git a/tests/integration/test_external_provider_api.py b/tests/integration/test_external_provider_api.py new file mode 100644 index 0000000..7154625 --- /dev/null +++ b/tests/integration/test_external_provider_api.py @@ -0,0 +1,295 @@ +import json +from collections.abc import AsyncIterator +from typing import Any + +import httpx +import pytest +from fastapi.testclient import TestClient +from pydantic import ValidationError + +from llm_router.app import create_app +from llm_router.backends import ( + BackendOutOfMemoryError, + BackendResult, + BackendUnavailableError, + DispatchingBackend, + ExternalDispatchRefusedError, + MockInferenceBackend, + VLLMBackend, +) +from llm_router.config import Settings +from llm_router.models import ChatCompletionRequest, ModelProfile, RouteDecision, TaskClass +from llm_router.routing import NoEligibleModelError, Router, default_model_profiles + +HEADERS = {"Authorization": "Bearer external-key"} + + +def request(privacy: str = "public", **routing: object) -> ChatCompletionRequest: + return ChatCompletionRequest.model_validate( + { + "messages": [{"role": "user", "content": "Analyze deeply"}], + "routing": {"privacy": privacy, **routing}, + } + ) + + +def decision(*, local: bool) -> RouteDecision: + return RouteDecision( + profile=ModelProfile( + id="general-local" if local else "approved-external-fallback", + revision="rev", + local=local, + context_limit=8192, + supported_tasks=frozenset(TaskClass), + quality=0.9, + ), + task=TaskClass.REASONING, + reason="test", + score=1.0, + candidate_count=1, + ) + + +class Named(MockInferenceBackend): + def __init__(self, name: str) -> None: + self.name = name + self.calls = 0 + + async def generate( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> BackendResult: + self.calls += 1 + return BackendResult(text=self.name, prompt_tokens=1, completion_tokens=1) + + async def stream( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> AsyncIterator[str]: + self.calls += 1 + yield self.name + + +@pytest.mark.asyncio +async def test_decisions_are_dispatched_to_the_engine_that_serves_them() -> None: + local, external = Named("local"), Named("external") + backend = DispatchingBackend(local=local, external=external) + + assert (await backend.generate(request(), decision(local=True))).text == "local" + assert (await backend.generate(request(), decision(local=False))).text == "external" + assert [delta async for delta in backend.stream(request(), decision(local=False))] == [ + "external" + ] + + +@pytest.mark.asyncio +@pytest.mark.parametrize("privacy", ["private", "restricted"]) +async def test_the_dispatch_boundary_refuses_non_public_data_whatever_routing_decided( + privacy: str, +) -> None: + external = Named("external") + backend = DispatchingBackend(local=Named("local"), external=external) + + with pytest.raises(ExternalDispatchRefusedError, match=privacy): + await backend.generate(request(privacy), decision(local=False)) + with pytest.raises(ExternalDispatchRefusedError): + async for _ in backend.stream(request(privacy), decision(local=False)): + pass + + # The provider was never contacted. + assert external.calls == 0 + + +@pytest.mark.asyncio +async def test_an_external_route_with_no_provider_is_unavailable_not_sent_locally() -> None: + local = Named("local") + backend = DispatchingBackend(local=local, external=None) + + with pytest.raises(BackendUnavailableError, match="no external provider"): + await backend.generate(request(), decision(local=False)) + assert local.calls == 0 + assert await backend.healthy() is True + + +@pytest.mark.asyncio +async def test_the_proxy_client_authenticates_and_uses_the_proxy_health_path() -> None: + seen: list[httpx.Request] = [] + + def handler(received: httpx.Request) -> httpx.Response: + seen.append(received) + if received.url.path == "/health/liveliness": + return httpx.Response(200) + return httpx.Response( + 200, + json={ + "choices": [{"message": {"content": "answer"}}], + "usage": {"prompt_tokens": 3, "completion_tokens": 1}, + }, + ) + + async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client: + proxy = VLLMBackend( + base_url="http://litellm:4000", + client=client, + api_key="proxy-key", + health_path="/health/liveliness", + ) + result = await proxy.generate(request(), decision(local=False)) + healthy = await proxy.healthy() + + completion = next(item for item in seen if item.url.path == "/v1/chat/completions") + assert result.text == "answer" and healthy is True + assert completion.headers["Authorization"] == "Bearer proxy-key" + assert json.loads(completion.content)["model"] == "approved-external-fallback" + + +def test_enabling_external_fallback_without_a_provider_is_refused_at_start_up() -> None: + with pytest.raises(ValidationError, match="ROUTER_EXTERNAL_BASE_URL must be set"): + Settings(api_keys="k", backend="vllm", external_fallback_enabled=True) + + # Fine once the proxy is named, and fine in mock mode where nothing leaves. + Settings( + api_keys="k", + backend="vllm", + external_fallback_enabled=True, + external_base_url="http://litellm:4000", + ) + Settings(api_keys="k", external_fallback_enabled=True) + + +def test_a_fallback_route_only_considers_models_served_elsewhere() -> None: + router = Router(default_model_profiles(), external_fallback_enabled=True) + eligible = request(allow_external_fallback=True) + + fallback = router.select(eligible, fallback_from="high-capability") + + assert fallback.profile.id == "approved-external-fallback" + assert "fell back from high-capability after the local engine failed" in fallback.reason + with pytest.raises(NoEligibleModelError): + router.select(request("private", allow_external_fallback=True), fallback_from="x") + with pytest.raises(NoEligibleModelError): + router.select(request(), fallback_from="x") + + +class LocalDown(MockInferenceBackend): + """The local engine fails; anything routed to the external model succeeds.""" + + def __init__(self, error: BackendUnavailableError) -> None: + self.error = error + self.served: list[str] = [] + + async def generate( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> BackendResult: + self.served.append(decision.profile.id) + if decision.profile.local: + raise self.error + return BackendResult(text="from the provider", prompt_tokens=2, completion_tokens=3) + + +def post(client: TestClient, **routing: object) -> Any: + return client.post( + "/v1/chat/completions", + headers=HEADERS, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Analyze deeply"}], + "routing": routing, + }, + ) + + +def build(backend: MockInferenceBackend, *, external: bool = True) -> TestClient: + settings = Settings(api_keys="external-key", external_fallback_enabled=external) + return TestClient(create_app(settings, backend=backend)) + + +@pytest.mark.parametrize( + "error", + [ + BackendUnavailableError("inference engine unreachable"), + BackendOutOfMemoryError("inference engine ran out of GPU memory"), + ], +) +def test_an_eligible_request_falls_back_when_the_local_engine_fails( + error: BackendUnavailableError, +) -> None: + backend = LocalDown(error) + with build(backend) as client: + response = post(client, privacy="public", allow_external_fallback=True) + metrics = client.get("/metrics").text + + body = response.json() + assert response.status_code == 200 + assert backend.served == ["high-capability", "approved-external-fallback"] + assert body["model"] == "approved-external-fallback" + assert body["choices"][0]["message"]["content"] == "from the provider" + # Attributed in the response and counted, never silent. + assert "fell back from high-capability" in body["routing"]["reason"] + assert response.headers["X-Route-Model"] == "approved-external-fallback" + assert ( + 'router_fallbacks_total{cause="' + type(error).__name__ + '",' + 'from_model="high-capability",to_model="approved-external-fallback"} 1.0' + ) in metrics + assert 'router_external_fallback_total{model="approved-external-fallback"} 1.0' in metrics + + +@pytest.mark.parametrize( + "routing", + [ + {"privacy": "private", "allow_external_fallback": True}, + {"privacy": "restricted", "allow_external_fallback": True}, + {"privacy": "public"}, + ], +) +def test_a_request_that_was_never_entitled_to_external_does_not_fall_back( + routing: dict[str, object], +) -> None: + backend = LocalDown(BackendUnavailableError("inference engine unreachable")) + with build(backend) as client: + response = post(client, **routing) + + assert response.status_code == 502 + assert "approved-external-fallback" not in backend.served + + +def test_the_operator_gate_also_governs_fallback() -> None: + backend = LocalDown(BackendUnavailableError("inference engine unreachable")) + with build(backend, external=False) as client: + response = post(client, privacy="public", allow_external_fallback=True) + + assert response.status_code == 502 + assert backend.served == ["high-capability"] + + +def test_a_fallback_response_is_cached_under_the_model_that_produced_it() -> None: + backend = LocalDown(BackendUnavailableError("inference engine unreachable")) + with build(backend) as client: + post(client, privacy="public", allow_external_fallback=True) + repeat = post(client, privacy="public", allow_external_fallback=True) + + assert repeat.headers["X-Cache"] == "exact" + assert repeat.headers["X-Route-Model"] == "approved-external-fallback" + + +def test_each_target_has_its_own_circuit() -> None: + backend = LocalDown(BackendUnavailableError("inference engine unreachable")) + settings = Settings( + api_keys="external-key", external_fallback_enabled=True, engine_failure_threshold=1 + ) + with TestClient(create_app(settings, backend=backend)) as client: + first = post(client, privacy="public", allow_external_fallback=True) + # The local circuit is now open; the provider's is not. + second = client.post( + "/v1/chat/completions", + headers=HEADERS, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Analyze deeply, again"}], + "routing": {"privacy": "public", "allow_external_fallback": True}, + }, + ) + metrics = client.get("/metrics").text + + assert first.status_code == 200 and second.status_code == 200 + assert second.json()["model"] == "approved-external-fallback" + assert 'router_engine_circuit_open{engine="mock"} 1.0' in metrics + assert 'router_engine_circuit_open{engine="external"} 0.0' in metrics diff --git a/tests/unit/test_deployment_manifests.py b/tests/unit/test_deployment_manifests.py index cad80a9..b379884 100644 --- a/tests/unit/test_deployment_manifests.py +++ b/tests/unit/test_deployment_manifests.py @@ -16,6 +16,21 @@ SECRET_MARKERS = ("password:", "token:", "apiKey:", "api_key:", "BEGIN PRIVATE KEY") +def secret_material(content: str) -> list[str]: + """Lines that hold a secret, as opposed to pointing at where one is kept. + + A secret-shaped key may reference the environment, which the secret + manager populates. It may never carry a value of its own. + """ + + found: list[str] = [] + for line in content.splitlines(): + for marker in SECRET_MARKERS: + if marker in line and not line.split(marker, 1)[1].strip().startswith("os.environ/"): + found.append(line.strip()) + return found + + def load_documents() -> list[dict[str, Any]]: documents: list[dict[str, Any]] = [] for path in MANIFESTS: @@ -85,8 +100,7 @@ def test_no_manifest_contains_secret_material() -> None: for path in MANIFESTS: content = path.read_text(encoding="utf-8") assert "kind: Secret\n" not in content - for marker in SECRET_MARKERS: - assert marker not in content, f"{path.name} contains {marker}" + assert secret_material(content) == [], path.name def test_gateway_reads_credentials_from_the_secret_manager() -> None: @@ -165,3 +179,83 @@ def test_engine_ingress_is_restricted_to_the_gateway() -> None: sources = policy["spec"]["ingress"][0]["from"] assert sources == [{"podSelector": {"matchLabels": {"app.kubernetes.io/name": "llm-gateway"}}}] + + +def test_a_secret_shaped_key_may_reference_the_environment_but_not_hold_a_value() -> None: + reference = " api_key: os.environ/EXTERNAL_PROVIDER_API_KEY" + literal = " api_key: sk-live-not-a-reference" + + assert secret_material(reference) == [] + assert secret_material("\n".join((reference, literal))) == [literal.strip()] + assert secret_material("-----BEGIN PRIVATE KEY-----") == ["-----BEGIN PRIVATE KEY-----"] + + +def test_proxy_configuration_matches_the_committed_litellm_config() -> None: + config_map = next(document for document in DOCUMENTS if document["kind"] == "ConfigMap") + + embedded = yaml.safe_load(config_map["data"]["config.yaml"]) + committed = yaml.safe_load(Path("config/litellm.yaml").read_text(encoding="utf-8")) + assert embedded == committed + + +def test_every_external_model_in_the_catalog_has_a_proxy_alias() -> None: + catalog = yaml.safe_load(Path("config/registry.yaml").read_text(encoding="utf-8")) + proxy = yaml.safe_load(Path("config/litellm.yaml").read_text(encoding="utf-8")) + + external = {model["id"] for model in catalog["models"] if model.get("local") is False} + aliases = {entry["model_name"] for entry in proxy["model_list"]} + assert external and external == aliases + + +def test_the_proxy_keeps_prompts_out_of_its_logs_and_leaves_retries_to_the_gateway() -> None: + proxy = yaml.safe_load(Path("config/litellm.yaml").read_text(encoding="utf-8")) + + assert proxy["litellm_settings"]["turn_off_message_logging"] is True + assert proxy["litellm_settings"]["num_retries"] == 0 + for entry in proxy["model_list"]: + assert entry["litellm_params"]["api_key"].startswith("os.environ/") + + +def test_only_the_gateway_reaches_the_proxy_and_only_the_proxy_reaches_out() -> None: + policies = { + document["metadata"]["name"]: document["spec"] + for document in DOCUMENTS + if document["kind"] == "NetworkPolicy" + } + proxy = policies["litellm-proxy"] + + assert proxy["ingress"][0]["from"] == [ + {"podSelector": {"matchLabels": {"app.kubernetes.io/name": "llm-gateway"}}} + ] + blocks = [ + rule["ipBlock"] for entry in proxy["egress"] for rule in entry["to"] if "ipBlock" in rule + ] + assert blocks == [ + {"cidr": "0.0.0.0/0", "except": ["10.0.0.0/8", "172.16.0.0/12", "192.168.0.0/16"]} + ] + # No other workload has an egress rule that leaves the namespace by address. + for name, spec in policies.items(): + if name == "litellm-proxy": + continue + for entry in spec.get("egress", []): + assert all("ipBlock" not in rule for rule in entry["to"]), name + gateway_targets = { + rule["podSelector"]["matchLabels"]["app.kubernetes.io/name"] + for entry in policies["llm-gateway"]["egress"] + for rule in entry["to"] + } + assert gateway_targets == {"vllm-serve", "redis", "litellm-proxy"} + + +def test_the_proxy_reads_both_of_its_keys_from_the_secret_manager() -> None: + external_secret = next( + document for document in DOCUMENTS if document["kind"] == "ExternalSecret" + ) + proxy = next( + document for document in WORKLOADS if document["metadata"]["name"] == "litellm-proxy" + ) + provided = {item["secretKey"] for item in external_secret["spec"]["data"]} + + for variable in pod_spec(proxy)["containers"][0]["env"]: + assert "value" not in variable + assert variable["valueFrom"]["secretKeyRef"]["key"] in provided