From ff3114dc0b8ccc5fc76b0c75a4848755009e1fec Mon Sep 17 00:00:00 2001 From: Yash-Chindam Date: Sat, 3 Oct 2026 19:09:18 +0530 Subject: [PATCH] feat: export engine, KV-cache, and accelerator telemetry Section 15 asks for KV-cache occupancy, batch size, GPU utilization and memory, and model-load duration, and section 7.4 makes exporting engine metrics an engine responsibility. None of it was reaching the gateway: the model-load collector was declared and never observed, and quality was only ever predicted, never compared with anything. The gateway now reads the serving path's own exposition (vLLM's vllm:* series plus DCGM_FI_DEV_* from a GPU exporter beside it) on each scrape and republishes the subset section 15 names, so one scrape answers for the whole path and no poller runs when nobody collects. A missing sample stays missing rather than being reported as zero, and an unreachable engine costs its series, not an error. Cold start is measured at the readiness transition instead of configured, and reopens on later recoveries so a reload after an eviction or a lost node is measured too. Live traffic is ungraded, so observed quality uses the one signal production exposes: whether a request declaring routing.structured returned parseable JSON, compared against the quality routing predicted. Co-Authored-By: Claude Opus 5 --- README.md | 33 ++- src/llm_router/app.py | 31 ++- src/llm_router/engine_stats.py | 181 +++++++++++++++++ src/llm_router/models.py | 3 + src/llm_router/observability.py | 121 ++++++++++- .../integration/test_engine_telemetry_api.py | 188 ++++++++++++++++++ tests/unit/test_engine_stats.py | 157 +++++++++++++++ tests/unit/test_observability.py | 122 ++++++++++++ 8 files changed, 833 insertions(+), 3 deletions(-) create mode 100644 src/llm_router/engine_stats.py create mode 100644 tests/integration/test_engine_telemetry_api.py create mode 100644 tests/unit/test_engine_stats.py 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)