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

decode_jsonl_span(raw)

One line of the sink format: name, start/end unix ns, flat attributes.

decode_otlp_json(document)

Decode the JSON form of ExportTraceServiceRequest (file exporter output).

decode_otlp_protobuf(data)

Decode an ExportTraceServiceRequest with the generated classes.

gunzip_capped(body, cap)

Inflate a gzip body, or None when its output would exceed cap.

json_decode_estimate(content)

What decoding this OTLP JSON export holds at once, estimated on its bytes before anything is parsed.

json_document_counts(document)

The spans and attribute values decode_otlp_json() would build from a parsed document, counted without building them.

otlp_protobuf_available()

parse_listen_address(listen)

HOST:PORT for the receiver; the host defaults to loopback.

parse_otlp_protobuf(data)

The ExportTraceServiceRequest message, before any span is built.

protobuf_parse_estimate(counts, content_bytes)

An upper bound on what parsing an export holds, from its wire counts.

protobuf_rates()

What a parse costs with the installed protobuf backend.

read_span_file(path)

Read OTLP JSON (one document, or one per line) or sink JSONL; say which.

retained_bytes(spans, clocks)

What each span's record will hold, estimated as the decode budget is (SPAN_BYTES, VALUE_BYTES and 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.

span_capability_event(context, *, receiver, ...)

The span capability record for a run that asked for a receiver.

span_clock_domain(resource, fallback_host)

The exporter's wall clock, named by its host; never a shared clock.

span_record(raw, *, session_id, run_id, ...)

spans_from_message(message)

spans_to_correlation_events(spans, *[, ...])

Map request spans onto v2 request and stage records.

Classes

OtlpSpanReceiver(*, listen, session_id, run_id)

Accept OTLP/HTTP trace exports on a local port while a profile runs.

ProtobufRates(message, element, byte, ...)

What a protobuf parse holds, in bytes, for each thing the wire scan counts (WireCounts).

RawSpan(name[, trace_id, span_id, ...])

A decoded span before it is stamped with run identity.

ReceiverLimits([max_connections, ...])

Admission bounds, applied before a body is read, decoded or queued.

ReceiverStats(requests, spans, ...)

Exceptions

OtlpProtobufUnavailable

The infer-otlp extra is not installed.

ProtobufDecodeError

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: object

A 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: ValueError

The bytes are not an OTLP trace export request.

exception stormlog.infer.vllm_spans.OtlpProtobufUnavailable[source]

Bases: RuntimeError

The infer-otlp extra is not installed.

stormlog.infer.vllm_spans.otlp_protobuf_available()[source]
Return type:

bool

stormlog.infer.vllm_spans.protobuf_rates()[source]

What a parse costs with the installed protobuf backend.

Return type:

ProtobufRates

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:
Return type:

int

stormlog.infer.vllm_spans.decode_otlp_protobuf(data)[source]

Decode an ExportTraceServiceRequest with the generated classes.

Parameters:

data (bytes)

Return type:

list[RawSpan]

stormlog.infer.vllm_spans.parse_otlp_protobuf(data)[source]

The ExportTraceServiceRequest message, 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:

RawSpan

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_BYTES and 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 as 127.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:

VllmSpanRecord

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
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]
class stormlog.infer.vllm_spans.ProtobufRates(message, element, byte, unknown_byte)[source]

Bases: object

What 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: object

Admission bounds, applied before a body is read, decoded or queued.

A connection over max_connections is answered 503 and closed without a handler thread. Each request, its request line, headers and body, must arrive within request_deadline_seconds of 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 most max_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 than max_spans_per_body spans 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 *Value key). 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 + 1 bytes: one byte over the cap, or input left unconsumed, refuses the body without inflating it whole. A stream cut before its trailer raises ValueError; a malformed one raises zlib.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: object

Accept 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.

start()[source]
Return type:

None

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 stop has run

drain()[source]
Return type:

list[VllmSpanRecord]

config_record()[source]
Return type:

dict[str, Any]

capability_metadata()[source]
Return type:

dict[str, Any]

stormlog.infer.vllm_spans.span_capability_event(context, *, receiver, listen, error)[source]

The span capability record for a run that asked for a receiver.

supported names every ingest path, enabled the ones the receiver could serve (protobuf only with the infer-otlp extra), and collected the ones that delivered at least one span. A receiver that could not listen is unavailable, with the error kept.

Parameters:
Return type:

CapabilityEvent

stormlog.infer.vllm_spans.parse_listen_address(listen)[source]

HOST:PORT for 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:
Return type:

list[CorrelationEvent]