From 9e5fa4b764fa53cafe47b276ebf663f1101b7ee8 Mon Sep 17 00:00:00 2001 From: Yash-Chindam Date: Sun, 30 Aug 2026 15:23:29 +0530 Subject: [PATCH 1/4] feat: add inference and routing metrics collectors Implements the section 15 metric surface: latency, time-to-first-token, time-per-output-token, token counters, queue and in-flight gauges, routing attribution, fallback frequency, queue-delay prediction error, cache outcomes, and model load duration. --- pyproject.toml | 1 + src/llm_router/observability.py | 148 +++++++++++++++++++++++++ tests/unit/test_observability.py | 181 +++++++++++++++++++++++++++++++ 3 files changed, 330 insertions(+) create mode 100644 src/llm_router/observability.py create mode 100644 tests/unit/test_observability.py diff --git a/pyproject.toml b/pyproject.toml index 525f1ec..3bbdc88 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -10,6 +10,7 @@ readme = "README.md" requires-python = ">=3.11" dependencies = [ "fastapi>=0.141.1,<1", + "prometheus-client>=0.26.0,<1", "pydantic-settings>=2.15.0,<3", "uvicorn[standard]>=0.52.4,<1", ] diff --git a/src/llm_router/observability.py b/src/llm_router/observability.py new file mode 100644 index 0000000..821d5a4 --- /dev/null +++ b/src/llm_router/observability.py @@ -0,0 +1,148 @@ +"""Inference and routing metrics defined in section 15 of the design specification.""" + +from prometheus_client import CONTENT_TYPE_LATEST, CollectorRegistry, Counter, Gauge, Histogram +from prometheus_client import generate_latest as render_registry + +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) + + +class Metrics: + """Prometheus collectors for request, routing, and engine telemetry.""" + + def __init__(self, registry: CollectorRegistry | None = None) -> None: + self.registry = registry if registry is not None else CollectorRegistry() + + self.requests_total = Counter( + "router_requests_total", + "Completed inference requests by selected model, task, and outcome.", + ["model", "task", "outcome"], + registry=self.registry, + ) + self.request_latency_seconds = Histogram( + "router_request_latency_seconds", + "End-to-end request latency observed by the gateway.", + ["model"], + buckets=LATENCY_BUCKETS, + registry=self.registry, + ) + self.time_to_first_token_seconds = Histogram( + "router_time_to_first_token_seconds", + "Time from admission to the first generated token.", + ["model"], + buckets=TOKEN_LATENCY_BUCKETS, + registry=self.registry, + ) + self.time_per_output_token_seconds = Histogram( + "router_time_per_output_token_seconds", + "Mean generation time per output token.", + ["model"], + buckets=TOKEN_LATENCY_BUCKETS, + registry=self.registry, + ) + self.tokens_total = Counter( + "router_tokens_total", + "Prompt and completion tokens processed.", + ["model", "kind"], + registry=self.registry, + ) + self.inflight_requests = Gauge( + "router_inflight_requests", + "Requests currently executing against an inference backend.", + registry=self.registry, + ) + self.queued_requests = Gauge( + "router_queued_requests", + "Requests waiting for an admission slot.", + registry=self.registry, + ) + self.rejections_total = Counter( + "router_rejections_total", + "Requests rejected before generation, by rejection type.", + ["type"], + registry=self.registry, + ) + self.routes_total = Counter( + "router_routes_total", + "Routing decisions by selected model, task, and privacy class.", + ["model", "task", "privacy"], + registry=self.registry, + ) + self.external_fallback_total = Counter( + "router_external_fallback_total", + "Requests dispatched to an approved external provider.", + ["model"], + registry=self.registry, + ) + self.predicted_quality = Histogram( + "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), + registry=self.registry, + ) + self.queue_delay_prediction_error_ms = Histogram( + "router_queue_delay_prediction_error_ms", + "Absolute error between predicted and observed queue delay.", + ["model"], + buckets=(1, 5, 10, 25, 50, 100, 250, 500, 1000), + registry=self.registry, + ) + self.cache_events_total = Counter( + "router_cache_events_total", + "Cache lookups by cache name and result.", + ["cache", "result"], + registry=self.registry, + ) + self.model_load_seconds = Histogram( + "router_model_load_seconds", + "Observed model load and cold-start duration.", + ["model"], + buckets=LATENCY_BUCKETS, + registry=self.registry, + ) + + def record_route(self, decision: RouteDecision, privacy: str) -> None: + self.routes_total.labels( + model=decision.profile.id, task=decision.task.value, privacy=privacy + ).inc() + self.predicted_quality.labels(model=decision.profile.id).observe(decision.profile.quality) + if not decision.profile.local: + self.external_fallback_total.labels(model=decision.profile.id).inc() + + def record_completion( + self, + decision: RouteDecision, + *, + latency_seconds: float, + queue_seconds: float, + prompt_tokens: int, + completion_tokens: int, + outcome: str = "success", + ) -> None: + model = decision.profile.id + self.requests_total.labels(model=model, task=decision.task.value, outcome=outcome).inc() + self.request_latency_seconds.labels(model=model).observe(latency_seconds) + self.time_to_first_token_seconds.labels(model=model).observe(queue_seconds) + if completion_tokens > 0: + generation_seconds = max(latency_seconds - queue_seconds, 0.0) + self.time_per_output_token_seconds.labels(model=model).observe( + generation_seconds / completion_tokens + ) + self.tokens_total.labels(model=model, kind="prompt").inc(prompt_tokens) + self.tokens_total.labels(model=model, kind="completion").inc(completion_tokens) + predicted_ms = float(decision.profile.estimated_queue_ms) + self.queue_delay_prediction_error_ms.labels(model=model).observe( + abs(predicted_ms - queue_seconds * 1000) + ) + + def record_rejection(self, rejection_type: str) -> None: + self.rejections_total.labels(type=rejection_type).inc() + + def record_cache_event(self, cache: str, result: str) -> None: + self.cache_events_total.labels(cache=cache, result=result).inc() + + def render(self) -> tuple[bytes, str]: + return render_registry(self.registry), CONTENT_TYPE_LATEST diff --git a/tests/unit/test_observability.py b/tests/unit/test_observability.py new file mode 100644 index 0000000..971ec75 --- /dev/null +++ b/tests/unit/test_observability.py @@ -0,0 +1,181 @@ +import pytest + +from llm_router.models import ModelProfile, RouteDecision, TaskClass +from llm_router.observability import Metrics + + +def build_decision() -> RouteDecision: + profile = ModelProfile( + id="general-local", + revision="mock-general@sha256:dev", + local=True, + context_limit=32768, + supported_tasks=frozenset({TaskClass.SUMMARIZATION}), + quality=0.89, + estimated_queue_ms=35, + ) + return RouteDecision( + profile=profile, + task=TaskClass.SUMMARIZATION, + reason="policy", + score=85.15, + candidate_count=2, + ) + + +def external_decision() -> RouteDecision: + profile = ModelProfile( + id="approved-external-fallback", + revision="external-policy-v1", + local=False, + context_limit=128000, + supported_tasks=frozenset(TaskClass), + quality=0.98, + estimated_queue_ms=45, + ) + return RouteDecision( + profile=profile, task=TaskClass.REASONING, reason="policy", score=1.0, candidate_count=1 + ) + + +def sample_value(metrics: Metrics, name: str, labels: dict[str, str]) -> float: + value = metrics.registry.get_sample_value(name, labels) + assert value is not None, f"missing sample {name}{labels}" + return value + + +def test_record_route_counts_decision_and_predicted_quality() -> None: + metrics = Metrics() + decision = build_decision() + + metrics.record_route(decision, privacy="private") + + assert ( + sample_value( + metrics, + "router_routes_total", + {"model": "general-local", "task": "summarization", "privacy": "private"}, + ) + == 1 + ) + assert sample_value(metrics, "router_predicted_quality_count", {"model": "general-local"}) == 1 + + +def test_record_route_counts_external_fallback_separately() -> None: + metrics = Metrics() + + metrics.record_route(external_decision(), privacy="public") + + assert ( + sample_value( + metrics, + "router_external_fallback_total", + {"model": "approved-external-fallback"}, + ) + == 1 + ) + + +def test_record_route_does_not_count_local_models_as_fallback() -> None: + metrics = Metrics() + + metrics.record_route(build_decision(), privacy="private") + + assert ( + metrics.registry.get_sample_value( + "router_external_fallback_total", {"model": "general-local"} + ) + is None + ) + + +def test_record_completion_tracks_tokens_latency_and_prediction_error() -> None: + metrics = Metrics() + decision = build_decision() + + metrics.record_completion( + decision, + latency_seconds=0.5, + queue_seconds=0.035, + prompt_tokens=12, + completion_tokens=4, + ) + + assert ( + sample_value( + metrics, + "router_requests_total", + {"model": "general-local", "task": "summarization", "outcome": "success"}, + ) + == 1 + ) + assert ( + sample_value(metrics, "router_tokens_total", {"model": "general-local", "kind": "prompt"}) + == 12 + ) + assert ( + sample_value( + metrics, "router_tokens_total", {"model": "general-local", "kind": "completion"} + ) + == 4 + ) + assert ( + sample_value( + metrics, "router_time_per_output_token_seconds_count", {"model": "general-local"} + ) + == 1 + ) + assert sample_value( + metrics, "router_queue_delay_prediction_error_ms_sum", {"model": "general-local"} + ) == pytest.approx(0.0, abs=1e-6) + + +def test_record_completion_skips_token_latency_without_output_tokens() -> None: + metrics = Metrics() + decision = build_decision() + + metrics.record_completion( + decision, + latency_seconds=0.2, + queue_seconds=0.01, + prompt_tokens=5, + completion_tokens=0, + outcome="error", + ) + + assert ( + metrics.registry.get_sample_value( + "router_time_per_output_token_seconds_count", {"model": "general-local"} + ) + is None + ) + assert ( + sample_value( + metrics, + "router_requests_total", + {"model": "general-local", "task": "summarization", "outcome": "error"}, + ) + == 1 + ) + + +def test_rejection_and_cache_events_are_labelled() -> None: + metrics = Metrics() + + metrics.record_rejection("quota_exceeded") + metrics.record_cache_event("exact", "hit") + + assert sample_value(metrics, "router_rejections_total", {"type": "quota_exceeded"}) == 1 + assert ( + sample_value(metrics, "router_cache_events_total", {"cache": "exact", "result": "hit"}) == 1 + ) + + +def test_render_returns_prometheus_exposition_payload() -> None: + metrics = Metrics() + metrics.record_rejection("overloaded") + + payload, content_type = metrics.render() + + assert b"router_rejections_total" in payload + assert content_type.startswith("text/plain") From e044f7a80863ee194f68b208512620d02815e12b Mon Sep 17 00:00:00 2001 From: Yash-Chindam Date: Sun, 30 Aug 2026 15:24:21 +0530 Subject: [PATCH 2/4] feat: expose Prometheus metrics from the gateway Records routing attribution, queue and in-flight depth, completion latency and tokens, and rejection counts, and serves them from an unauthenticated /metrics endpoint intended for in-cluster scraping. --- src/llm_router/app.py | 37 ++++++++++++++++++++++-- tests/integration/test_api.py | 53 +++++++++++++++++++++++++++++++++++ 2 files changed, 88 insertions(+), 2 deletions(-) diff --git a/src/llm_router/app.py b/src/llm_router/app.py index 76d6a2c..1ce0b33 100644 --- a/src/llm_router/app.py +++ b/src/llm_router/app.py @@ -23,6 +23,7 @@ ChatMessage, Usage, ) +from llm_router.observability import Metrics from llm_router.routing import NoEligibleModelError, Router, default_model_profiles @@ -30,6 +31,7 @@ def create_app( settings: Settings | None = None, *, backend: InferenceBackend | None = None, + metrics: Metrics | None = None, ) -> FastAPI: runtime_settings = settings or get_settings() router = Router( @@ -42,6 +44,7 @@ def create_app( ) quota = SlidingWindowQuota(runtime_settings.quota_requests_per_minute) inference_backend = backend or MockInferenceBackend() + telemetry = metrics if metrics is not None else Metrics() @asynccontextmanager async def lifespan(app: FastAPI) -> AsyncIterator[None]: @@ -77,10 +80,12 @@ async def authenticate(authorization: str | None = Header(default=None)) -> str: @app.exception_handler(NoEligibleModelError) async def no_model_handler(_: Request, error: NoEligibleModelError) -> JSONResponse: + telemetry.record_rejection("no_eligible_model") return JSONResponse(status_code=422, content={"error": {"message": str(error)}}) @app.exception_handler(AdmissionRejectedError) async def admission_handler(_: Request, error: AdmissionRejectedError) -> JSONResponse: + telemetry.record_rejection("overloaded") return JSONResponse( status_code=503, headers={"Retry-After": "1"}, @@ -89,6 +94,7 @@ async def admission_handler(_: Request, error: AdmissionRejectedError) -> JSONRe @app.exception_handler(QuotaExceededError) async def quota_handler(_: Request, error: QuotaExceededError) -> JSONResponse: + telemetry.record_rejection("quota_exceeded") return JSONResponse( status_code=429, headers={"Retry-After": "60"}, @@ -105,6 +111,11 @@ async def readiness(request: Request) -> dict[str, str]: raise HTTPException(status_code=503, detail="not ready") return {"status": "ready"} + @app.get("/metrics") + async def prometheus_metrics() -> Response: + payload, content_type = telemetry.render() + return Response(content=payload, media_type=content_type) + @app.get("/v1/models", dependencies=[Depends(authenticate)]) async def models() -> dict[str, object]: visible = [ @@ -126,10 +137,32 @@ async def chat_completions( response: Response, subject: str = Depends(authenticate), ) -> ChatCompletionResponse: + started = time.perf_counter() await quota.consume(subject) decision = router.select(payload) - async with admission.slot(): - result = await inference_backend.generate(payload, decision) + telemetry.record_route(decision, privacy=payload.routing.privacy.value) + + telemetry.queued_requests.inc() + try: + async with admission.slot(): + telemetry.queued_requests.dec() + queue_seconds = time.perf_counter() - started + telemetry.inflight_requests.inc() + try: + result = await inference_backend.generate(payload, decision) + finally: + telemetry.inflight_requests.dec() + except AdmissionRejectedError: + telemetry.queued_requests.dec() + raise + + telemetry.record_completion( + decision, + latency_seconds=time.perf_counter() - started, + queue_seconds=queue_seconds, + prompt_tokens=result.prompt_tokens, + completion_tokens=result.completion_tokens, + ) response.headers["X-Route-Model"] = decision.profile.id response.headers["X-Route-Revision"] = decision.profile.revision diff --git a/tests/integration/test_api.py b/tests/integration/test_api.py index 72cc4ae..3ee1fd6 100644 --- a/tests/integration/test_api.py +++ b/tests/integration/test_api.py @@ -89,3 +89,56 @@ def test_external_model_requires_public_data_and_explicit_opt_in() -> None: ) assert response.status_code == 200 assert response.json()["model"] == "approved-external-fallback" + + +def test_metrics_endpoint_reports_route_and_completion_telemetry() -> None: + with build_client() as client: + client.post( + "/v1/chat/completions", + headers={"Authorization": "Bearer integration-key"}, + json={ + "model": "auto", + "messages": [{"role": "user", "content": "Classify this ticket"}], + }, + ) + response = client.get("/metrics") + + assert response.status_code == 200 + assert response.headers["content-type"].startswith("text/plain") + body = response.text + assert "router_requests_total" in body + assert ( + 'router_routes_total{model="small-specialist",privacy="private",task="classification"}' + in body + ) + assert "router_tokens_total" in body + assert "router_inflight_requests 0.0" in body + assert "router_queued_requests 0.0" in body + + +def test_metrics_endpoint_counts_quota_rejections() -> None: + with build_client(quota=1) as client: + headers = {"Authorization": "Bearer integration-key"} + body = {"model": "auto", "messages": [{"role": "user", "content": "hello"}]} + assert client.post("/v1/chat/completions", headers=headers, json=body).status_code == 200 + assert client.post("/v1/chat/completions", headers=headers, json=body).status_code == 429 + metrics = client.get("/metrics").text + + assert 'router_rejections_total{type="quota_exceeded"} 1.0' in metrics + + +def test_metrics_endpoint_counts_policy_rejections() -> None: + with build_client() as client: + response = client.post( + "/v1/chat/completions", + headers={"Authorization": "Bearer integration-key"}, + json={ + "model": "approved-external-fallback", + "messages": [{"role": "user", "content": "hello"}], + "routing": {"privacy": "restricted", "allow_external_fallback": True}, + }, + ) + assert response.status_code == 422 + metrics = client.get("/metrics").text + + assert 'router_rejections_total{type="no_eligible_model"} 1.0' in metrics From 30d3e796d818bfbb42af2d90770b50eedaad0dba Mon Sep 17 00:00:00 2001 From: Yash-Chindam Date: Sun, 30 Aug 2026 15:24:49 +0530 Subject: [PATCH 3/4] test: verify metrics exposition end to end Confirms a live server publishes routed request counters, latency histograms, token totals, and predicted quality after a real completion. --- tests/e2e/api.spec.ts | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/tests/e2e/api.spec.ts b/tests/e2e/api.spec.ts index f4f415f..6293f7e 100644 --- a/tests/e2e/api.spec.ts +++ b/tests/e2e/api.spec.ts @@ -41,3 +41,24 @@ test("rejects invalid credentials", async ({ playwright }, testInfo) => { expect(response.status()).toBe(401); await anonymous.dispose(); }); + +test("publishes scrape-ready metrics for a completed request", async ({ request }) => { + const completion = await request.post("/v1/chat/completions", { + data: { + model: "auto", + messages: [{ role: "user", content: "Summarize the quarterly report" }], + }, + }); + expect(completion.status()).toBe(200); + const routedModel = completion.headers()["x-route-model"]; + + const metrics = await request.get("/metrics"); + expect(metrics.status()).toBe(200); + expect(metrics.headers()["content-type"]).toContain("text/plain"); + + const body = await metrics.text(); + expect(body).toContain(`router_requests_total{model="${routedModel}"`); + expect(body).toContain("router_request_latency_seconds_bucket"); + expect(body).toContain("router_tokens_total"); + expect(body).toContain("router_predicted_quality_sum"); +}); From c1033d74f05f06fe9df5252f608e1532ce64bd03 Mon Sep 17 00:00:00 2001 From: Yash-Chindam Date: Sun, 30 Aug 2026 15:24:50 +0530 Subject: [PATCH 4/4] docs: document the exported metric surface --- README.md | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/README.md b/README.md index 6af6b3b..c113725 100644 --- a/README.md +++ b/README.md @@ -61,6 +61,25 @@ Successful PR CI runs are merged automatically only for trusted same-repository and Dependabot. Forks, drafts, and untrusted author associations are deliberately skipped; repository branch-protection and review requirements continue to apply. +## Observability + +`GET /metrics` returns Prometheus exposition text and is intentionally unauthenticated so +in-cluster scrapers can read it; restrict it with network policy rather than a bearer token. + +| Metric | Purpose | +|---|---| +| `router_request_latency_seconds` | End-to-end latency histogram per model. | +| `router_time_to_first_token_seconds` | Admission-to-first-token delay. | +| `router_time_per_output_token_seconds` | Mean generation time per output token. | +| `router_tokens_total` | Prompt and completion tokens per model. | +| `router_inflight_requests` / `router_queued_requests` | Live capacity and queue depth. | +| `router_routes_total` | Requests per route with task and privacy class. | +| `router_external_fallback_total` | Fallback frequency. | +| `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. | + ## Runtime settings All settings use the `ROUTER_` prefix.