stormlog.infer.vllm_spans
Ingest vLLM’s OpenTelemetry request spans.
Spans arrive in three ways: an OTLP/HTTP receiver that stormlog infer
profile runs for the length of a run, an OTLP JSON file written by a
collector’s file exporter, or the one-span-per-line JSONL that a small sink
writes. Every path produces the same infer.vllm_span record with the
native attributes untouched.
OTLP protobuf bodies, which is what vLLM’s exporter sends, are decoded with
the generated classes from opentelemetry-proto (the infer-otlp
extra). Without that package the receiver still runs, accepts OTLP JSON,
and records protobuf as supported but not enabled. The latency attributes a
span carries are phase residency measured on the engine’s clock; the
mapping to the v2 correlation model keeps them as reported durations and
marks the derived stage windows as estimates.
Functions
|
One line of the sink format: name, start/end unix ns, flat attributes. |
|
Decode the JSON form of |
|
Decode an |
|
Inflate a gzip body, or None when its output would exceed |
|
What decoding this OTLP JSON export holds at once, estimated on its bytes before anything is parsed. |
|
The spans and attribute values |
|
|
|
The |
|
An upper bound on what parsing an export holds, from its wire counts. |
What a parse costs with the installed protobuf backend. |
|
|
Read OTLP JSON (one document, or one per line) or sink JSONL; say which. |
|
What each span's record will hold, estimated as the decode budget is ( |
|
The span capability record for a run that asked for a receiver. |
|
The exporter's wall clock, named by its host; never a shared clock. |
|
|
|
|
|
Map request spans onto v2 request and stage records. |
Classes
|
Accept OTLP/HTTP trace exports on a local port while a profile runs. |
|
What a protobuf parse holds, in bytes, for each thing the wire scan counts ( |
|
A decoded span before it is stamped with run identity. |
|
Admission bounds, applied before a body is read, decoded or queued. |
|
Exceptions
The |
|
The bytes are not an OTLP trace export request. |
- class stormlog.infer.vllm_spans.RawSpan(name, trace_id=None, span_id=None, parent_span_id=None, kind=None, start_unix_ns=None, end_unix_ns=None, attributes=<factory>, resource=<factory>, scope=<factory>, status=None, dropped=<factory>)[source]
Bases:
objectA decoded span before it is stamped with run identity.
- Parameters:
name (str)
trace_id (str | None)
span_id (str | None)
parent_span_id (str | None)
kind (str | None)
start_unix_ns (int | None)
end_unix_ns (int | None)
attributes (dict[str, Any])
resource (dict[str, Any])
scope (dict[str, Any])
status (dict[str, Any] | None)
dropped (dict[str, int])
- name: str
- 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]
- resource: dict[str, Any]
- scope: dict[str, Any]
- status: dict[str, Any] | None = None
- dropped: dict[str, int]
- exception stormlog.infer.vllm_spans.ProtobufDecodeError[source]
Bases:
ValueErrorThe bytes are not an OTLP trace export request.
Bases:
RuntimeErrorThe
infer-otlpextra is not installed.
- stormlog.infer.vllm_spans.protobuf_rates()[source]
What a parse costs with the installed protobuf backend.
- Return type:
- stormlog.infer.vllm_spans.protobuf_parse_estimate(counts, content_bytes)[source]
An upper bound on what parsing an export holds, from its wire counts.
- Parameters:
counts (WireCounts)
content_bytes (int)
- Return type:
int
- stormlog.infer.vllm_spans.decode_otlp_protobuf(data)[source]
Decode an
ExportTraceServiceRequestwith the generated classes.- Parameters:
data (bytes)
- Return type:
list[RawSpan]
- stormlog.infer.vllm_spans.parse_otlp_protobuf(data)[source]
The
ExportTraceServiceRequestmessage, before any span is built.- Parameters:
data (bytes | bytearray)
- Return type:
Any
- stormlog.infer.vllm_spans.spans_from_message(message)[source]
- Parameters:
message (Any)
- Return type:
list[RawSpan]
- stormlog.infer.vllm_spans.decode_otlp_json(document)[source]
Decode the JSON form of
ExportTraceServiceRequest(file exporter output).- Parameters:
document (Any)
- Return type:
list[RawSpan]
- stormlog.infer.vllm_spans.decode_jsonl_span(raw)[source]
One line of the sink format: name, start/end unix ns, flat attributes.
- Parameters:
raw (dict[str, Any])
- Return type:
- stormlog.infer.vllm_spans.read_span_file(path)[source]
Read OTLP JSON (one document, or one per line) or sink JSONL; say which.
- Parameters:
path (str | Path)
- Return type:
tuple[str, list[RawSpan]]
- stormlog.infer.vllm_spans.retained_bytes(spans, clocks)[source]
What each span’s record will hold, estimated as the decode budget is (
SPAN_BYTES,VALUE_BYTESand the size of its text), its status and the request ID it keeps apart from its attributes included; a resource, a scope and a clock domain shared by several spans are charged once, to the first.- Parameters:
spans (Sequence[RawSpan])
clocks (Mapping[int, str])
- Return type:
list[int]
- stormlog.infer.vllm_spans.span_clock_domain(resource, fallback_host)[source]
The exporter’s wall clock, named by its host; never a shared clock.
vLLM’s resource carries no
host.name, so for spans the receiver collects the domain is named by the peer address the export came from, such as127.0.0.1/unix_epoch_ns. Without a boot ID it never counts as the client’s clock, even on one machine.- Parameters:
resource (dict[str, Any])
fallback_host (str)
- Return type:
str
- stormlog.infer.vllm_spans.span_record(raw, *, session_id, run_id, source, clock_domain, received_at_ns=None)[source]
- Parameters:
raw (RawSpan)
session_id (str)
run_id (str)
source (str)
clock_domain (str)
received_at_ns (int | None)
- Return type:
- class stormlog.infer.vllm_spans.ReceiverStats(requests: 'int' = 0, spans: 'int' = 0, decode_failures: 'int' = 0, unsupported_media: 'int' = 0, protobuf_unavailable: 'int' = 0, grpc_attempts: 'int' = 0, oversized: 'int' = 0, too_large: 'int' = 0, bad_requests: 'int' = 0, handler_errors: 'int' = 0, after_stop: 'int' = 0, refused_connections: 'int' = 0, busy: 'int' = 0, header_timeouts: 'int' = 0, body_timeouts: 'int' = 0, scan_timeouts: 'int' = 0, too_many_spans: 'int' = 0, dropped_queue_full: 'int' = 0, by_media: 'dict[str, int]' = <factory>)[source]
Bases:
object- Parameters:
requests (int)
spans (int)
decode_failures (int)
unsupported_media (int)
protobuf_unavailable (int)
grpc_attempts (int)
oversized (int)
too_large (int)
bad_requests (int)
handler_errors (int)
after_stop (int)
refused_connections (int)
busy (int)
header_timeouts (int)
body_timeouts (int)
scan_timeouts (int)
too_many_spans (int)
dropped_queue_full (int)
by_media (dict[str, int])
- requests: int = 0
- spans: int = 0
- decode_failures: int = 0
- unsupported_media: int = 0
- grpc_attempts: int = 0
- oversized: int = 0
- too_large: int = 0
- bad_requests: int = 0
- handler_errors: int = 0
- after_stop: int = 0
- refused_connections: int = 0
- busy: int = 0
- header_timeouts: int = 0
- body_timeouts: int = 0
- scan_timeouts: int = 0
- too_many_spans: int = 0
- dropped_queue_full: int = 0
- by_media: dict[str, int]
- class stormlog.infer.vllm_spans.ProtobufRates(message, element, byte, unknown_byte)[source]
Bases:
objectWhat a protobuf parse holds, in bytes, for each thing the wire scan counts (
WireCounts).- Parameters:
message (int)
element (int)
byte (int)
unknown_byte (int)
- message: int
- element: int
- byte: int
- unknown_byte: int
- class stormlog.infer.vllm_spans.ReceiverLimits(max_connections=8, request_deadline_seconds=10.0, max_inflight_bytes=134217728, max_spans_per_body=10000, max_queued_spans=100000, max_queued_bytes=67108864)[source]
Bases:
objectAdmission bounds, applied before a body is read, decoded or queued.
A connection over
max_connectionsis answered 503 and closed without a handler thread. Each request, its request line, headers and body, must arrive withinrequest_deadline_secondsof when the receiver starts waiting for it, and a protobuf body’s wire scan finish within it too; a kept-alive connection idle that long is closed. The exports being read and decoded at once may be charged at mostmax_inflight_bytes, each step charged before it runs: the body, with the most a gzip body can inflate to, then an estimate of decoding it (json_decode_estimate, or for protobuf its parse, from its messages counted on the wire, and then the spans it holds). An export that does not fit now is answered 503; one that never could, 413. A body with more thanmax_spans_per_bodyspans is refused, a protobuf one on its wire counts before it is parsed, and one that does not fit the queue whole is refused with 503, so the exporter can resend it; spans are never queued in part.- Parameters:
max_connections (int)
request_deadline_seconds (float)
max_inflight_bytes (int)
max_spans_per_body (int)
max_queued_spans (int)
max_queued_bytes (int)
- max_connections: int = 8
- request_deadline_seconds: float = 10.0
- max_inflight_bytes: int = 134217728
- max_spans_per_body: int = 10000
- max_queued_spans: int = 100000
- max_queued_bytes: int = 67108864
- stormlog.infer.vllm_spans.json_decode_estimate(content)[source]
What decoding this OTLP JSON export holds at once, estimated on its bytes before anything is parsed.
Counted: structural tokens (every value is opened by one, or follows a comma or colon), spans (each has a
"name") and attribute values (each has one*Valuekey). A key written with an escape, such as"na\u006de", is missed here;json_document_counts()counts the spans and values exactly once the document is parsed, and the receiver charges what this missed before any span is built.- Parameters:
content (bytes | bytearray)
- Return type:
int
- stormlog.infer.vllm_spans.json_document_counts(document)[source]
The spans and attribute values
decode_otlp_json()would build from a parsed document, counted without building them.- Parameters:
document (Any)
- Return type:
tuple[int, int]
- stormlog.infer.vllm_spans.gunzip_capped(body, cap)[source]
Inflate a gzip body, or None when its output would exceed
cap.A gzip member a few hundred kilobytes long can hold gigabytes of zeros, so the decoder is asked for at most
cap + 1bytes: one byte over the cap, or input left unconsumed, refuses the body without inflating it whole. A stream cut before its trailer raisesValueError; a malformed one raiseszlib.error.- Parameters:
body (bytes | bytearray)
cap (int)
- Return type:
bytes | None
- class stormlog.infer.vllm_spans.OtlpSpanReceiver(*, listen, session_id, run_id, limits=None)[source]
Bases:
objectAccept OTLP/HTTP trace exports on a local port while a profile runs.
- Parameters:
listen (str)
session_id (str)
run_id (str)
limits (ReceiverLimits | None)
- property listen: str
- property enabled: list[str]
The receiver paths this process can serve.
- stop()[source]
Stop for good: nothing is queued after this returns.
New requests are refused first, then the listener closes, every accepted connection is shut down, and handlers still inside a request get a short grace to finish, so what they decoded is in the queue for the final drain and nothing can arrive after it.
- Return type:
None
- property stopped: bool
the listener is closed for good.
- Type:
True once
stophas run
- drain()[source]
- Return type:
list[VllmSpanRecord]
- stormlog.infer.vllm_spans.span_capability_event(context, *, receiver, listen, error)[source]
The span capability record for a run that asked for a receiver.
supportednames every ingest path,enabledthe ones the receiver could serve (protobuf only with theinfer-otlpextra), andcollectedthe ones that delivered at least one span. A receiver that could not listen is unavailable, with the error kept.- Parameters:
context (CorrelationContext)
receiver (OtlpSpanReceiver | None)
listen (str)
error (str | None)
- Return type:
- stormlog.infer.vllm_spans.parse_listen_address(listen)[source]
HOST:PORTfor the receiver; the host defaults to loopback.- Parameters:
listen (str)
- Return type:
tuple[str, int]
- stormlog.infer.vllm_spans.spans_to_correlation_events(spans, *, producer_id='vllm.otel')[source]
Map request spans onto v2 request and stage records.
The request record carries the span’s own timestamps as reported by vLLM. Each stage window is placed from the span start by adding the reported durations in scheduler order, so the windows are estimates; the native duration attribute is kept in each stage’s metadata.
- Parameters:
spans (Iterable[VllmSpanRecord])
producer_id (str)
- Return type:
list[CorrelationEvent]