"""Records for vLLM native telemetry inside an inference artifact.
A scrape record holds one ``/metrics`` response in the compact form from
:mod:`stormlog.infer.vllm_metrics`, stamped on the client's clock. A span
record holds one OpenTelemetry span exactly as vLLM exported it. Both are
appended to the client artifact next to the request events, and both are
aggregate or engine-side evidence: nothing in them says which request used
which GPU time.
"""
from __future__ import annotations
import math
from collections.abc import Iterable
from dataclasses import dataclass, field
from typing import Any
from .vllm_metrics import CompactScrape, Discovery
VLLM_SCHEMA_VERSION = 1
SCRAPE_EVENT_TYPE = "infer.vllm_scrape"
SPAN_EVENT_TYPE = "infer.vllm_span"
OBSERVATION_SCOPE = "engine_aggregate"
MARKER_PHASE_START = "phase_start"
MARKER_INTERVAL = "interval"
MARKER_PHASE_END = "phase_end"
MARKER_MANUAL = "manual"
MARKERS = (MARKER_PHASE_START, MARKER_INTERVAL, MARKER_PHASE_END, MARKER_MANUAL)
SCRAPE_OK = "ok"
SCRAPE_ERROR = "error"
SCRAPE_STATES = (SCRAPE_OK, SCRAPE_ERROR)
SPAN_SOURCE_RECEIVER = "otlp_http_receiver"
SPAN_SOURCE_OTLP_JSON = "otlp_json_file"
SPAN_SOURCE_JSONL = "jsonl_file"
SPAN_SOURCES = (SPAN_SOURCE_RECEIVER, SPAN_SOURCE_OTLP_JSON, SPAN_SOURCE_JSONL)
def _nonempty(value: object, name: str) -> None:
if not isinstance(value, str) or not value:
raise ValueError(f"{name} must be a non-empty string")
def _optional_nonempty(value: object, name: str) -> None:
if value is not None:
_nonempty(value, name)
def _is_int(value: object) -> bool:
return isinstance(value, int) and not isinstance(value, bool)
def _positive_int(value: object, name: str) -> None:
if not _is_int(value) or not isinstance(value, int) or value <= 0:
raise ValueError(f"{name} must be a positive integer")
def _optional_nonnegative_int(value: object, name: str) -> None:
if value is None:
return
if not isinstance(value, int) or isinstance(value, bool) or value < 0:
raise ValueError(f"{name} must be a non-negative integer or null")
def _string_dict(value: object, name: str) -> None:
if not isinstance(value, dict) or any(not isinstance(k, str) for k in value):
raise ValueError(f"{name} must be an object with string keys")
[docs]
@dataclass(frozen=True)
class VllmScrapeRecord:
"""One ``/metrics`` response, kept whole, stamped on the client clock."""
session_id: str
run_id: str
observed_at_ns: int
source_url: str
marker: str
interval_ms: int
status: str
clock_domain: str
case_id: str | None = None
phase: str | None = None
duration_ms: float | None = None
http_status: int | None = None
error: str | None = None
content_digest: str | None = None
content_bytes: int | None = None
scrape: CompactScrape | None = None
discovery: Discovery | None = None
def __post_init__(self) -> None:
for name in ("session_id", "run_id", "source_url", "clock_domain"):
_nonempty(getattr(self, name), name)
_positive_int(self.observed_at_ns, "observed_at_ns")
_positive_int(self.interval_ms, "interval_ms")
if self.marker not in MARKERS:
raise ValueError("unsupported scrape marker")
if self.status not in SCRAPE_STATES:
raise ValueError("unsupported scrape status")
_optional_nonempty(self.case_id, "case_id")
_optional_nonempty(self.phase, "phase")
_optional_nonempty(self.error, "error")
_optional_nonempty(self.content_digest, "content_digest")
_optional_nonnegative_int(self.http_status, "http_status")
_optional_nonnegative_int(self.content_bytes, "content_bytes")
self._validate_duration()
self._validate_outcome()
def _validate_duration(self) -> None:
duration = self.duration_ms
if duration is None:
return
if not isinstance(duration, (int, float)) or isinstance(duration, bool):
raise ValueError("duration_ms must be a non-negative number or null")
if not math.isfinite(duration) or duration < 0:
raise ValueError("duration_ms must be a non-negative number or null")
def _validate_outcome(self) -> None:
if self.status == SCRAPE_OK:
if self.scrape is None or self.discovery is None:
raise ValueError("an ok scrape carries its series and discovery")
if self.error is not None:
raise ValueError("an ok scrape has no error")
elif self.scrape is not None or self.discovery is not None:
raise ValueError("a failed scrape carries no series")
elif self.error is None:
raise ValueError("a failed scrape says why")
[docs]
def to_record(self) -> dict[str, Any]:
return {
"schema_version": VLLM_SCHEMA_VERSION,
"event_type": SCRAPE_EVENT_TYPE,
"observation_scope": OBSERVATION_SCOPE,
"session_id": self.session_id,
"run_id": self.run_id,
"observed_at_ns": self.observed_at_ns,
"timestamp_ns": self.observed_at_ns,
"source_url": self.source_url,
"marker": self.marker,
"interval_ms": self.interval_ms,
"status": self.status,
"clock_domain": self.clock_domain,
"case_id": self.case_id,
"phase": self.phase,
"duration_ms": self.duration_ms,
"http_status": self.http_status,
"error": self.error,
"content_digest": self.content_digest,
"content_bytes": self.content_bytes,
"scrape": self.scrape.to_record() if self.scrape is not None else None,
"discovery": (
self.discovery.to_record() if self.discovery is not None else None
),
}
[docs]
@classmethod
def from_record(cls, record: dict[str, Any]) -> VllmScrapeRecord:
_check_envelope(record, SCRAPE_EVENT_TYPE)
if record.get("observation_scope") != OBSERVATION_SCOPE:
raise ValueError("scrape records are engine-aggregate evidence")
scrape = record.get("scrape")
discovery = record.get("discovery")
return cls(
session_id=record["session_id"],
run_id=record["run_id"],
observed_at_ns=record["observed_at_ns"],
source_url=record["source_url"],
marker=record["marker"],
interval_ms=record["interval_ms"],
status=record["status"],
clock_domain=record["clock_domain"],
case_id=record.get("case_id"),
phase=record.get("phase"),
duration_ms=record.get("duration_ms"),
http_status=record.get("http_status"),
error=record.get("error"),
content_digest=record.get("content_digest"),
content_bytes=record.get("content_bytes"),
scrape=CompactScrape.from_record(scrape) if scrape is not None else None,
discovery=(
_discovery_from_record(discovery) if discovery is not None else None
),
)
def _discovery_from_record(record: dict[str, Any]) -> Discovery:
return Discovery(
present=tuple(record["present"]),
absent=tuple(record["absent"]),
optional_absent=tuple(record["optional_absent"]),
optional_present=tuple(record.get("optional_present", ())),
deprecated_present=tuple(record["deprecated_present"]),
removed_present=tuple(record["removed_present"]),
unknown=tuple(record["unknown"]),
engines=tuple(record["engines"]),
model_names=tuple(record["model_names"]),
process_start_ns=record["process_start_ns"],
)
def _check_envelope(record: dict[str, Any], event_type: str) -> None:
if (
type(record.get("schema_version")) is not int
or record.get("schema_version") != VLLM_SCHEMA_VERSION
or record.get("event_type") != event_type
):
raise ValueError(f"unsupported {event_type} record")
[docs]
@dataclass(frozen=True)
class VllmSpanRecord:
"""One span as exported, with its attributes under their native names.
The timestamps are the exporter's wall clock, named by ``clock_domain``.
``request_id`` is the client-chosen ``X-Request-Id`` recovered from
``gen_ai.request.id`` when the span carries one; the join to a Stormlog
request uses it, never a rebuilt string.
"""
session_id: str
run_id: str
source: str
name: str
clock_domain: str
received_at_ns: int | None = None
trace_id: str | None = None
span_id: str | None = None
parent_span_id: str | None = None
kind: str | None = None
start_unix_ns: int | None = None
end_unix_ns: int | None = None
attributes: dict[str, Any] = field(default_factory=dict)
resource: dict[str, Any] = field(default_factory=dict)
scope: dict[str, Any] = field(default_factory=dict)
status: dict[str, Any] | None = None
dropped: dict[str, int] = field(default_factory=dict)
request_id: str | None = None
def __post_init__(self) -> None:
for name in ("session_id", "run_id", "name", "clock_domain"):
_nonempty(getattr(self, name), name)
if self.source not in SPAN_SOURCES:
raise ValueError("unsupported span source")
for name in ("trace_id", "span_id", "parent_span_id", "kind", "request_id"):
_optional_nonempty(getattr(self, name), name)
for name in ("received_at_ns", "start_unix_ns", "end_unix_ns"):
_optional_nonnegative_int(getattr(self, name), name)
if (
self.start_unix_ns is not None
and self.end_unix_ns is not None
and self.end_unix_ns < self.start_unix_ns
):
raise ValueError("end_unix_ns must be >= start_unix_ns")
for name in ("attributes", "resource", "scope", "dropped"):
_string_dict(getattr(self, name), name)
if self.status is not None:
_string_dict(self.status, "status")
@property
def duration_ns(self) -> int | None:
if self.start_unix_ns is None or self.end_unix_ns is None:
return None
return self.end_unix_ns - self.start_unix_ns
[docs]
def to_record(self) -> dict[str, Any]:
return {
"schema_version": VLLM_SCHEMA_VERSION,
"event_type": SPAN_EVENT_TYPE,
"session_id": self.session_id,
"run_id": self.run_id,
"source": self.source,
"name": self.name,
"clock_domain": self.clock_domain,
"received_at_ns": self.received_at_ns,
"timestamp_ns": self.start_unix_ns,
"trace_id": self.trace_id,
"span_id": self.span_id,
"parent_span_id": self.parent_span_id,
"kind": self.kind,
"start_unix_ns": self.start_unix_ns,
"end_unix_ns": self.end_unix_ns,
"attributes": dict(self.attributes),
"resource": dict(self.resource),
"scope": dict(self.scope),
"status": dict(self.status) if self.status is not None else None,
"dropped": dict(self.dropped),
"request_id": self.request_id,
}
[docs]
@classmethod
def from_record(cls, record: dict[str, Any]) -> VllmSpanRecord:
_check_envelope(record, SPAN_EVENT_TYPE)
return cls(
session_id=record["session_id"],
run_id=record["run_id"],
source=record["source"],
name=record["name"],
clock_domain=record["clock_domain"],
received_at_ns=record.get("received_at_ns"),
trace_id=record.get("trace_id"),
span_id=record.get("span_id"),
parent_span_id=record.get("parent_span_id"),
kind=record.get("kind"),
start_unix_ns=record.get("start_unix_ns"),
end_unix_ns=record.get("end_unix_ns"),
attributes=dict(record.get("attributes") or {}),
resource=dict(record.get("resource") or {}),
scope=dict(record.get("scope") or {}),
status=record.get("status"),
dropped=dict(record.get("dropped") or {}),
request_id=record.get("request_id"),
)
[docs]
def request_id_from_span_id(native_id: object) -> str | None:
"""Recover the ``X-Request-Id`` vLLM embedded in ``gen_ai.request.id``.
vLLM names a chat completion ``chatcmpl-<X-Request-Id>`` and a text
completion ``cmpl-<X-Request-Id>-<index>``, with the per-prompt index only
on the completions path. A trailing ``-<digits>`` is stripped only when it
is one; the request ids Stormlog sends end in ``_<n>``, never ``-<n>``.
Without the client's header vLLM uses a random id, which no Stormlog
request owns.
"""
if not isinstance(native_id, str):
return None
for prefix in ("chatcmpl-", "cmpl-"):
if native_id.startswith(prefix):
body = native_id[len(prefix) :]
head, _sep, tail = body.rpartition("-")
return head if head and tail.isdigit() else body
return None
[docs]
def load_vllm_records(
records: Iterable[dict[str, Any]],
) -> tuple[list[VllmScrapeRecord], list[VllmSpanRecord]]:
"""Parse the vLLM records of an artifact; an invalid one is an error."""
scrapes: list[VllmScrapeRecord] = []
spans: list[VllmSpanRecord] = []
for index, record in enumerate(records):
event_type = record.get("event_type")
try:
if event_type == SCRAPE_EVENT_TYPE:
scrapes.append(VllmScrapeRecord.from_record(record))
elif event_type == SPAN_EVENT_TYPE:
spans.append(VllmSpanRecord.from_record(record))
except (KeyError, TypeError, ValueError) as exc:
raise ValueError(f"invalid {event_type} record {index}: {exc}") from exc
scrapes.sort(key=lambda item: item.observed_at_ns)
return scrapes, spans