diff --git a/README.md b/README.md index 2f9f520..c13f339 100644 --- a/README.md +++ b/README.md @@ -160,7 +160,38 @@ in-cluster scrapers can read it; restrict it with network policy rather than a b | `router_queue_delay_prediction_error_ms` | Predicted versus observed queue delay. | | `router_rejections_total` | Quota, overload, and policy rejections. | | `router_cache_events_total` | Cache lookups by cache and result. | -| `router_model_load_seconds` | Model load and cold-start duration. | +| `router_model_load_seconds` | Measured engine load and cold-start duration. | +| `router_observed_quality` / `router_quality_prediction_error` | Observed quality where it is checkable, against what routing predicted. | +| `router_structured_output_total` | Structured-output responses by validity. | + +Engine and accelerator state is pulled through from the serving path on each scrape, so the +gateway stays the single scrape target and no poller runs when nobody is collecting. The +gateway reads the engine's own `/metrics` (vLLM's `vllm:*` series, plus `DCGM_FI_DEV_*` from a +GPU exporter beside it) and republishes: + +| Metric | Purpose | +|---|---| +| `router_engine_running_requests` | Requests the engine is decoding — its live batch size. | +| `router_engine_batch_size` | Batch size sampled per scrape; the average is `_sum / _count`. | +| `router_engine_waiting_requests` | Requests queued inside the engine, not yet batched. | +| `router_engine_kv_cache_occupancy_ratio` | Fraction of the KV cache allocated. | +| `router_engine_preemptions_total` | Requests preempted under KV-cache pressure. | +| `router_gpu_utilization_ratio` | Accelerator utilization. | +| `router_gpu_memory_used_bytes` / `router_gpu_memory_total_bytes` | Framebuffer memory in use and installed. | + +An unreachable engine costs a scrape its engine series, never an error: a missing sample stays +missing rather than being reported as zero. + +Live traffic is ungraded, so the only quality signal observable in production is whether a +request that declared `routing.structured` actually returned parseable JSON. That is recorded +as an observed quality of 1 or 0 and compared with the quality routing predicted for the model +it picked. Full task-level quality comes from the offline harness below, never from sampled +traffic. + +Cold start is measured rather than configured: `/readyz` is the one place that sees the engine +go from loading to serving, so the duration of that window is observed there. The window +reopens on every later recovery, so a reload after an out-of-memory eviction or a lost node is +measured too — not only the first start. ## Caching diff --git a/src/llm_router/app.py b/src/llm_router/app.py index 3937e7a..36fd1cd 100644 --- a/src/llm_router/app.py +++ b/src/llm_router/app.py @@ -37,6 +37,8 @@ semantic_cache_eligible, ) from llm_router.config import Settings, get_settings +from llm_router.engine_stats import ColdStartTracker, EngineStatsCollector +from llm_router.evaluation import structured_output_valid from llm_router.models import ( ChatCompletionChoice, ChatCompletionRequest, @@ -78,6 +80,7 @@ def create_app( cache_store: CacheStore | None = None, registry: Registry | None = None, redis_client: RedisLike | None = None, + engine_stats: EngineStatsCollector | None = None, ) -> FastAPI: runtime_settings = settings or get_settings() catalog = registry if registry is not None else _load_catalog(runtime_settings.registry_path) @@ -114,6 +117,15 @@ def create_app( else MockInferenceBackend() ) 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) + if engine_client is not None + else None + ) + cold_start = ColdStartTracker() + # One gateway deployment faces one engine target, so engine-level telemetry + # and cold starts are attributed to that target rather than to a model. + engine_label = runtime_settings.backend exact_cache: CacheStore = ( cache_store or ( @@ -210,12 +222,25 @@ async def health() -> dict[str, str]: async def readiness(request: Request) -> dict[str, str]: if not getattr(request.app.state, "ready", False): raise HTTPException(status_code=503, detail="not ready") - if not await inference_backend.healthy(): + healthy = await inference_backend.healthy() + # The probe is the one place that sees the engine go from loading to + # serving, so cold start is measured here instead of being configured. + loaded_seconds = cold_start.observe(healthy=healthy, now=time.perf_counter()) + if loaded_seconds is not None: + telemetry.record_model_load(engine_label, loaded_seconds) + if not healthy: raise HTTPException(status_code=503, detail="inference backend is unhealthy") return {"status": "ready"} @app.get("/metrics") async def prometheus_metrics() -> Response: + # Engine and accelerator state is pulled through on scrape so the + # gateway stays the single scrape target for the whole serving path and + # no background poller runs when nobody is collecting. + if engine_telemetry is not None: + stats = await engine_telemetry.sample() + if stats is not None: + telemetry.record_engine_stats(stats, engine=engine_label) payload, content_type = telemetry.render() return Response(content=payload, media_type=content_type) @@ -437,6 +462,8 @@ async def _stream_completion( text = "".join(collected) prompt_tokens = max(1, len(prompt) // 4) completion_tokens = max(1, len(text) // 4) + if payload.routing.structured: + telemetry.record_structured_output(decision, valid=structured_output_valid(text)) telemetry.record_completion( decision, latency_seconds=time.perf_counter() - started, @@ -529,6 +556,8 @@ async def chat_completions( telemetry.queued_requests.dec() raise + if payload.routing.structured: + telemetry.record_structured_output(decision, valid=structured_output_valid(result.text)) telemetry.record_completion( decision, latency_seconds=time.perf_counter() - started, diff --git a/src/llm_router/engine_stats.py b/src/llm_router/engine_stats.py new file mode 100644 index 0000000..0ead957 --- /dev/null +++ b/src/llm_router/engine_stats.py @@ -0,0 +1,181 @@ +"""Engine and accelerator telemetry scraped from the serving path (7.4, 15). + +The gateway can only see a request from the outside. KV-cache occupancy, the +engine's current batch, preemptions, and GPU utilization exist inside the vLLM +engine and its GPU exporter. This module parses the Prometheus exposition those +already publish and hands the subset section 15 requires back to the gateway's +own registry, so one scrape of `/metrics` answers for the whole serving path +rather than only for the hop the gateway performed. + +Nothing here fails a request: telemetry is best effort, and an unreachable +engine yields no sample rather than an error. +""" + +import math +from dataclasses import dataclass, field + +import httpx + +# vLLM publishes engine state under the `vllm:` prefix; a DCGM exporter running +# beside the engine publishes accelerator state under `DCGM_FI_DEV_*`. Both are +# read from the same endpoint so the serving pod stays a single scrape target. +RUNNING_REQUESTS = "vllm:num_requests_running" +WAITING_REQUESTS = "vllm:num_requests_waiting" +KV_CACHE_USAGE = "vllm:gpu_cache_usage_perc" +PREEMPTIONS = "vllm:num_preemptions_total" +GPU_UTILIZATION = "DCGM_FI_DEV_GPU_UTIL" +GPU_MEMORY_USED = "DCGM_FI_DEV_FB_USED" +GPU_MEMORY_TOTAL = "DCGM_FI_DEV_FB_TOTAL" + +GPU_LABELS = ("gpu", "device", "UUID") + + +@dataclass(frozen=True) +class EngineStats: + """One observation of engine and accelerator state. + + Every field is optional because an engine may not publish it, and a missing + sample must stay missing rather than be reported as zero. + """ + + running_requests: float | None = None + waiting_requests: float | None = None + kv_cache_usage_ratio: float | None = None + preemptions_total: float | None = None + gpu_utilization_ratio: dict[str, float] = field(default_factory=dict) + gpu_memory_used_bytes: dict[str, float] = field(default_factory=dict) + gpu_memory_total_bytes: dict[str, float] = field(default_factory=dict) + + @property + def empty(self) -> bool: + return self == EngineStats() + + +def _parse_labels(block: str) -> dict[str, str]: + labels: dict[str, str] = {} + for part in block.split(","): + name, separator, value = part.partition("=") + if not separator: + continue + labels[name.strip()] = value.strip().strip('"') + return labels + + +def parse_exposition(payload: str) -> dict[str, list[tuple[dict[str, str], float]]]: + """Parse Prometheus text exposition into samples keyed by metric name. + + Comments, blank lines, unparseable values, and NaN are skipped: an engine + reporting NaN for a ratio it has not computed yet is not a measurement. + """ + + samples: dict[str, list[tuple[dict[str, str], float]]] = {} + for line in payload.splitlines(): + stripped = line.strip() + if not stripped or stripped.startswith("#"): + continue + head, _, raw_value = stripped.rpartition(" ") + if not head: + continue + try: + value = float(raw_value) + except ValueError: + continue + if math.isnan(value) or math.isinf(value): + continue + name, bracket, remainder = head.partition("{") + labels = _parse_labels(remainder.rstrip("}")) if bracket else {} + samples.setdefault(name.strip(), []).append((labels, value)) + return samples + + +def _scalar( + samples: dict[str, list[tuple[dict[str, str], float]]], name: str, *, total: bool = False +) -> float | None: + """Collapse a metric to one number, summing when replicas report separately.""" + + found = samples.get(name) + if not found: + return None + values = [value for _, value in found] + return sum(values) if total or len(values) > 1 else values[0] + + +def _by_gpu(samples: dict[str, list[tuple[dict[str, str], float]]], name: str) -> dict[str, float]: + """Index per-accelerator samples by whichever device label the exporter used.""" + + indexed: dict[str, float] = {} + for position, (labels, value) in enumerate(samples.get(name, [])): + key = next((labels[label] for label in GPU_LABELS if labels.get(label)), str(position)) + indexed[key] = value + return indexed + + +def parse_engine_stats(payload: str) -> EngineStats: + """Read engine and accelerator state out of a Prometheus exposition payload.""" + + samples = parse_exposition(payload) + utilization = _by_gpu(samples, GPU_UTILIZATION) + return EngineStats( + running_requests=_scalar(samples, RUNNING_REQUESTS, total=True), + waiting_requests=_scalar(samples, WAITING_REQUESTS, total=True), + kv_cache_usage_ratio=_scalar(samples, KV_CACHE_USAGE), + preemptions_total=_scalar(samples, PREEMPTIONS, total=True), + # DCGM reports utilization as a percentage; section 15 reports a ratio. + gpu_utilization_ratio={gpu: value / 100 for gpu, value in utilization.items()}, + # DCGM reports framebuffer memory in mebibytes. + gpu_memory_used_bytes={ + gpu: value * 1024 * 1024 for gpu, value in _by_gpu(samples, GPU_MEMORY_USED).items() + }, + gpu_memory_total_bytes={ + gpu: value * 1024 * 1024 for gpu, value in _by_gpu(samples, GPU_MEMORY_TOTAL).items() + }, + ) + + +@dataclass +class ColdStartTracker: + """Measures engine load time across the readiness transition (13, 15). + + Section 13 asks for cold-start time to be documented rather than assumed. + The readiness probe already observes the engine going from unready to + serving, so the duration is measured there. The window reopens on every + later recovery, so a reload after an out-of-memory eviction or a lost node + is measured too rather than only the first start. + """ + + _unready_since: float | None = None + + def observe(self, *, healthy: bool, now: float) -> float | None: + """Return the completed load duration, or None while nothing completed.""" + + if not healthy: + if self._unready_since is None: + self._unready_since = now + return None + if self._unready_since is None: + return None + elapsed = now - self._unready_since + self._unready_since = None + return elapsed + + +@dataclass +class EngineStatsCollector: + """Pulls engine telemetry on demand, tolerating an unreachable engine.""" + + base_url: str + client: httpx.AsyncClient + path: str = "/metrics" + timeout_seconds: float = 2.0 + + async def sample(self) -> EngineStats | None: + try: + response = await self.client.get( + f"{self.base_url}{self.path}", timeout=self.timeout_seconds + ) + except httpx.HTTPError: + return None + if response.status_code >= 400: + return None + stats = parse_engine_stats(response.text) + return None if stats.empty else stats diff --git a/src/llm_router/models.py b/src/llm_router/models.py index c7b521e..8e00377 100644 --- a/src/llm_router/models.py +++ b/src/llm_router/models.py @@ -32,6 +32,9 @@ class RoutingOptions(BaseModel): latency_tier: Literal["interactive", "standard", "batch"] = "standard" quality_floor: float = Field(default=0.0, ge=0.0, le=1.0) allow_external_fallback: bool = False + # Declaring a structured-output requirement lets the gateway check the one + # quality signal live traffic exposes: whether the response actually parsed. + structured: bool = False class ChatCompletionRequest(BaseModel): diff --git a/src/llm_router/observability.py b/src/llm_router/observability.py index 821d5a4..0457683 100644 --- a/src/llm_router/observability.py +++ b/src/llm_router/observability.py @@ -3,10 +3,13 @@ from prometheus_client import CONTENT_TYPE_LATEST, CollectorRegistry, Counter, Gauge, Histogram from prometheus_client import generate_latest as render_registry +from llm_router.engine_stats import EngineStats from llm_router.models import RouteDecision LATENCY_BUCKETS = (0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0) TOKEN_LATENCY_BUCKETS = (0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5) +QUALITY_BUCKETS = (0.5, 0.6, 0.7, 0.8, 0.85, 0.9, 0.95, 1.0) +BATCH_SIZE_BUCKETS = (1, 2, 4, 8, 16, 32, 64, 128, 256) class Metrics: @@ -80,7 +83,30 @@ def __init__(self, registry: CollectorRegistry | None = None) -> None: "router_predicted_quality", "Predicted quality of the selected model at decision time.", ["model"], - buckets=(0.5, 0.6, 0.7, 0.8, 0.85, 0.9, 0.95, 1.0), + buckets=QUALITY_BUCKETS, + registry=self.registry, + ) + self.observed_quality = Histogram( + "router_observed_quality", + # Live traffic is ungraded, so the only quality signal observable in + # production is whether a structured request produced valid output. + # Full quality comes from the offline harness in section 16. + "Observed quality of served responses where it is checkable.", + ["model", "signal"], + buckets=QUALITY_BUCKETS, + registry=self.registry, + ) + self.quality_prediction_error = Histogram( + "router_quality_prediction_error", + "Absolute error between predicted and observed quality.", + ["model", "signal"], + buckets=(0.01, 0.05, 0.1, 0.2, 0.3, 0.5, 1.0), + registry=self.registry, + ) + self.structured_output_total = Counter( + "router_structured_output_total", + "Structured-output responses by validity.", + ["model", "result"], registry=self.registry, ) self.queue_delay_prediction_error_ms = Histogram( @@ -103,6 +129,55 @@ def __init__(self, registry: CollectorRegistry | None = None) -> None: buckets=LATENCY_BUCKETS, registry=self.registry, ) + self.engine_running_requests = Gauge( + "router_engine_running_requests", + "Requests the engine is currently decoding, which is its live batch size.", + ["engine"], + registry=self.registry, + ) + self.engine_waiting_requests = Gauge( + "router_engine_waiting_requests", + "Requests queued inside the engine and not yet batched.", + ["engine"], + registry=self.registry, + ) + self.engine_batch_size = Histogram( + "router_engine_batch_size", + "Engine batch size sampled per scrape; the average is sum over count.", + ["engine"], + buckets=BATCH_SIZE_BUCKETS, + registry=self.registry, + ) + self.engine_kv_cache_occupancy_ratio = Gauge( + "router_engine_kv_cache_occupancy_ratio", + "Fraction of the engine's KV cache currently allocated.", + ["engine"], + registry=self.registry, + ) + self.engine_preemptions_total = Gauge( + "router_engine_preemptions_total", + "Requests the engine preempted under KV-cache pressure, as reported by the engine.", + ["engine"], + registry=self.registry, + ) + self.gpu_utilization_ratio = Gauge( + "router_gpu_utilization_ratio", + "Accelerator utilization reported by the GPU exporter.", + ["gpu"], + registry=self.registry, + ) + self.gpu_memory_used_bytes = Gauge( + "router_gpu_memory_used_bytes", + "Accelerator framebuffer memory in use.", + ["gpu"], + registry=self.registry, + ) + self.gpu_memory_total_bytes = Gauge( + "router_gpu_memory_total_bytes", + "Accelerator framebuffer memory installed.", + ["gpu"], + registry=self.registry, + ) def record_route(self, decision: RouteDecision, privacy: str) -> None: self.routes_total.labels( @@ -144,5 +219,49 @@ def record_rejection(self, rejection_type: str) -> None: def record_cache_event(self, cache: str, result: str) -> None: self.cache_events_total.labels(cache=cache, result=result).inc() + def record_structured_output(self, decision: RouteDecision, *, valid: bool) -> None: + """Record the one quality signal live traffic actually exposes. + + A structured request either parsed or it did not, so validity is an + observed quality of 1 or 0 and can be compared with what routing + predicted for the model it chose. + """ + + model = decision.profile.id + observed = 1.0 if valid else 0.0 + self.structured_output_total.labels( + model=model, result="valid" if valid else "invalid" + ).inc() + self.observed_quality.labels(model=model, signal="structured_validity").observe(observed) + self.quality_prediction_error.labels(model=model, signal="structured_validity").observe( + abs(decision.profile.quality - observed) + ) + + def record_model_load(self, model: str, seconds: float) -> None: + """Record a measured cold start, from first unready observation to ready.""" + + self.model_load_seconds.labels(model=model).observe(seconds) + + def record_engine_stats(self, stats: EngineStats, *, engine: str) -> None: + """Publish engine and accelerator state, leaving absent samples absent.""" + + if stats.running_requests is not None: + self.engine_running_requests.labels(engine=engine).set(stats.running_requests) + self.engine_batch_size.labels(engine=engine).observe(stats.running_requests) + if stats.waiting_requests is not None: + self.engine_waiting_requests.labels(engine=engine).set(stats.waiting_requests) + if stats.kv_cache_usage_ratio is not None: + self.engine_kv_cache_occupancy_ratio.labels(engine=engine).set( + stats.kv_cache_usage_ratio + ) + if stats.preemptions_total is not None: + self.engine_preemptions_total.labels(engine=engine).set(stats.preemptions_total) + for gpu, ratio in stats.gpu_utilization_ratio.items(): + self.gpu_utilization_ratio.labels(gpu=gpu).set(ratio) + for gpu, used in stats.gpu_memory_used_bytes.items(): + self.gpu_memory_used_bytes.labels(gpu=gpu).set(used) + for gpu, total in stats.gpu_memory_total_bytes.items(): + self.gpu_memory_total_bytes.labels(gpu=gpu).set(total) + def render(self) -> tuple[bytes, str]: return render_registry(self.registry), CONTENT_TYPE_LATEST diff --git a/tests/integration/test_engine_telemetry_api.py b/tests/integration/test_engine_telemetry_api.py new file mode 100644 index 0000000..8024c32 --- /dev/null +++ b/tests/integration/test_engine_telemetry_api.py @@ -0,0 +1,188 @@ +from collections.abc import AsyncIterator + +import httpx +from fastapi.testclient import TestClient + +from llm_router.app import create_app +from llm_router.backends import BackendResult, MockInferenceBackend +from llm_router.config import Settings +from llm_router.engine_stats import EngineStatsCollector +from llm_router.models import ChatCompletionRequest, RouteDecision + +HEADERS = {"Authorization": "Bearer telemetry-key"} +EXPOSITION = """ +vllm:num_requests_running{model_name="general-local"} 9.0 +vllm:num_requests_waiting{model_name="general-local"} 2.0 +vllm:gpu_cache_usage_perc{model_name="general-local"} 0.64 +DCGM_FI_DEV_GPU_UTIL{gpu="0"} 91.0 +DCGM_FI_DEV_FB_USED{gpu="0"} 1024.0 +""" + + +class StructuredBackend(MockInferenceBackend): + """Returns whatever text the test needs so validity can be asserted.""" + + def __init__(self, text: str) -> None: + self.text = text + + async def generate( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> BackendResult: + return BackendResult(text=self.text, prompt_tokens=8, completion_tokens=4) + + async def stream( + self, request: ChatCompletionRequest, decision: RouteDecision + ) -> AsyncIterator[str]: + yield self.text + + +class FlappingBackend(MockInferenceBackend): + """Reports unhealthy until a later probe, so a cold start can be measured.""" + + def __init__(self, healthy_from_probe: int) -> None: + self.probes = 0 + self.healthy_from_probe = healthy_from_probe + + async def healthy(self) -> bool: + self.probes += 1 + return self.probes >= self.healthy_from_probe + + +def build_client( + *, backend: object | None = None, engine_stats: EngineStatsCollector | None = None +) -> TestClient: + settings = Settings(api_keys="telemetry-key") + return TestClient( + create_app( + settings, + backend=backend, # type: ignore[arg-type] + engine_stats=engine_stats, + ) + ) + + +def collector_for(payload: str, *, status: int = 200) -> EngineStatsCollector: + transport = httpx.MockTransport(lambda _: httpx.Response(status, text=payload)) + return EngineStatsCollector( + base_url="http://engine:8001", client=httpx.AsyncClient(transport=transport) + ) + + +def test_metrics_scrape_republishes_engine_and_gpu_state() -> None: + with build_client(engine_stats=collector_for(EXPOSITION)) as client: + body = client.get("/metrics").text + + assert 'router_engine_running_requests{engine="mock"} 9.0' in body + assert 'router_engine_waiting_requests{engine="mock"} 2.0' in body + assert 'router_engine_kv_cache_occupancy_ratio{engine="mock"} 0.64' in body + assert 'router_gpu_utilization_ratio{gpu="0"} 0.91' in body + assert 'router_gpu_memory_used_bytes{gpu="0"} 1.073741824e+09' in body + # The live batch is also sampled into the histogram so an average exists. + assert 'router_engine_batch_size_count{engine="mock"} 1.0' in body + + +def test_metrics_stay_available_when_the_engine_cannot_be_scraped() -> None: + with build_client(engine_stats=collector_for("", status=503)) as client: + response = client.get("/metrics") + + assert response.status_code == 200 + assert "router_requests_total" in response.text + # The family is declared, so absence has to be asserted on the samples. + assert "router_engine_running_requests{engine=" not in response.text + + +def test_metrics_need_no_engine_collector_at_all() -> None: + with build_client() as client: + response = client.get("/metrics") + + assert response.status_code == 200 + assert "router_rejections_total" in response.text + + +def test_a_structured_request_that_returns_json_records_observed_quality() -> None: + with build_client(backend=StructuredBackend('{"invoice": "INV-1"}')) as client: + completion = client.post( + "/v1/chat/completions", + headers=HEADERS, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Extract the invoice fields"}], + "routing": {"privacy": "private", "structured": True}, + }, + ) + body = client.get("/metrics").text + + assert completion.status_code == 200 + assert 'router_structured_output_total{model="small-specialist",result="valid"} 1.0' in body + assert 'signal="structured_validity"' in body + + +def test_a_structured_request_that_returns_prose_is_observed_quality_zero() -> None: + with build_client(backend=StructuredBackend("the invoice number is INV-1")) as client: + client.post( + "/v1/chat/completions", + headers=HEADERS, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Extract the invoice fields"}], + "routing": {"privacy": "private", "structured": True}, + }, + ) + body = client.get("/metrics").text + + assert 'router_structured_output_total{model="small-specialist",result="invalid"} 1.0' in body + + +def test_an_unstructured_request_records_no_validity_signal() -> None: + with build_client(backend=StructuredBackend("plain prose")) as client: + client.post( + "/v1/chat/completions", + headers=HEADERS, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Summarize the report"}], + "routing": {"privacy": "private"}, + }, + ) + body = client.get("/metrics").text + + assert "router_structured_output_total{" not in body + + +def test_a_streamed_structured_response_is_checked_too() -> None: + with build_client(backend=StructuredBackend('{"claim": "CLM-9"}')) as client: + with client.stream( + "POST", + "/v1/chat/completions", + headers=HEADERS, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Extract the claim id"}], + "stream": True, + "routing": {"privacy": "private", "structured": True}, + }, + ) as response: + assert response.status_code == 200 + list(response.iter_lines()) + body = client.get("/metrics").text + + assert 'router_structured_output_total{model="small-specialist",result="valid"} 1.0' in body + + +def test_readiness_measures_the_cold_start_it_observes() -> None: + backend = FlappingBackend(healthy_from_probe=3) + with build_client(backend=backend) as client: + assert client.get("/readyz").status_code == 503 + assert client.get("/readyz").status_code == 503 + assert client.get("/readyz").status_code == 200 + body = client.get("/metrics").text + + assert 'router_model_load_seconds_count{model="mock"} 1.0' in body + + +def test_an_engine_ready_on_the_first_probe_reports_no_cold_start() -> None: + with build_client(backend=MockInferenceBackend()) as client: + assert client.get("/readyz").status_code == 200 + body = client.get("/metrics").text + + assert "router_model_load_seconds_count" not in body diff --git a/tests/unit/test_engine_stats.py b/tests/unit/test_engine_stats.py new file mode 100644 index 0000000..3be72c0 --- /dev/null +++ b/tests/unit/test_engine_stats.py @@ -0,0 +1,157 @@ +import httpx +import pytest + +from llm_router.engine_stats import ( + ColdStartTracker, + EngineStats, + EngineStatsCollector, + parse_engine_stats, + parse_exposition, +) + +VLLM_EXPOSITION = """ +# HELP vllm:num_requests_running Number of requests currently running. +# TYPE vllm:num_requests_running gauge +vllm:num_requests_running{model_name="general-local"} 12.0 +vllm:num_requests_waiting{model_name="general-local"} 3.0 +vllm:gpu_cache_usage_perc{model_name="general-local"} 0.73 +vllm:num_preemptions_total{model_name="general-local"} 4.0 +DCGM_FI_DEV_GPU_UTIL{gpu="0",UUID="GPU-abc"} 87.0 +DCGM_FI_DEV_FB_USED{gpu="0",UUID="GPU-abc"} 20480.0 +DCGM_FI_DEV_FB_TOTAL{gpu="0",UUID="GPU-abc"} 24576.0 +""" + + +def test_exposition_parser_keeps_labels_and_skips_comments() -> None: + samples = parse_exposition(VLLM_EXPOSITION) + + assert samples["vllm:num_requests_running"] == [({"model_name": "general-local"}, 12.0)] + assert "# HELP vllm:num_requests_running" not in samples + assert len(samples) == 7 + + +def test_exposition_parser_skips_unparseable_and_non_finite_values() -> None: + samples = parse_exposition( + "\n".join( + ( + "vllm:gpu_cache_usage_perc NaN", + "vllm:num_requests_waiting +Inf", + "vllm:num_requests_running not_a_number", + "malformed_line_without_value", + "", + "vllm:num_preemptions_total 2.0", + ) + ) + ) + + assert samples == {"vllm:num_preemptions_total": [({}, 2.0)]} + + +def test_engine_stats_convert_exporter_units() -> None: + stats = parse_engine_stats(VLLM_EXPOSITION) + + assert stats.running_requests == 12.0 + assert stats.waiting_requests == 3.0 + assert stats.kv_cache_usage_ratio == 0.73 + assert stats.preemptions_total == 4.0 + # DCGM reports a percentage and mebibytes; section 15 reports a ratio and bytes. + assert stats.gpu_utilization_ratio == {"0": pytest.approx(0.87)} + assert stats.gpu_memory_used_bytes == {"0": 20480.0 * 1024 * 1024} + assert stats.gpu_memory_total_bytes == {"0": 24576.0 * 1024 * 1024} + assert stats.empty is False + + +def test_engine_stats_sum_replicas_that_report_separately() -> None: + stats = parse_engine_stats( + "\n".join( + ( + 'vllm:num_requests_running{model_name="a"} 4.0', + 'vllm:num_requests_running{model_name="b"} 6.0', + ) + ) + ) + + assert stats.running_requests == 10.0 + + +def test_missing_samples_stay_absent_rather_than_zero() -> None: + stats = parse_engine_stats('vllm:num_requests_running{model_name="a"} 1.0') + + assert stats.running_requests == 1.0 + assert stats.kv_cache_usage_ratio is None + assert stats.waiting_requests is None + assert stats.gpu_utilization_ratio == {} + + +def test_an_engine_reporting_nothing_useful_is_empty() -> None: + assert parse_engine_stats("# nothing but comments\n").empty is True + assert EngineStats().empty is True + + +def test_gpu_samples_fall_back_to_position_without_a_device_label() -> None: + stats = parse_engine_stats('DCGM_FI_DEV_GPU_UTIL 50.0\nDCGM_FI_DEV_GPU_UTIL{other="x"} 70.0') + + assert stats.gpu_utilization_ratio == {"0": pytest.approx(0.5), "1": pytest.approx(0.7)} + + +@pytest.mark.asyncio +async def test_collector_samples_a_reachable_engine() -> None: + transport = httpx.MockTransport(lambda _: httpx.Response(200, text=VLLM_EXPOSITION)) + async with httpx.AsyncClient(transport=transport) as client: + collector = EngineStatsCollector(base_url="http://engine:8001", client=client) + + stats = await collector.sample() + + assert stats is not None + assert stats.kv_cache_usage_ratio == 0.73 + + +@pytest.mark.asyncio +async def test_collector_returns_nothing_when_the_engine_is_unreachable() -> None: + def explode(_: httpx.Request) -> httpx.Response: + raise httpx.ConnectError("refused") + + async with httpx.AsyncClient(transport=httpx.MockTransport(explode)) as client: + collector = EngineStatsCollector(base_url="http://engine:8001", client=client) + + assert await collector.sample() is None + + +@pytest.mark.asyncio +async def test_collector_returns_nothing_for_an_error_response_or_empty_payload() -> None: + async with httpx.AsyncClient( + transport=httpx.MockTransport(lambda _: httpx.Response(503)) + ) as client: + assert await EngineStatsCollector("http://engine:8001", client).sample() is None + + async with httpx.AsyncClient( + transport=httpx.MockTransport(lambda _: httpx.Response(200, text="# empty\n")) + ) as client: + assert await EngineStatsCollector("http://engine:8001", client).sample() is None + + +def test_cold_start_measures_the_unready_to_ready_window() -> None: + tracker = ColdStartTracker() + + assert tracker.observe(healthy=False, now=10.0) is None + assert tracker.observe(healthy=False, now=12.0) is None + assert tracker.observe(healthy=True, now=41.5) == pytest.approx(31.5) + + +def test_an_engine_already_warm_reports_no_cold_start() -> None: + tracker = ColdStartTracker() + + assert tracker.observe(healthy=True, now=5.0) is None + assert tracker.observe(healthy=True, now=6.0) is None + + +def test_every_later_reload_is_measured_too() -> None: + tracker = ColdStartTracker() + tracker.observe(healthy=False, now=0.0) + + first = tracker.observe(healthy=True, now=20.0) + tracker.observe(healthy=False, now=100.0) + second = tracker.observe(healthy=True, now=130.0) + + assert first == pytest.approx(20.0) + assert second == pytest.approx(30.0) diff --git a/tests/unit/test_observability.py b/tests/unit/test_observability.py index 971ec75..5d4d942 100644 --- a/tests/unit/test_observability.py +++ b/tests/unit/test_observability.py @@ -1,5 +1,6 @@ import pytest +from llm_router.engine_stats import EngineStats from llm_router.models import ModelProfile, RouteDecision, TaskClass from llm_router.observability import Metrics @@ -179,3 +180,124 @@ def test_render_returns_prometheus_exposition_payload() -> None: assert b"router_rejections_total" in payload assert content_type.startswith("text/plain") + + +def test_engine_stats_publish_kv_cache_batch_and_gpu_state() -> None: + metrics = Metrics() + stats = EngineStats( + running_requests=12.0, + waiting_requests=3.0, + kv_cache_usage_ratio=0.73, + preemptions_total=4.0, + gpu_utilization_ratio={"0": 0.87}, + gpu_memory_used_bytes={"0": 2048.0}, + gpu_memory_total_bytes={"0": 4096.0}, + ) + + metrics.record_engine_stats(stats, engine="vllm") + + assert sample_value(metrics, "router_engine_running_requests", {"engine": "vllm"}) == 12 + assert sample_value(metrics, "router_engine_waiting_requests", {"engine": "vllm"}) == 3 + assert sample_value( + metrics, "router_engine_kv_cache_occupancy_ratio", {"engine": "vllm"} + ) == pytest.approx(0.73) + assert sample_value(metrics, "router_engine_preemptions_total", {"engine": "vllm"}) == 4 + assert sample_value(metrics, "router_gpu_utilization_ratio", {"gpu": "0"}) == pytest.approx( + 0.87 + ) + assert sample_value(metrics, "router_gpu_memory_used_bytes", {"gpu": "0"}) == 2048 + assert sample_value(metrics, "router_gpu_memory_total_bytes", {"gpu": "0"}) == 4096 + + +def test_batch_size_average_is_recoverable_from_the_histogram() -> None: + metrics = Metrics() + + metrics.record_engine_stats(EngineStats(running_requests=4.0), engine="vllm") + metrics.record_engine_stats(EngineStats(running_requests=8.0), engine="vllm") + + total = sample_value(metrics, "router_engine_batch_size_sum", {"engine": "vllm"}) + count = sample_value(metrics, "router_engine_batch_size_count", {"engine": "vllm"}) + assert total / count == pytest.approx(6.0) + + +def test_absent_engine_samples_publish_no_series() -> None: + metrics = Metrics() + + metrics.record_engine_stats(EngineStats(running_requests=2.0), engine="vllm") + + assert ( + metrics.registry.get_sample_value( + "router_engine_kv_cache_occupancy_ratio", {"engine": "vllm"} + ) + is None + ) + assert metrics.registry.get_sample_value("router_gpu_utilization_ratio", {"gpu": "0"}) is None + + +def test_structured_validity_is_recorded_as_observed_quality() -> None: + metrics = Metrics() + decision = build_decision() + + metrics.record_structured_output(decision, valid=True) + + assert ( + sample_value( + metrics, + "router_structured_output_total", + {"model": "general-local", "result": "valid"}, + ) + == 1 + ) + assert ( + sample_value( + metrics, + "router_observed_quality_sum", + {"model": "general-local", "signal": "structured_validity"}, + ) + == 1.0 + ) + # The model predicted 0.89 and the response parsed, so the error is 0.11. + assert sample_value( + metrics, + "router_quality_prediction_error_sum", + {"model": "general-local", "signal": "structured_validity"}, + ) == pytest.approx(0.11) + + +def test_invalid_structured_output_is_observed_quality_zero() -> None: + metrics = Metrics() + decision = build_decision() + + metrics.record_structured_output(decision, valid=False) + + assert ( + sample_value( + metrics, + "router_structured_output_total", + {"model": "general-local", "result": "invalid"}, + ) + == 1 + ) + assert ( + sample_value( + metrics, + "router_observed_quality_sum", + {"model": "general-local", "signal": "structured_validity"}, + ) + == 0.0 + ) + assert sample_value( + metrics, + "router_quality_prediction_error_sum", + {"model": "general-local", "signal": "structured_validity"}, + ) == pytest.approx(0.89) + + +def test_model_load_duration_is_observed() -> None: + metrics = Metrics() + + metrics.record_model_load("vllm", 31.5) + + assert sample_value( + metrics, "router_model_load_seconds_sum", {"model": "vllm"} + ) == pytest.approx(31.5)