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
27 changes: 27 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -227,6 +227,30 @@ go from loading to serving, so the duration of that window is observed there. Th
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.

### Tracing

Every chat completion is one OpenTelemetry span carrying what is needed to explain the route:
tenant, effective and declared privacy class, task, model and revision, adapter, cache result,
score, candidate count, the route reason, and token usage. The response quotes it as
`X-Trace-Id`. Only the OpenTelemetry API is a runtime dependency, so tracing is a no-op until a
collector is configured:

```bash
python -m pip install -e ".[tracing]"
ROUTER_OTLP_ENDPOINT=http://otel-collector:4318/v1/traces
```

Prompts are redacted by privacy class, evaluated after any tenant floor has been applied:

| Class | Recorded |
|---|---|
| `restricted` | Length only. No digest: a digest of a short or templated prompt can be reversed by guessing. |
| `private` | Length and a SHA-256 digest, so repeats can be correlated. |
| `public` | Length and digest; a bounded prefix of the content only with `ROUTER_TRACE_PROMPT_CONTENT=true`. |

Completions are never recorded. A failed request records its error type and not its message,
because an engine error can echo the request it rejected.

## Caching

Cache keys bind the tenant, the active model-catalog fingerprint, and the generation
Expand Down Expand Up @@ -278,6 +302,9 @@ All settings use the `ROUTER_` prefix.
| `ROUTER_QUOTA_REQUESTS_PER_MINUTE` | `120` | Per-token sliding-window quota. |
| `ROUTER_EXTERNAL_FALLBACK_ENABLED` | `false` | Operator gate for external fallback. |
| `ROUTER_REDIS_URL` | _(empty)_ | Shared cache and quota state; in-process when empty. |
| `ROUTER_TENANT_KEYS` | _(empty)_ | `tenant:key` bindings; bare `ROUTER_API_KEYS` keys use the default tenant. |
| `ROUTER_OTLP_ENDPOINT` | _(empty)_ | OTLP/HTTP trace collector; tracing is a no-op when empty. |
| `ROUTER_TRACE_PROMPT_CONTENT` | `false` | Records a bounded prefix of `public` prompts only. |
| `ROUTER_BACKEND` | `mock` | `mock` or `vllm`. |
| `ROUTER_VLLM_BASE_URL` | `http://127.0.0.1:8001` | vLLM OpenAI-compatible endpoint. |
| `ROUTER_BACKEND_TIMEOUT_SECONDS` | `60` | Per-request engine timeout. |
Expand Down
12 changes: 12 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ requires-python = ">=3.11"
dependencies = [
"fastapi>=0.141.1,<1",
"httpx>=0.28.1,<1",
"opentelemetry-api>=1.30,<2",
"prometheus-client>=0.26.0,<1",
"pyyaml>=6.0.3,<7",
"pydantic-settings>=2.15.0,<3",
Expand All @@ -21,8 +22,13 @@ dependencies = [
redis = [
"redis>=8.1.0,<9",
]
tracing = [
"opentelemetry-exporter-otlp-proto-http>=1.30,<2",
"opentelemetry-sdk>=1.30,<2",
]
dev = [
"mypy>=2.3.1,<3",
"opentelemetry-sdk>=1.30,<2",
"pytest>=9.1.1,<10",
"pytest-asyncio>=1.4.0,<2",
"pytest-cov>=7.1.0,<8",
Expand Down Expand Up @@ -68,3 +74,9 @@ packages = ["llm_router"]
module = ["redis.*"]
ignore_missing_imports = true
follow_imports = "skip"

[[tool.mypy.overrides]]
# The OTLP exporter ships in the optional tracing extra and is imported lazily,
# only when a collector endpoint is configured.
module = ["opentelemetry.exporter.*"]
ignore_missing_imports = true
112 changes: 83 additions & 29 deletions src/llm_router/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
import httpx
from fastapi import Depends, FastAPI, Header, HTTPException, Request, Response, status
from fastapi.responses import JSONResponse, StreamingResponse
from opentelemetry.trace import TracerProvider

from llm_router.admission import (
AdmissionController,
Expand Down Expand Up @@ -60,6 +61,7 @@
strictest_privacy,
)
from llm_router.routing import NoEligibleModelError, Router, default_model_profiles
from llm_router.tracing import RequestSpan, Tracing, build_tracer_provider


@dataclass(frozen=True)
Expand Down Expand Up @@ -102,6 +104,7 @@ def create_app(
registry: Registry | None = None,
redis_client: RedisLike | None = None,
engine_stats: EngineStatsCollector | None = None,
tracer_provider: TracerProvider | 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 @@ -144,6 +147,10 @@ def create_app(
else None
)
cold_start = ColdStartTracker()
tracing = Tracing(
tracer_provider or build_tracer_provider(runtime_settings),
record_prompt_content=runtime_settings.trace_prompt_content,
)
# 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
Expand Down Expand Up @@ -479,54 +486,85 @@ async def _stream_completion(
decision: RouteDecision,
started: float,
queue_seconds: float,
span: RequestSpan,
) -> AsyncIterator[str]:
completion_id = f"chatcmpl-{uuid.uuid4().hex}"
created = int(time.time())
model_id = decision.profile.id
collected: list[str] = []

yield _chunk(completion_id, created, model_id, delta={"role": "assistant"})
async for delta in inference_backend.stream(payload, decision):
collected.append(delta)
yield _chunk(completion_id, created, model_id, delta={"content": delta})
yield _chunk(completion_id, created, model_id, delta={}, finish_reason="stop")
yield "data: [DONE]\n\n"
# The stream outlives the handler, so it owns the span from here and
# ends it whether generation finishes, fails, or the client leaves.
try:
yield _chunk(completion_id, created, model_id, delta={"role": "assistant"})
async for delta in inference_backend.stream(payload, decision):
collected.append(delta)
yield _chunk(completion_id, created, model_id, delta={"content": delta})
yield _chunk(completion_id, created, model_id, delta={}, finish_reason="stop")
yield "data: [DONE]\n\n"

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,
queue_seconds=queue_seconds,
prompt_tokens=prompt_tokens,
completion_tokens=completion_tokens,
)
await _store_cache(
payload,
prompt,
cache_key,
subject,
decision,
BackendResult(
text=text, prompt_tokens=prompt_tokens, completion_tokens=completion_tokens
),
)
text = "".join(collected)
prompt_tokens = max(1, len(prompt) // 4)
completion_tokens = max(1, len(text) // 4)
span.set_usage(prompt_tokens=prompt_tokens, completion_tokens=completion_tokens)
if payload.routing.structured:
telemetry.record_structured_output(decision, valid=structured_output_valid(text))
telemetry.record_completion(
decision,
latency_seconds=time.perf_counter() - started,
queue_seconds=queue_seconds,
prompt_tokens=prompt_tokens,
completion_tokens=completion_tokens,
)
await _store_cache(
payload,
prompt,
cache_key,
subject,
decision,
BackendResult(
text=text, prompt_tokens=prompt_tokens, completion_tokens=completion_tokens
),
)
except Exception as error:
span.fail(error)
raise
finally:
span.end()

@app.post("/v1/chat/completions", response_model=None)
async def chat_completions(
payload: ChatCompletionRequest,
response: Response,
principal: Annotated[Principal, Depends(authenticate)],
) -> ChatCompletionResponse | StreamingResponse:
span = tracing.start_request()
try:
result = await _complete(payload, response, principal, span)
except Exception as error:
span.fail(error)
raise
finally:
if not span.handed_off:
span.end()
trace_id = span.trace_id
if trace_id is not None:
# A response returned directly does not inherit injected headers.
target = result if isinstance(result, Response) else response
target.headers["X-Trace-Id"] = trace_id
return result

async def _complete(
payload: ChatCompletionRequest,
response: Response,
principal: Principal,
span: RequestSpan,
) -> ChatCompletionResponse | StreamingResponse:
started = time.perf_counter()
tenant = _entitlement(principal)
# Quota and cache are scoped to the tenant, not the credential, so
# rotating a key neither resets a quota nor orphans a cache.
subject = principal.tenant_id
await _consume_quota(subject, tenant)

declared_privacy = payload.routing.privacy
if tenant is not None:
Expand All @@ -543,15 +581,25 @@ async def chat_completions(
privacy_raised_from = (
declared_privacy if payload.routing.privacy is not declared_privacy else None
)
# Described before the quota is charged so a rejected request is still
# attributed to its tenant, and after the floor so the trace is redacted
# under the effective class rather than the declared one.
span.set_request(payload, tenant_id=subject, declared_privacy=declared_privacy)
await _consume_quota(subject, tenant)

prompt = payload.prompt
cache_key = build_cache_key(
payload, tenant=subject, model_revision=catalog_version, prompt=prompt
)

cached = await _lookup_cache(payload, prompt, cache_key, subject)
span.set_cache("miss" if cached is None else cached[0])
if cached is not None:
hit_name, entry = cached
span.set_served_from_cache(model_id=entry.model_id, model_revision=entry.model_revision)
span.set_usage(
prompt_tokens=entry.prompt_tokens, completion_tokens=entry.completion_tokens
)
if payload.stream:
return _cached_stream(entry, hit_name)
response.headers["X-Cache"] = hit_name
Expand Down Expand Up @@ -583,6 +631,7 @@ async def chat_completions(
),
privacy_raised_from=privacy_raised_from,
)
span.set_route(decision)
if runtime_settings.cache_enabled:
decision_cache.set(prompt, decision.task)
telemetry.record_route(decision, privacy=payload.routing.privacy.value)
Expand All @@ -601,6 +650,7 @@ async def chat_completions(
telemetry.inflight_requests.inc()
try:
if payload.stream:
span.handed_off = True
return StreamingResponse(
_stream_completion(
payload,
Expand All @@ -610,6 +660,7 @@ async def chat_completions(
decision,
started,
queue_seconds,
span,
),
media_type="text/event-stream",
headers=_route_headers(decision, "miss"),
Expand All @@ -621,6 +672,9 @@ async def chat_completions(
telemetry.queued_requests.dec()
raise

span.set_usage(
prompt_tokens=result.prompt_tokens, completion_tokens=result.completion_tokens
)
if payload.routing.structured:
telemetry.record_structured_output(decision, valid=structured_output_valid(result.text))
telemetry.record_completion(
Expand Down
5 changes: 5 additions & 0 deletions src/llm_router/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ class Settings(BaseSettings):
registry_path: str = "config/registry.yaml"
routing_policy_version: str = "v1"
redis_url: str = ""
# Traces are exported only when a collector endpoint is set. Prompt content
# is never recorded unless an operator opts in, and then only for public
# requests; restricted and private prompts stay out of traces regardless.
otlp_endpoint: str = ""
trace_prompt_content: bool = False
cache_enabled: bool = True
cache_ttl_seconds: float = Field(default=300.0, gt=0)
cache_max_entries: int = Field(default=1024, ge=1)
Expand Down
Loading
Loading