Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 32 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
31 changes: 30 additions & 1 deletion src/llm_router/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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 (
Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
181 changes: 181 additions & 0 deletions src/llm_router/engine_stats.py
Original file line number Diff line number Diff line change
@@ -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
3 changes: 3 additions & 0 deletions src/llm_router/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
Loading
Loading