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
19 changes: 19 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
]
Expand Down
37 changes: 35 additions & 2 deletions src/llm_router/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,15 @@
ChatMessage,
Usage,
)
from llm_router.observability import Metrics
from llm_router.routing import NoEligibleModelError, Router, default_model_profiles


def create_app(
settings: Settings | None = None,
*,
backend: InferenceBackend | None = None,
metrics: Metrics | None = None,
) -> FastAPI:
runtime_settings = settings or get_settings()
router = Router(
Expand All @@ -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]:
Expand Down Expand Up @@ -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"},
Expand All @@ -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"},
Expand All @@ -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 = [
Expand All @@ -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
Expand Down
148 changes: 148 additions & 0 deletions src/llm_router/observability.py
Original file line number Diff line number Diff line change
@@ -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
21 changes: 21 additions & 0 deletions tests/e2e/api.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
});
53 changes: 53 additions & 0 deletions tests/integration/test_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading
Loading