"""Scrape vLLM's ``/metrics`` during a profile and keep each response whole.
The profiler scrapes at the start and end of every phase and on a cadence
inside it. Each scrape becomes one ``infer.vllm_scrape`` record on the
client's clock, so no clock alignment is needed to place it against the
phase windows. A scrape that fails is recorded as a failure with its reason,
never as an empty set of series. Deltas and rates are computed at analysis
time, where resets and restarts can be recognised.
"""
from __future__ import annotations
import asyncio
import hashlib
import http.client
import threading
import time
import urllib.error
import urllib.parse
import urllib.request
from collections.abc import Callable
from dataclasses import dataclass
from functools import partial
from typing import Any, TypeVar
from ..scrub import redact_url
from .correlation_events import CapabilityEvent, CorrelationContext
from .vllm_metrics import (
CATALOG,
VERIFIED_VLLM_VERSION,
CompactScrape,
Discovery,
ScrapeTooLarge,
compact_scrape,
discover,
parse_prometheus_text,
)
from .vllm_telemetry import (
MARKER_INTERVAL,
SCRAPE_ERROR,
SCRAPE_OK,
VllmScrapeRecord,
)
DEFAULT_METRICS_PATH = "/metrics"
AUTO_METRICS_URL = "auto"
CAPABILITY_COMPONENT = "vllm.metrics"
# The phase-end scrape of a run being stopped waits at most this long, so a
# server that stopped answering cannot hold Ctrl+C back.
INTERRUPT_SCRAPE_TIMEOUT_SECONDS = 2.0
# A response is read at most this far, and parsed into records only below this
# many series, so a misbehaving endpoint cannot grow the client's memory. A
# vLLM 0.30.0 response for one model is about 90 KB and 360 series.
MAX_SCRAPE_BYTES = 8 * 1024 * 1024
MAX_SCRAPE_SERIES = 20_000
[docs]
def resolve_metrics_url(endpoint: str, requested: str | None) -> str | None:
"""None when scraping is off; ``auto`` is the endpoint's origin plus ``/metrics``."""
if requested is None:
return None
if requested != AUTO_METRICS_URL:
return requested
parts = urllib.parse.urlsplit(endpoint)
if parts.scheme not in {"http", "https"} or not parts.netloc:
raise ValueError("--vllm-metrics needs a URL when the endpoint is not http(s)")
return f"{parts.scheme}://{parts.netloc}{DEFAULT_METRICS_PATH}"
[docs]
def metrics_api_key(
endpoint: str,
url: str,
api_key: str | None,
on_warning: Callable[[str], None] | None = None,
) -> str | None:
"""The bearer token for scrapes: the endpoint's, and only on its origin.
A metrics URL on another scheme, host or port is scraped without
credentials, so the token given for the inference endpoint never goes
anywhere else; the run says so once.
"""
if api_key is None or _origin(endpoint) == _origin(url) is not None:
return api_key
if on_warning is not None:
on_warning(
f"the bearer token is not sent to {redact_url(url) or url}: it is not "
"the endpoint's origin, so scrapes go out without credentials"
)
return None
_DEFAULT_PORTS = {"http": 80, "https": 443}
def _origin(url: str) -> tuple[str, str, int | None] | None:
"""Lowercase scheme and host plus the effective port, so ``https://h`` and
``https://h:443`` are one origin; None for a URL without a usable port."""
parts = urllib.parse.urlsplit(url)
scheme = parts.scheme.lower()
try:
port = parts.port
except ValueError:
return None
if port is None:
port = _DEFAULT_PORTS.get(scheme)
return (scheme, (parts.hostname or "").lower(), port)
[docs]
@dataclass(frozen=True)
class FetchResult:
"""What one GET of the metrics URL returned, or why it did not."""
text: str | None
http_status: int | None
error: str | None
duration_ms: float
class _NoRedirect(urllib.request.HTTPRedirectHandler):
"""Refuse every redirect, so the 3xx is answered as an error.
urllib would follow it and forward the Authorization header to wherever
it points, so a /metrics on the endpoint's origin that redirects would
hand the endpoint's token to another origin. A scrape is of one URL:
the redirect is recorded with its target and never followed.
"""
def redirect_request(
self,
req: urllib.request.Request,
fp: Any,
code: int,
msg: str,
headers: Any,
newurl: str,
) -> None:
return None
_OPENER = urllib.request.build_opener(_NoRedirect())
[docs]
def fetch_metrics(
url: str,
*,
timeout_seconds: float,
api_key: str | None = None,
max_bytes: int = MAX_SCRAPE_BYTES,
) -> FetchResult:
"""GET the metrics text; every failure becomes a result, never an exception.
The body is read at most ``max_bytes + 1`` bytes far: a longer response
is a failed scrape, and the rest of it is never read.
"""
headers = {"Accept": "text/plain"}
if api_key:
headers["Authorization"] = f"Bearer {api_key}"
request = urllib.request.Request(url, headers=headers, method="GET")
started = time.perf_counter()
try:
with _OPENER.open(request, timeout=timeout_seconds) as response:
status = int(response.status)
try:
raw = _read_capped(response, max_bytes)
except http.client.HTTPException as exc:
# A body cut short of its Content-Length (IncompleteRead) is
# an http.client error, not an OSError: still a failed
# scrape, never the run's end.
elapsed = (time.perf_counter() - started) * 1000.0
return FetchResult(
None, status, f"{type(exc).__name__}: {exc}", elapsed
)
elapsed = (time.perf_counter() - started) * 1000.0
if len(raw) > max_bytes:
error = f"oversized: the response is over {max_bytes} bytes"
return FetchResult(None, status, error, elapsed)
body = raw.decode("utf-8", errors="replace")
return FetchResult(body, status, None, elapsed)
except urllib.error.HTTPError as exc:
elapsed = (time.perf_counter() - started) * 1000.0
return FetchResult(None, exc.code, _http_error_text(exc), elapsed)
except (urllib.error.URLError, OSError, ValueError) as exc:
elapsed = (time.perf_counter() - started) * 1000.0
return FetchResult(None, None, f"{type(exc).__name__}: {exc}", elapsed)
def _read_capped(response: Any, max_bytes: int) -> bytes:
"""At most ``max_bytes + 1`` bytes of the body.
A bounded read returns short at EOF instead of raising, so a body cut
short of its Content-Length is raised here as ``IncompleteRead``, as an
unbounded read would.
"""
raw = bytes(response.read(max_bytes + 1))
remaining = getattr(response, "length", None)
if len(raw) <= max_bytes and remaining:
raise http.client.IncompleteRead(raw, remaining)
return raw
def _http_error_text(exc: urllib.error.HTTPError) -> str:
"""``HTTP 503``, or for a refused redirect also where it pointed."""
location = exc.headers.get("Location") if exc.headers is not None else None
if 300 <= exc.code < 400 and location:
return f"HTTP {exc.code}: redirect to {redact_url(location) or location} not followed"
return f"HTTP {exc.code}"
def _fetch_and_parse(
url: str,
timeout_seconds: float,
api_key: str | None,
max_bytes: int = MAX_SCRAPE_BYTES,
max_series: int = MAX_SCRAPE_SERIES,
) -> tuple[FetchResult, CompactScrape | None]:
"""Fetch and parse, touching no scraper state, so the work can run on a
thread whose result may be dropped."""
result = fetch_metrics(
url, timeout_seconds=timeout_seconds, api_key=api_key, max_bytes=max_bytes
)
if result.text is None:
return result, None
try:
families = parse_prometheus_text(result.text, max_series=max_series)
return result, compact_scrape(families)
except ScrapeTooLarge as exc:
return _failed_fetch(result, f"oversized: {exc}"), None
except ValueError as exc:
return _failed_fetch(result, f"unparseable response: {exc}"), None
def _failed_fetch(result: FetchResult, error: str) -> FetchResult:
return FetchResult(None, result.http_status, error, result.duration_ms)
T = TypeVar("T")
async def _off_loop(func: Callable[[], T]) -> T:
"""Run ``func`` on a daemon thread; cancelling the await abandons it.
``asyncio.to_thread`` would make a cancelled caller, and then the
interpreter's exit, wait for a fetch blocked on a silent endpoint until
its socket timeout. A daemon thread is dropped instead, and its result
is discarded because nothing awaits it any more.
"""
loop = asyncio.get_running_loop()
future: asyncio.Future[T] = loop.create_future()
def deliver(outcome: Callable[[], None]) -> None:
if not future.done():
outcome()
def run() -> None:
try:
value = func()
except BaseException as exc:
outcome = partial(future.set_exception, exc)
else:
outcome = partial(future.set_result, value)
try:
loop.call_soon_threadsafe(deliver, outcome)
except RuntimeError:
pass # the loop is closed: nobody is waiting
threading.Thread(target=run, name="stormlog-vllm-scrape", daemon=True).start()
return await future
[docs]
class VllmMetricsScraper:
"""Turn metrics responses into records and remember what they exposed."""
def __init__(
self,
*,
url: str,
interval_seconds: float,
timeout_seconds: float,
session_id: str,
run_id: str,
clock_domain: str,
api_key: str | None = None,
on_warning: Callable[[str], None] | None = None,
max_scrape_bytes: int = MAX_SCRAPE_BYTES,
max_scrape_series: int = MAX_SCRAPE_SERIES,
) -> None:
self.url = url
self.source_url = redact_url(url) or url
self.max_scrape_bytes = max_scrape_bytes
self.max_scrape_series = max_scrape_series
self.interval_seconds = interval_seconds
self.interval_ms = max(1, round(interval_seconds * 1000))
self.timeout_seconds = timeout_seconds
self.session_id = session_id
self.run_id = run_id
self.clock_domain = clock_domain
self.api_key = api_key
self.on_warning = on_warning
self.ok_scrapes = 0
self.failed_scrapes = 0
self.first_discovery: Discovery | None = None
self.last_error: str | None = None
# Fetches still running on their threads, abandoned ones included.
self._fetching = 0
self._fetching_lock = threading.Lock()
[docs]
def scrape(
self,
*,
marker: str,
case_id: str | None = None,
phase: str | None = None,
timeout_seconds: float | None = None,
) -> VllmScrapeRecord:
"""Fetch and parse once, here; the record says what happened either way.
``timeout_seconds`` overrides the scraper's own for one scrape, for
the one taken on the way out of an interrupted run.
"""
observed_at_ns = time.time_ns()
result, compact = self._fetch(timeout_seconds)()
return self._record(observed_at_ns, marker, case_id, phase, result, compact)
[docs]
async def scrape_async(
self,
*,
marker: str,
case_id: str | None = None,
phase: str | None = None,
timeout_seconds: float | None = None,
) -> VllmScrapeRecord:
"""``scrape`` with the fetch on a thread a cancelled caller abandons.
Counting and the record happen back on the loop, so a fetch dropped
by cancellation leaves no trace: no counter moves and no record is
written for it.
"""
observed_at_ns = time.time_ns()
fetch = self._fetch(timeout_seconds)
def tracked() -> tuple[FetchResult, CompactScrape | None]:
try:
return fetch()
finally:
with self._fetching_lock:
self._fetching -= 1
with self._fetching_lock:
self._fetching += 1
result, compact = await _off_loop(tracked)
return self._record(observed_at_ns, marker, case_id, phase, result, compact)
[docs]
def fetching(self) -> bool:
"""A fetch is still running on its thread, one whose caller gave up
waiting for it included: a socket timeout bounds each read, not the
whole response, so a trickling one can run on for long."""
with self._fetching_lock:
return self._fetching > 0
[docs]
def abandoned(
self,
*,
marker: str,
case_id: str | None,
phase: str | None,
observed_at_ns: int,
deadline_seconds: float,
reason: str | None = None,
) -> VllmScrapeRecord:
"""The record of a scrape given up at an overall deadline.
The fetch itself may still be reading on its thread; its result is
dropped, so this failed record is the only trace of the scrape.
"""
error = reason or (
f"the {deadline_seconds:g} s deadline of the interrupted run passed "
"while the response was still arriving"
)
result = FetchResult(
None, None, f"abandoned: {error}", deadline_seconds * 1000.0
)
return self._failed(observed_at_ns, marker, case_id, phase, result)
def _timeout(self, override: float | None) -> float:
return self.timeout_seconds if override is None else override
def _fetch(
self, timeout_seconds: float | None
) -> Callable[[], tuple[FetchResult, CompactScrape | None]]:
"""One bounded fetch-and-parse, as a call that touches no scraper state."""
return partial(
_fetch_and_parse,
self.url,
self._timeout(timeout_seconds),
self.api_key,
self.max_scrape_bytes,
self.max_scrape_series,
)
def _record(
self,
observed_at_ns: int,
marker: str,
case_id: str | None,
phase: str | None,
result: FetchResult,
compact: CompactScrape | None,
) -> VllmScrapeRecord:
if compact is None:
return self._failed(observed_at_ns, marker, case_id, phase, result)
found = discover(compact)
if self.first_discovery is None:
self.first_discovery = found
self.ok_scrapes += 1
encoded = (result.text or "").encode("utf-8")
return VllmScrapeRecord(
session_id=self.session_id,
run_id=self.run_id,
observed_at_ns=observed_at_ns,
source_url=self.source_url,
marker=marker,
interval_ms=self.interval_ms,
status=SCRAPE_OK,
clock_domain=self.clock_domain,
case_id=case_id,
phase=phase,
duration_ms=result.duration_ms,
http_status=result.http_status,
content_digest=hashlib.sha256(encoded).hexdigest(),
content_bytes=len(encoded),
scrape=compact,
discovery=found,
)
def _failed(
self,
observed_at_ns: int,
marker: str,
case_id: str | None,
phase: str | None,
result: FetchResult,
) -> VllmScrapeRecord:
self.failed_scrapes += 1
error = result.error or "metrics request failed"
if self.last_error is None and self.on_warning is not None:
self.on_warning(f"vLLM metrics scrape of {self.source_url} failed: {error}")
self.last_error = error
return VllmScrapeRecord(
session_id=self.session_id,
run_id=self.run_id,
observed_at_ns=observed_at_ns,
source_url=self.source_url,
marker=marker,
interval_ms=self.interval_ms,
status=SCRAPE_ERROR,
clock_domain=self.clock_domain,
case_id=case_id,
phase=phase,
duration_ms=result.duration_ms,
http_status=result.http_status,
error=error,
)
[docs]
async def interval_loop(
self,
*,
append: Callable[[dict[str, Any]], None],
case_id: str,
phase: str,
stop_event: asyncio.Event,
) -> None:
"""Scrape every interval until the phase ends; the fetch runs off the loop.
Cancelling the loop while a fetch is in flight abandons that fetch
(see ``scrape_async``), so a stop never waits on a silent endpoint.
"""
while True:
try:
await asyncio.wait_for(stop_event.wait(), timeout=self.interval_seconds)
except asyncio.TimeoutError:
pass
else:
return
record = await self.scrape_async(
marker=MARKER_INTERVAL, case_id=case_id, phase=phase
)
append(record.to_record())
[docs]
def config_record(self) -> dict[str, Any]:
return {
"url": self.source_url,
"interval_seconds": self.interval_seconds,
"timeout_seconds": self.timeout_seconds,
"authorization": "bearer" if self.api_key else None,
"max_scrape_bytes": self.max_scrape_bytes,
"max_scrape_series": self.max_scrape_series,
}
[docs]
def capability_event(self, context: CorrelationContext) -> CapabilityEvent:
"""What the engine exposed, as a v2 capability record for the artifact.
``supported`` lists the catalog's normalized fields, ``enabled`` the
ones the first successful scrape exposed, and ``collected`` repeats
``enabled`` once at least one scrape succeeded. An endpoint that never
answered is unavailable with empty lists.
"""
found = self.first_discovery
available = found is not None and self.ok_scrapes > 0
enabled = (
sorted(CATALOG[name].field for name in found.present)
if found is not None and available
else []
)
return CapabilityEvent(
context=context,
event_id=f"capability:{CAPABILITY_COMPONENT}",
component=CAPABILITY_COMPONENT,
available=available,
supported=(
sorted(entry.field for entry in CATALOG.values()) if available else []
),
enabled=enabled,
collected=list(enabled),
metadata=self._capability_metadata(found),
)
def _capability_metadata(self, found: Discovery | None) -> dict[str, Any]:
metadata: dict[str, Any] = {
"url": self.source_url,
"verified_vllm_version": VERIFIED_VLLM_VERSION,
"ok_scrapes": self.ok_scrapes,
"failed_scrapes": self.failed_scrapes,
"last_error": self.last_error,
"observation_scope": "engine_aggregate",
"per_request_attribution": "none",
}
if found is not None:
metadata.update(
{
"unknown_series": list(found.unknown),
"deprecated_series": list(found.deprecated_present),
"removed_series": list(found.removed_present),
"optional_absent": list(found.optional_absent),
"optional_present": list(found.optional_present),
"engines": list(found.engines),
}
)
return metadata