"""Cheap signals that decide whether a window of scrapes looks like an incident.
``evaluate_signal`` is a pure function over one window of ``/metrics`` scrapes
the caller chose. It returns the kind's value, whether the window had enough
data to decide, and the verdict against the shared threshold table. It never
diagnoses: a value over its threshold says a mechanism is *suspected* in an
engine's aggregate metrics, which cover every client's traffic. Confirming
it, and saying whose requests it hurt, is the diagnoser's job.
Only kinds that ``/metrics`` alone can decide are evaluated; every other kind
answers ``requires_hook``, ``requires_trace`` or ``requires_client``.
"""
from __future__ import annotations
import math
from collections.abc import Mapping, Sequence
from dataclasses import dataclass, field, replace
from types import MappingProxyType
from typing import Any
from . import diagnosis_vocabulary as kinds
from .diagnosis_thresholds import (
DEFAULT_THRESHOLDS,
KV_PREEMPTIONS,
PREFIX_HIT_RATIO_DROP,
PREFIX_MIN_QUERIED,
QUEUE_MEDIAN_WAITING,
THRESHOLDS_VERSION,
resolve_threshold,
)
from .scrape_window import (
REASON_NO_OBSERVATIONS,
WindowCheck,
check_window,
counter_window,
gauge_median,
gauge_window,
histogram_quantile_bounds,
)
from .vllm_telemetry import VllmScrapeRecord
WAITING = "vllm:num_requests_waiting"
WAITING_BY_REASON = "vllm:num_requests_waiting_by_reason"
QUEUE_TIME = "vllm:request_queue_time_seconds"
PREEMPTIONS = "vllm:num_preemptions_total"
KV_USAGE = "vllm:kv_cache_usage_perc"
PREFIX_HITS = "vllm:prefix_cache_hits_total"
PREFIX_QUERIES = "vllm:prefix_cache_queries_total"
REASON_REQUIRES_HOOK = "requires_hook"
REASON_REQUIRES_TRACE = "requires_trace"
REASON_REQUIRES_CLIENT = "requires_client"
REASON_REQUIRES_REFERENCE = "requires_reference"
REASON_TOO_FEW_QUERIES = "too_few_queried_tokens"
# Metrics describe the engine's whole traffic; a value is never one client's.
SCOPE = "engine_global"
_REQUIREMENTS: Mapping[str, str] = MappingProxyType(
{
kinds.MIXED_PREFILL_INTERFERENCE: REASON_REQUIRES_HOOK,
kinds.HOST_STALL: REASON_REQUIRES_HOOK,
kinds.CAPTURE_PAUSE: REASON_REQUIRES_HOOK,
kinds.RANK_DELAY: REASON_REQUIRES_TRACE,
kinds.TRANSFER_DEGRADATION: REASON_REQUIRES_TRACE,
kinds.CLIENT_ADMISSION: REASON_REQUIRES_CLIENT,
**{kind: REASON_REQUIRES_CLIENT for kind in kinds.WORKLOAD_KINDS},
}
)
[docs]
@dataclass(frozen=True)
class SignalConfig:
"""How to evaluate one window.
``thresholds`` overrides entries of the shared table by key; a result
says when it did. A key the table lacks, or a value that is not a finite
number (a NaN, a string, a bool), is refused: a NaN would never be
exceeded and a misspelt key never read, both silently. ``min_scrapes``
must be at least 2. ``reference`` is the prefix-cache hit ratio a
window is compared with, which only the caller can know.
"""
engine: str | None = None
min_scrapes: int = 2
thresholds: Mapping[str, float] = field(default_factory=dict)
reference: float | None = None
def __post_init__(self) -> None:
if self.reference is not None and not 0.0 <= self.reference <= 1.0:
raise ValueError("reference must be a hit ratio between 0 and 1")
if isinstance(self.min_scrapes, bool) or not isinstance(self.min_scrapes, int):
raise ValueError("min_scrapes must be an integer")
if self.min_scrapes < 2:
# A window of one scrape has no length: a gauge would decide on
# one sample and a counter could not decide at all.
raise ValueError("min_scrapes must be at least 2")
unknown = sorted(set(self.thresholds) - set(DEFAULT_THRESHOLDS))
if unknown:
raise ValueError(f"unknown threshold keys: {', '.join(unknown)}")
if not all(_finite_number(value) for value in self.thresholds.values()):
raise ValueError("threshold overrides must be finite numbers")
def _finite_number(value: Any) -> bool:
"""A real, finite number: not a bool, a string or a NaN."""
return (
isinstance(value, (int, float))
and not isinstance(value, bool)
and math.isfinite(value)
)
[docs]
@dataclass(frozen=True)
class SignalValue:
"""One kind's value over a window, and the verdict against its threshold.
``exceeds`` is None whenever ``sufficient`` is False; ``reason`` is then
the first of the reasons listed in ``detail["reasons"]``.
"""
value: float | None
sufficient: bool
reason: str | None
exceeds: bool | None
threshold: float | None
thresholds_version: str
threshold_overridden: bool
detail: Mapping[str, Any]
[docs]
def evaluate_signal(
kind: str,
scrapes: Sequence[VllmScrapeRecord],
config: SignalConfig | None = None,
) -> SignalValue:
"""Evaluate ``kind`` over a window of consecutive scrapes, in time order.
Raises:
ValueError: for a kind outside the closed vocabulary.
"""
kinds.check_kind(kind)
config = config or SignalConfig()
evaluate = _EVALUATORS.get(kind)
if evaluate is None:
return _insufficient(None, None, False, (_REQUIREMENTS[kind],), {})
check = check_window(scrapes, engine=config.engine, min_scrapes=config.min_scrapes)
return evaluate(scrapes, config, check)
def _queue(
scrapes: Sequence[VllmScrapeRecord], config: SignalConfig, check: WindowCheck
) -> SignalValue:
threshold, overridden = resolve_threshold(QUEUE_MEDIAN_WAITING, config.thresholds)
gauge = gauge_window(scrapes, WAITING, engine=config.engine)
value = gauge_median(gauge)
reasons = (*check.reasons, *gauge.reasons)
detail = {
**_window_detail(check),
"max": gauge.max,
"samples": gauge.n,
"waiting_by_reason": _waiting_by_reason(scrapes, config.engine),
"queue_time_p90_s": _queue_time_p90(scrapes, config.engine),
}
if reasons or value is None:
return _insufficient(value, threshold, overridden, reasons, detail)
return _decided(value, threshold, overridden, detail)
def _waiting_by_reason(
scrapes: Sequence[VllmScrapeRecord], engine: str | None
) -> dict[str, float | None]:
"""The median of each waiting reason; a reason the server lacks is None."""
return {
reason: gauge_median(
gauge_window(
scrapes, WAITING_BY_REASON, labels={"reason": reason}, engine=engine
)
)
for reason in ("capacity", "deferred")
}
def _queue_time_p90(
scrapes: Sequence[VllmScrapeRecord], engine: str | None
) -> list[float | None] | None:
bounds = histogram_quantile_bounds(scrapes, QUEUE_TIME, 0.9, engine=engine)
if bounds.count_delta is None or REASON_NO_OBSERVATIONS in bounds.reasons:
return None
return [bounds.lo, bounds.hi]
def _kv(
scrapes: Sequence[VllmScrapeRecord], config: SignalConfig, check: WindowCheck
) -> SignalValue:
threshold, overridden = resolve_threshold(KV_PREEMPTIONS, config.thresholds)
counter = counter_window(scrapes, PREEMPTIONS, engine=config.engine)
usage = gauge_window(scrapes, KV_USAGE, engine=config.engine)
detail = {
**_window_detail(check),
"rate_per_s": counter.rate_per_s,
"kv_cache_usage_max": usage.max,
}
reasons = (*check.reasons, *counter.reasons)
if reasons or counter.delta is None:
return _insufficient(counter.delta, threshold, overridden, reasons, detail)
return _decided(counter.delta, threshold, overridden, detail)
def _prefix(
scrapes: Sequence[VllmScrapeRecord], config: SignalConfig, check: WindowCheck
) -> SignalValue:
threshold, overridden = resolve_threshold(PREFIX_HIT_RATIO_DROP, config.thresholds)
least, _ = resolve_threshold(PREFIX_MIN_QUERIED, config.thresholds)
hits = counter_window(scrapes, PREFIX_HITS, engine=config.engine)
queries = counter_window(scrapes, PREFIX_QUERIES, engine=config.engine)
reasons = [*check.reasons, *hits.reasons, *queries.reasons]
ratio = _hit_ratio(hits.delta, queries.delta) if not reasons else None
if not reasons:
reasons.extend(_volume_reasons(ratio, queries.delta, least))
detail = {
**_window_detail(check),
"hits": hits.delta,
"queries": queries.delta,
"reference": config.reference,
}
reference = config.reference
if reference is None or reasons or ratio is None:
if not reasons:
reasons.append(REASON_REQUIRES_REFERENCE)
return _insufficient(ratio, threshold, overridden, reasons, detail)
# The value is the ratio; the verdict is on its fall below the reference.
detail["drop"] = reference - ratio
return replace(
_decided(reference - ratio, threshold, overridden, detail), value=ratio
)
def _volume_reasons(
ratio: float | None, queried: float | None, least: float
) -> list[str]:
"""Why a window's hit ratio decides nothing: no tokens queried, or too few."""
if ratio is None:
return [REASON_NO_OBSERVATIONS]
return [REASON_TOO_FEW_QUERIES] if (queried or 0.0) < least else []
def _hit_ratio(hits: float | None, queries: float | None) -> float | None:
"""Hits per queried token, or None when nothing was queried."""
if hits is None or not queries:
return None
return hits / queries
def _window_detail(check: WindowCheck) -> dict[str, Any]:
return {
"scope": SCOPE,
"scrapes": check.scrapes,
# Failed scrapes inside the window: a caller need not count again.
"failed_scrapes": check.failed,
"window_seconds": check.seconds,
"placement": check.placement,
}
def _decided(
value: float, threshold: float, overridden: bool, detail: dict[str, Any]
) -> SignalValue:
detail["reasons"] = []
return SignalValue(
value,
True,
None,
value >= threshold,
threshold,
THRESHOLDS_VERSION,
overridden,
MappingProxyType(detail),
)
def _insufficient(
value: float | None,
threshold: float | None,
overridden: bool,
reasons: Sequence[str],
detail: dict[str, Any],
) -> SignalValue:
unique = list(dict.fromkeys(reasons))
detail["reasons"] = unique
detail.setdefault("scope", SCOPE)
return SignalValue(
value,
False,
unique[0] if unique else None,
None,
threshold,
THRESHOLDS_VERSION,
overridden,
MappingProxyType(detail),
)
_EVALUATORS = MappingProxyType(
{
kinds.QUEUE_SATURATION: _queue,
kinds.KV_PREEMPTION_PRESSURE: _kv,
kinds.PREFIX_CACHE_LOSS: _prefix,
}
)
__all__ = [
"REASON_REQUIRES_CLIENT",
"REASON_REQUIRES_HOOK",
"REASON_REQUIRES_REFERENCE",
"REASON_REQUIRES_TRACE",
"REASON_TOO_FEW_QUERIES",
"SignalConfig",
"SignalValue",
"evaluate_signal",
]