From c1a6ba6452c54a31499efa5f7b7887227923f7f6 Mon Sep 17 00:00:00 2001 From: nhobin219 <92616895+nhobin219@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:46:42 +0000 Subject: [PATCH 1/2] feat(examples): metrics and StreamMetricExporter in the OTel example examples/otel/metrics.py publishes OTel metric collections to a stream, one row per data point. It covers sums, gauges and both histogram kinds, keeps exemplars, and asks for delta temporality for counters and histograms so a window's total is a sum(). The services count orders by outcome and record http.server.request.duration inside their spans. export.py regroups metric rows into MetricsData for OTel's OTLP metric exporter. just demo otel moves to otel-gui 3.0.0, the first release to accept metrics, and warms its metrics decoder too: the lazy-load race that crashes it is unchanged in 3.0.0 (5 of 5 fresh instances died with all three signals at once; none after warming). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi --- CHANGELOG.md | 17 +++ examples/README.md | 51 +++++-- examples/__main__.py | 4 +- examples/otel/demo.py | 73 ++++++++-- examples/otel/export.py | 189 ++++++++++++++++++++++++- examples/otel/gui.py | 24 ++-- examples/otel/metrics.py | 281 +++++++++++++++++++++++++++++++++++++ examples/otel/services.py | 84 +++++++++-- tests/test_otel_example.py | 198 +++++++++++++++++++++++++- 9 files changed, 864 insertions(+), 57 deletions(-) create mode 100644 examples/otel/metrics.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b31940..8b3921c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,23 @@ All notable changes are recorded here. Versions follow The 0.1.0 entry describes what the library is rather than what changed, since there was nothing to have changed from. Everything above it is ordinary. +## Unreleased + +### Added + +- **Metrics in the OTel example** (`examples/otel/metrics.py`): + - **`StreamMetricExporter`**, an OTel metric exporter that publishes each + collection to a stream, one row per data point. It covers sums, gauges + and both kinds of histogram, with exemplars kept, so a metric leads to + the trace it measured. + - **Delta temporality** for counters and histograms, so a window's total + is a `sum()` over the stored table. + - **The services record metrics:** an order count by outcome and + `http.server.request.duration`. + - **`otel/export.py` re-exports metrics as OTLP**, and **`just demo otel` + shows them in otel-gui's Metrics tab.** The dashboard moves to otel-gui + 3.0.0, the first release that accepts metrics. + ## 0.13.1 — 2026-10-03 ### Changed diff --git a/examples/README.md b/examples/README.md index 8f9a39d..cc4d085 100644 --- a/examples/README.md +++ b/examples/README.md @@ -38,7 +38,7 @@ just demo NAME [ARGS] # run one, with all its processes | `live` | as `trades`, with a live-only broker | no log: nothing to replay | | `fastapi` | `fastapi_app.py` as the broker, then as `trades` | the broker mounted in a FastAPI app | | `book` | broker, `book/producer.py`, `book/index.html` | Bitstamp's live order book as a keyed table log, kept by a browser page | -| `otel` | otel-gui, broker, `otel/services.py`, `otel/analytics.py`, `otel/export.py` | OpenTelemetry logs and traces through streams, in a dashboard | +| `otel` | otel-gui, broker, `otel/services.py`, `otel/analytics.py`, `otel/export.py` | OpenTelemetry logs, traces and metrics through streams, in a dashboard | | `clean` | | delete what the demos stored | `just demo` starts a demo's processes in order, waiting for each one that @@ -260,22 +260,39 @@ against E. An empty diff means the migration does what production does. Once it's empty, the cutover is C reading D's stream, and B retires with its history still a queryable table. -## OpenTelemetry logs and traces, in a dashboard +## OpenTelemetry logs, traces and metrics, in a dashboard ``` just demo otel # otel-gui, a broker, the services, analytics, and the OTLP exporter uv run python -m examples.otel.demo # the same pipeline once, printing what each part saw ``` -Two simulated services, checkout and payments, handle traced orders and log -through Python's `logging` and the OpenTelemetry SDK. Their log records and -spans are published to two streams, `logs` and `spans`. `otel/export.py` -follows both and re-exports every row as OTLP to [otel-gui](https://github.com/metafab/otel-gui), -a local dashboard where logs, traces and the service map fill in live. A +Two simulated services, checkout and payments, handle traced orders, log +through Python's `logging`, and record metrics, all through the OpenTelemetry +SDK. Their log records, spans and metric data points are published to three +streams, `logs`, `spans` and `metrics`. `otel/export.py` follows all three +and re-exports every row as OTLP to [otel-gui](https://github.com/metafab/otel-gui), +a local dashboard where logs, traces, metrics and the service map fill in live. A failed order is one trace across both services: payments' `POST /charge` span with a `card declined` event and error log, and checkout's request marked failed with an `order failed` warning. +**The metrics** are an order count by outcome (`shop.orders`) and each +service's request duration (`http.server.request.duration`, OTel's semantic +convention), exported every second. A stream row is one data point, carrying +its metric's name, unit, type and temporality. Counters and histograms are +exported as deltas, so each row is what happened in its own second and a +window's total is a `sum()`: + +```sql +SELECT attributes['outcome'].string_value AS outcome, sum(value_int) AS orders +FROM log WHERE name = 'shop.orders' GROUP BY outcome +``` + +A measurement made inside a span keeps that span's trace as an *exemplar*, +so a metric leads straight to a request: the failed order's count points at +the failed order's trace. + `otel/analytics.py` is real-time analytics on the same telemetry, in one `Stream.live` view of `spans`. Every few seconds it asks, in one SQL query, each service's error rate over the last 30 seconds against its error rate @@ -291,26 +308,29 @@ current, and the question is a query. | file | what it holds | |---|---| -| `otel/common.py` | what both signals share: `AnyValue`, ids, scope, and publishing from OTel's export thread | +| `otel/common.py` | what the signals share: `AnyValue`, ids, scope, and publishing from OTel's export thread | | `otel/logs.py` | the log record's schema, its row conversion, and `StreamLogExporter` | | `otel/spans.py` | the span's schema, its row conversion, and `StreamSpanExporter` | +| `otel/metrics.py` | the metric data point's schema, its row conversion, and `StreamMetricExporter` | | `otel/demo.py` | the two services, the broker, and the one-shot demo | -| `otel/export.py` | rows back to OTel records and spans, out through OTel's OTLP exporters | -| `otel/services.py` | the producer: two services logging and tracing through OTel, publishing to the broker | +| `otel/export.py` | rows back to OTel records, spans and metric exports, out through OTel's OTLP exporters | +| `otel/services.py` | the producer: two services logging, tracing and recording metrics through OTel, publishing to the broker | | `otel/analytics.py` | real-time analytics with `Stream.live`: error rate now against always, per service | | `otel/gui.py` | otel-gui, downloaded, checked and run: where `export.py` sends what it reads | **None of this is in streamcast.** The schemas and conversions are built from the column types any stream can declare: trace and span ids as hex binary, -attributes as a map, OTel's `AnyValue` as a struct, and a span's events and -links as lists of structs. The OTel packages are dev dependencies, for this +attributes as a map, OTel's `AnyValue` as a struct, a span's events and +links as lists of structs, and a histogram's buckets as lists. The OTel packages are dev dependencies, for this example only. **`otel/export.py` is an ordinary OTLP exporter.** OTel viewers are *receivers*: telemetry is pushed to them, and none subscribes to a WebSocket. So this subscribes to the streams, turns each row back into an SDK log record or span, and hands it to OpenTelemetry's own batch processors and OTLP/HTTP -exporters. The batching, the protobuf encoding and the retries are OTel's, +exporters. Metrics have no batch processor to hand a point to, so metric rows +are regrouped every half second into the resource, scope and metric nesting +an export carries, and handed to OTel's OTLP metric exporter. The batching, the protobuf encoding and the retries are OTel's, and it works with any OTLP/HTTP receiver. Point `--receiver` at an OTel Collector and it feeds whatever the Collector does. @@ -324,4 +344,7 @@ The one-shot run shows what a subscriber can do beyond a dashboard: - one failed request replayed by its trace id from both streams, with the id given as hex text in `where=`; - SQL over the stored tables: errors per service, the slowest request from - its root span, and ingest lag from `streamcast_ts`. + its root span, and ingest lag from `streamcast_ts`; +- SQL over the metrics: orders by outcome, mean request duration per service + from the histograms, and the failed order's exemplar, which is the failed + request's trace id. diff --git a/examples/__main__.py b/examples/__main__.py index 67c26ff..9e2c260 100644 --- a/examples/__main__.py +++ b/examples/__main__.py @@ -123,7 +123,7 @@ def trades( ], ), "otel": Demo( - "OpenTelemetry logs and traces through streams, in otel-gui", + "OpenTelemetry logs, traces and metrics through streams, in otel-gui", [ Process("otel-gui", module("examples.otel.gui"), 4318), Process( @@ -134,6 +134,8 @@ def trades( "logs=examples.otel.logs:SCHEMA", "--stream", "spans=examples.otel.spans:SCHEMA", + "--stream", + "metrics=examples.otel.metrics:SCHEMA", "--port", "8766", "--root", diff --git a/examples/otel/demo.py b/examples/otel/demo.py index 829e7f8..815dfc8 100644 --- a/examples/otel/demo.py +++ b/examples/otel/demo.py @@ -1,10 +1,10 @@ -"""OpenTelemetry logs and traces, through two streams, into tables you can query. +"""OpenTelemetry logs, traces and metrics, through streams, into tables you can query. uv run python -m examples.otel.demo # once, printing what it saw just demo otel # the same roles as processes, live Nothing here is part of streamcast. The OTel schemas and conversions live in -`logs.py` and `spans.py`, built from ordinary column types. The OTel packages +`logs.py`, `spans.py` and `metrics.py`, built from ordinary column types. The OTel packages are dev dependencies, for this example only. The three roles, in one process so a test can run it: the broker, the @@ -12,17 +12,20 @@ separate processes, talking over sockets, and nothing below changes but the URI. `just demo otel` runs them as separate processes. -1. **A broker** serves two streams, `logs` and `spans`, each with a log, - accepting remote publishers (`publish=True`). +1. **A broker** serves three streams, `logs`, `spans` and `metrics`, each + with a log, accepting remote publishers (`publish=True`). 2. **The producer**: two services, checkout and payments, handle traced - requests and log through the standard `logging` module. OTel's SDK turns - each into records and spans, published by the exporters in `logs.py` and - `spans.py`. + requests, log through the standard `logging` module, and count orders and + time requests. OTel's SDK turns each into records, spans and metric data + points, published by the exporters in `logs.py`, `spans.py` and + `metrics.py`. 3. **A live tail** follows the problems as they happen — `where=` a membership filter on `severity_text` — and **a replay** reads back one request by its trace id from both streams, `where={"trace_id": "<32 hex chars>"}`. -4. **The stored tables** answer what a log search and a trace view would: - errors per service, the slowest request, and ingest lag. +4. **The stored tables** answer what a log search, a trace view and a metrics + dashboard would: errors per service, the slowest request, ingest lag, + orders by outcome, mean request duration per service, and the failed + order's exemplar, which names the failed request's trace. """ from __future__ import annotations @@ -33,7 +36,7 @@ from typing import TYPE_CHECKING, Any import streamcast -from examples.otel import logs, spans +from examples.otel import logs, metrics, spans from examples.otel.services import exporting, traffic if TYPE_CHECKING: @@ -42,10 +45,13 @@ # -- the broker ------------------------------------------------------------------- -def streams(root: Path) -> tuple[streamcast.Stream, streamcast.Stream]: +def streams( + root: Path, +) -> tuple[streamcast.Stream, streamcast.Stream, streamcast.Stream]: return ( streamcast.Stream.new("logs", root=root, schema=logs.SCHEMA), streamcast.Stream.new("spans", root=root, schema=spans.SCHEMA), + streamcast.Stream.new("metrics", root=root, schema=metrics.SCHEMA), ) @@ -75,9 +81,9 @@ async def one_trace( async def main(root: Path) -> dict[str, Any]: """Run the whole demo under `root`, print it, and return what it saw.""" - log_stream, span_stream = streams(root) + log_stream, span_stream, metric_stream = streams(root) server = await streamcast.serve( - [log_stream, span_stream], "127.0.0.1", 0, publish=True + [log_stream, span_stream, metric_stream], "127.0.0.1", 0, publish=True ) base = f"ws://127.0.0.1:{server.sockets[0].getsockname()[1]}" try: @@ -116,13 +122,40 @@ async def main(root: Path) -> dict[str, Any]: log_stream, "SELECT max(streamcast_ts * 1000 - time_unix_nano) / 1e6 AS ms FROM log", ) + # Each row of `shop.orders` is one interval's count (delta temporality), + # so the total is a sum. + outcomes = query( + metric_stream, + "SELECT attributes['outcome'].string_value AS outcome, " + "sum(value_int)::BIGINT AS orders FROM log " + "WHERE name = 'shop.orders' GROUP BY outcome ORDER BY outcome", + ) + durations = query( + metric_stream, + "SELECT service, sum(count)::BIGINT AS requests, " + "sum(sum) / sum(count) * 1e3 AS mean_ms FROM log " + "WHERE name = 'http.server.request.duration' " + "GROUP BY service ORDER BY service", + ) + # The failed order was counted inside its request's span, so the SDK + # kept it as an exemplar: from the metric straight to the trace. + exemplars = query( + metric_stream, + "SELECT lower(hex(e.trace_id)) AS trace FROM (" + "SELECT unnest(exemplars) AS e FROM log WHERE name = 'shop.orders' " + "AND attributes['outcome'].string_value = 'failed')", + ) [stored_logs] = query(log_stream, "SELECT count(*) AS n FROM log") [stored_spans] = query(span_stream, "SELECT count(*) AS n FROM log") + [stored_points] = query(metric_stream, "SELECT count(*) AS n FROM log") finally: server.close() await server.wait_closed() - print(f"{stored_logs['n']} log records and {stored_spans['n']} spans stored") + print( + f"{stored_logs['n']} log records, {stored_spans['n']} spans and " + f"{stored_points['n']} metric data points stored" + ) print("live tail, errors and warnings as they happened:") for message in live: print( @@ -142,6 +175,14 @@ async def main(root: Path) -> dict[str, Any]: print(f"errors per service: {errors}") print(f"slowest request: {slowest['ms']:.1f} ms") print(f"ingest lag, worst: {lag['ms']:.1f} ms") + print(f"orders by outcome: {outcomes}") + for row in durations: + print( + f" {row['service']:>8} {row['requests']} requests, " + f"mean {row['mean_ms']:.2f} ms" + ) + + print(f"the failed order's exemplar: trace {[e['trace'] for e in exemplars]}") return { "logs": stored_logs["n"], @@ -153,6 +194,10 @@ async def main(root: Path) -> dict[str, Any]: "errors": errors, "slowest_ms": slowest["ms"], "lag_ms": lag["ms"], + "points": stored_points["n"], + "outcomes": outcomes, + "durations": durations, + "exemplars": [e["trace"] for e in exemplars], } diff --git a/examples/otel/export.py b/examples/otel/export.py index 0f3deac..e11aadd 100644 --- a/examples/otel/export.py +++ b/examples/otel/export.py @@ -10,6 +10,12 @@ local viewer like [otel-gui](https://github.com/metafab/otel-gui), a Collector, or a vendor's endpoint. +Metrics have no batch processor to hand a data point to: the SDK's metric +pipeline starts at instruments, not at points. So rows of `metrics.py`'s +schema are collected for half a second, grouped back into the +resource → scope → metric nesting an export carries, and handed to OTel's +OTLP metric exporter as one `MetricsData`. + It exists because OTel viewers are RECEIVERS: data is pushed to them, and none subscribes to a WebSocket. A stream is where the telemetry lives; this is how any of them reads it. @@ -25,15 +31,33 @@ import asyncio import contextlib import json +import math import signal from typing import TYPE_CHECKING, Any from opentelemetry._logs import LogRecord, SeverityNumber from opentelemetry.attributes import BoundedAttributes from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter +from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter from opentelemetry.sdk._logs import ReadWriteLogRecord from opentelemetry.sdk._logs.export import BatchLogRecordProcessor +from opentelemetry.sdk.metrics import Exemplar +from opentelemetry.sdk.metrics.export import ( + AggregationTemporality, + Buckets, + ExponentialHistogram, + ExponentialHistogramDataPoint, + Gauge, + Histogram, + HistogramDataPoint, + Metric, + MetricsData, + NumberDataPoint, + ResourceMetrics, + ScopeMetrics, + Sum, +) from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import Event, ReadableSpan from opentelemetry.sdk.trace.export import BatchSpanProcessor @@ -53,7 +77,7 @@ import streamcast if TYPE_CHECKING: - from collections.abc import Callable, Mapping + from collections.abc import Callable, Iterable, Mapping from opentelemetry.context import Context @@ -190,6 +214,157 @@ def span(row: Mapping[str, Any]) -> ReadableSpan: ) +# -- metric rows back to an export ------------------------------------------------ + + +def number(row: Mapping[str, Any]) -> float | None: + """`value_int` or `value_double`, whichever the row holds.""" + return row["value_int"] if row["value_int"] is not None else row["value_double"] + + +def exemplars(row: Mapping[str, Any]) -> list[Exemplar]: + # An exemplar always has a value; a null one was non-finite, stored as + # null because a stream refuses NaN and ±inf. NaN is the closest return. + return [ + Exemplar( + frozen(e["filtered_attributes"]), + math.nan if number(e) is None else number(e), # ty: ignore[invalid-argument-type] + e["time_unix_nano"], + ident(e["span_id"]) or None, + ident(e["trace_id"]) or None, + ) + for e in row["exemplars"] or () + ] + + +def buckets(stored: Mapping[str, Any] | None) -> Buckets: + stored = stored or {} + return Buckets(stored.get("offset") or 0, list(stored.get("bucket_counts") or ())) + + +def point( + row: Mapping[str, Any], +) -> NumberDataPoint | HistogramDataPoint | ExponentialHistogramDataPoint: + """One stored data point back as the SDK's: the inverse of `metrics.rows`.""" + shared = { + "attributes": frozen(row["attributes"]), + # None for a gauge, which has no start, as the SDK leaves it. + "start_time_unix_nano": row["start_time_unix_nano"], + "time_unix_nano": row["time_unix_nano"], + "exemplars": exemplars(row), + } + if row["type"] in ("sum", "gauge"): + return NumberDataPoint(value=number(row), **shared) # ty: ignore[invalid-argument-type] + + histogram = { + "count": row["count"], + "sum": row["sum"], + "min": row["min"], + "max": row["max"], + } + if row["type"] == "histogram": + return HistogramDataPoint( + bucket_counts=tuple(row["bucket_counts"] or ()), + explicit_bounds=tuple(row["explicit_bounds"] or ()), + **histogram, + **shared, + ) + + return ExponentialHistogramDataPoint( + scale=row["scale"], + zero_count=row["zero_count"], + positive=buckets(row["positive"]), + negative=buckets(row["negative"]), + flags=row["flags"] or 0, + **histogram, + **shared, + ) + + +def data( + row: Mapping[str, Any], points: list +) -> Sum | Gauge | Histogram | ExponentialHistogram: + """The metric's data, of the row's type, holding `points`.""" + kind = row["type"] + if kind == "gauge": + return Gauge(points) + + temporality = AggregationTemporality(row["temporality"]) + if kind == "sum": + return Sum(points, temporality, bool(row["is_monotonic"])) + + if kind == "histogram": + return Histogram(points, temporality) + + return ExponentialHistogram(points, temporality) + + +def metrics_data(rows: Iterable[Mapping[str, Any]]) -> MetricsData: + """Stored data points regrouped as one export: resource, scope, metric, points. + + A metric is identified as OTLP identifies one: by its resource and scope, + and its name, unit, type and temporality. Points keep their order. + """ + grouped: dict[str, dict[str, dict[tuple, list[Mapping[str, Any]]]]] = {} + for row in rows: + resource = json.dumps(row["resource"], sort_keys=True, default=str) + scope_key = json.dumps(row["scope"], sort_keys=True) + metric = ( + row["name"], + row["description"], + row["unit"], + row["type"], + row["temporality"], + row["is_monotonic"], + ) + by_scope = grouped.setdefault(resource, {}) + by_scope.setdefault(scope_key, {}).setdefault(metric, []).append(row) + + resource_metrics = [] + for resource, by_scope in grouped.items(): + scope_metrics = [] + for by_metric in by_scope.values(): + found = [] + for points in by_metric.values(): + first = points[0] + found.append( + Metric( + first["name"], + first["description"] or "", + first["unit"] or "", + data(first, [point(row) for row in points]), + ) + ) + + first = next(iter(by_metric.values()))[0] + scope_metrics.append( + ScopeMetrics(scope(first) or InstrumentationScope(""), found, "") + ) + + resource_metrics.append( + ResourceMetrics( + Resource(attributes(json.loads(resource))), scope_metrics, "" + ) + ) + + return MetricsData(resource_metrics) + + +async def export_metrics( + exporter: OTLPMetricExporter, + pending: list[Mapping[str, Any]], + every: float = 0.5, +) -> None: + """Every `every` seconds, what has arrived goes out as one OTLP export.""" + while True: + await asyncio.sleep(every) + if pending: + batch = pending[:] + pending.clear() + # Off the loop: the exporter's HTTP and retries block. + await asyncio.to_thread(exporter.export, metrics_data(batch)) + + # -- following the streams -------------------------------------------------------- @@ -228,6 +403,8 @@ async def main(broker: str, receiver: str) -> None: span_processor = BatchSpanProcessor( OTLPSpanExporter(endpoint=f"{receiver}/v1/traces"), schedule_delay_millis=500 ) + metric_exporter = OTLPMetricExporter(endpoint=f"{receiver}/v1/metrics") + metric_rows: list[Mapping[str, Any]] = [] print(f"exporting OTLP to {receiver}", flush=True) try: async with asyncio.TaskGroup() as group: @@ -237,14 +414,22 @@ async def main(broker: str, receiver: str) -> None: group.create_task( follow(f"{broker}/spans", lambda row: span_processor.on_end(span(row))) ) + group.create_task(follow(f"{broker}/metrics", metric_rows.append)) + group.create_task(export_metrics(metric_exporter, metric_rows)) finally: log_processor.shutdown() span_processor.shutdown() + if metric_rows: + metric_exporter.export(metrics_data(metric_rows)) + + metric_exporter.shutdown() if __name__ == "__main__": parser = argparse.ArgumentParser(description=__doc__.split("\n\n")[0]) - parser.add_argument("--broker", default=BROKER, help="serving /logs and /spans") + parser.add_argument( + "--broker", default=BROKER, help="serving /logs, /spans and /metrics" + ) parser.add_argument( "--receiver", default=RECEIVER, help="an OTLP/HTTP base URL, without /v1/..." ) diff --git a/examples/otel/gui.py b/examples/otel/gui.py index 6079774..0fcdb1b 100644 --- a/examples/otel/gui.py +++ b/examples/otel/gui.py @@ -3,7 +3,7 @@ just demo otel # this, a broker, the services, and the OTLP exporter [otel-gui](https://github.com/metafab/otel-gui) is a local OTLP receiver -with a dashboard for logs, traces and the service map. The first run +with a dashboard for logs, traces, metrics and the service map. The first run downloads its release for this platform, checks its SHA-256 against the published one, and caches it. It listens on 127.0.0.1, and Ctrl-C stops it. @@ -27,7 +27,8 @@ import urllib.request from pathlib import Path -VERSION = "2.1.0" +VERSION = "3.0.0" +"""3.0.0 is the first to take metrics: earlier ones answer `/v1/metrics` 501.""" ASSETS = { ("Linux", "x86_64"): "otel-gui-linux-x64", ("Linux", "aarch64"): "otel-gui-linux-arm64", @@ -95,16 +96,17 @@ def ready(port: int, gui: subprocess.Popen[bytes], log: Path) -> None: def warm(port: int) -> None: - """Have otel-gui load its trace and logs decoders, one after the other. - - otel-gui (2.1.0, and 3.0.0 unchanged) loads its trace and logs `.proto` - files lazily into one shared protobufjs Root, on the first request to - each. The exporter sends both at once, the two loads interleave, one - resolves before `resource.proto` is parsed, and the throw escapes into a - callback and kills the dashboard. An empty request to each, in turn, does - the loading before anything can race it. + """Have otel-gui load its trace, logs and metrics decoders, one at a time. + + otel-gui (2.1.0 and 3.0.0) loads each signal's `.proto` files lazily into + one shared protobufjs Root, on the first request to each. The exporter + sends all three at once, the loads interleave, one resolves before + `resource.proto` is parsed, and the throw escapes into a callback and + kills the dashboard: on 3.0.0, all three signals' first requests at once + killed it 5 times in 5. An empty request to each, in turn, does the + loading before anything can race it, and none died. """ - for signal_name in ("traces", "logs"): + for signal_name in ("traces", "logs", "metrics"): request = urllib.request.Request( # noqa: S310 f"http://127.0.0.1:{port}/v1/{signal_name}", data=b"", diff --git a/examples/otel/metrics.py b/examples/otel/metrics.py new file mode 100644 index 0000000..df8d85a --- /dev/null +++ b/examples/otel/metrics.py @@ -0,0 +1,281 @@ +"""OTel's metric data point as a stream's row, and an exporter that publishes it. + +**How a metric maps onto rows**, beyond what it shares with a log record: + +* One row per DATA POINT, not per export. An export nests resource, scope, + metric and points; the point is what you query, so each is a row carrying + its metric's name, unit, type and temporality beside its own values. +* `type` is one of `sum`, `gauge`, `histogram` and `exponential_histogram`, + and only that kind's columns are set. A summary is not here: the Python SDK + never produces one. +* A number is `value_int` or `value_double`, as OTLP keeps it: an int64 + counter must not round through a double. +* `temporality` is OTLP's number: 1 delta, 2 cumulative, null for a gauge. + The exporter asks for DELTA for counters and histograms, the SDK's + `delta` preference: each row is then what happened in its own interval, so + a window's total is a `sum()` and a histogram heatmap reads directly. + Up-down counters stay cumulative, as that preference keeps them: their + current value is the useful one. +* `exemplars` link a point to the traces it measured, by trace and span id. + The SDK records one only for a measurement made inside a sampled span. +* A non-finite double is stored as null, as in a log record. +""" + +from __future__ import annotations + +import math +from typing import TYPE_CHECKING, Any + +from opentelemetry.sdk.metrics import ( + Counter, + Histogram, + ObservableCounter, + ObservableGauge, + ObservableUpDownCounter, + UpDownCounter, +) +from opentelemetry.sdk.metrics.export import ( + AggregationTemporality, + ExponentialHistogram, + Gauge, + MetricExporter, + MetricExportResult, + Sum, +) + +from examples.otel import common + +if TYPE_CHECKING: + import asyncio + + from opentelemetry.sdk.metrics import Exemplar + from opentelemetry.sdk.metrics.export import MetricsData + + import streamcast + + +def _list_of(items: dict[str, Any]) -> dict[str, Any]: + return {"type": ["array", "null"], "items": items} + + +def _struct(properties: dict[str, Any]) -> dict[str, Any]: + return { + "type": ["object", "null"], + "properties": properties, + "required": [], + "additionalProperties": False, + } + + +COUNTS = _list_of({"type": "integer"}) +BUCKETS = _struct( + { + "offset": {"type": ["integer", "null"], "format": "int32"}, + "bucket_counts": COUNTS, + } +) +"""One side of an exponential histogram: counts from bucket `offset` up.""" + +SCHEMA: dict[str, Any] = { + "type": "object", + "properties": { + "time_unix_nano": {"type": "integer"}, + "start_time_unix_nano": {"type": ["integer", "null"]}, + "service": {"type": ["string", "null"]}, + "name": {"type": "string"}, + "description": {"type": ["string", "null"]}, + "unit": {"type": ["string", "null"]}, + "type": {"type": "string"}, + "temporality": {"type": ["integer", "null"], "format": "int32"}, + "is_monotonic": {"type": ["boolean", "null"]}, + "attributes": common.ATTRIBUTES, + # A sum or a gauge: exactly one of the two. + "value_int": {"type": ["integer", "null"]}, + "value_double": {"type": ["number", "null"]}, + # A histogram, either kind. + "count": {"type": ["integer", "null"]}, + "sum": {"type": ["number", "null"]}, + "min": {"type": ["number", "null"]}, + "max": {"type": ["number", "null"]}, + # An explicit-bucket histogram: `len(bounds) + 1` counts. + "bucket_counts": COUNTS, + "explicit_bounds": _list_of({"type": "number"}), + # An exponential histogram. + "scale": {"type": ["integer", "null"], "format": "int32"}, + "zero_count": {"type": ["integer", "null"]}, + "positive": BUCKETS, + "negative": BUCKETS, + "flags": {"type": ["integer", "null"], "format": "int32"}, + "exemplars": _list_of( + _struct( + { + "time_unix_nano": {"type": ["integer", "null"]}, + "value_int": {"type": ["integer", "null"]}, + "value_double": {"type": ["number", "null"]}, + "trace_id": common.TRACE_ID, + "span_id": common.SPAN_ID, + "filtered_attributes": common.ATTRIBUTES, + } + ) + ), + "scope": common.SCOPE, + "resource": common.ATTRIBUTES, + }, + "required": ["time_unix_nano", "name", "type"], +} + +DELTA = 1 +CUMULATIVE = 2 +"""OTLP's `AggregationTemporality` numbers, the same in Python's SDK.""" + +PREFERRED_TEMPORALITY: dict[type, AggregationTemporality] = { + Counter: AggregationTemporality.DELTA, + UpDownCounter: AggregationTemporality.CUMULATIVE, + Histogram: AggregationTemporality.DELTA, + ObservableCounter: AggregationTemporality.DELTA, + ObservableUpDownCounter: AggregationTemporality.CUMULATIVE, + ObservableGauge: AggregationTemporality.CUMULATIVE, +} +"""The SDK's `delta` temporality preference, as OTLP exporters spell it.""" + + +def _finite(value: float | None) -> float | None: + return value if value is None or math.isfinite(value) else None + + +def _number(value: float) -> dict[str, object]: + """`value_int` for an int, `value_double` for anything else, as OTLP does.""" + if isinstance(value, int) and not isinstance(value, bool): + return {"value_int": value, "value_double": None} + + return {"value_int": None, "value_double": _finite(float(value))} + + +def _exemplar(exemplar: Exemplar) -> dict[str, object]: + return { + "time_unix_nano": exemplar.time_unix_nano, + **_number(exemplar.value), + "trace_id": common.trace_id(exemplar.trace_id or 0), + "span_id": common.span_id(exemplar.span_id or 0), + "filtered_attributes": common.attributes(exemplar.filtered_attributes), + } + + +def _kind(data: object) -> str: + if isinstance(data, Sum): + return "sum" + + if isinstance(data, Gauge): + return "gauge" + + if isinstance(data, ExponentialHistogram): + return "exponential_histogram" + + return "histogram" + + +def rows(metrics_data: MetricsData) -> list[dict[str, object]]: + """One OTel metrics export as rows of `SCHEMA`, one per data point.""" + out: list[dict[str, object]] = [] + for resource_metrics in metrics_data.resource_metrics: + resource = dict(resource_metrics.resource.attributes) + for scope_metrics in resource_metrics.scope_metrics: + for metric in scope_metrics.metrics: + data = metric.data + kind = _kind(data) + temporality = getattr(data, "aggregation_temporality", None) + monotonic = getattr(data, "is_monotonic", None) + for point in data.data_points: + row: dict[str, object] = { + "time_unix_nano": point.time_unix_nano, + "start_time_unix_nano": getattr( + point, "start_time_unix_nano", None + ) + or None, + "service": resource.get("service.name"), + "name": metric.name, + "description": metric.description or None, + "unit": metric.unit or None, + "type": kind, + "temporality": None + if temporality is None + else temporality.value, + "is_monotonic": monotonic, + "attributes": common.attributes(point.attributes), + "exemplars": [_exemplar(e) for e in point.exemplars or ()] + or None, + "scope": common.scope(scope_metrics.scope), + "resource": common.attributes(resource), + } + if kind in ("sum", "gauge"): + row.update(_number(point.value)) # ty: ignore[unresolved-attribute] + else: + row.update( + count=point.count, # ty: ignore[unresolved-attribute] + sum=_finite(point.sum), # ty: ignore[unresolved-attribute] + min=_finite(point.min), # ty: ignore[unresolved-attribute] + max=_finite(point.max), # ty: ignore[unresolved-attribute] + ) + + if kind == "histogram": + row.update( + bucket_counts=list(point.bucket_counts), # ty: ignore[unresolved-attribute] + explicit_bounds=list(point.explicit_bounds), # ty: ignore[unresolved-attribute] + ) + + if kind == "exponential_histogram": + row.update( + scale=point.scale, # ty: ignore[unresolved-attribute] + zero_count=point.zero_count, # ty: ignore[unresolved-attribute] + flags=point.flags, # ty: ignore[unresolved-attribute] + positive={ + "offset": point.positive.offset, # ty: ignore[unresolved-attribute] + "bucket_counts": list(point.positive.bucket_counts), # ty: ignore[unresolved-attribute] + }, + negative={ + "offset": point.negative.offset, # ty: ignore[unresolved-attribute] + "bucket_counts": list(point.negative.bucket_counts), # ty: ignore[unresolved-attribute] + }, + ) + + out.append(row) + + return out + + +class StreamMetricExporter(MetricExporter): + """Publishes each collection of OTel metrics to a stream, a row per data point. + + Asks for delta temporality for counters and histograms (see the module + docstring); pass `preferred_temporality` to choose otherwise. + """ + + def __init__( + self, + publication: streamcast.Publication, + loop: asyncio.AbstractEventLoop, + preferred_temporality: dict[type, AggregationTemporality] | None = None, + ) -> None: + super().__init__( + preferred_temporality=preferred_temporality or PREFERRED_TEMPORALITY + ) + self._publication = publication + self._loop = loop + + def export( + self, + metrics_data: MetricsData, + timeout_millis: float = 10_000, # noqa: ARG002 + **kwargs: object, # noqa: ARG002 + ) -> MetricExportResult: + batch = rows(metrics_data) + if not batch or common.publish(self._publication, self._loop, batch): + return MetricExportResult.SUCCESS + + return MetricExportResult.FAILURE + + def force_flush(self, timeout_millis: float = 10_000) -> bool: # noqa: ARG002 + return True # nothing buffered here; see `StreamLogExporter` + + def shutdown(self, timeout_millis: float = 30_000, **kwargs: object) -> None: + pass diff --git a/examples/otel/services.py b/examples/otel/services.py index a3af6ee..6b87f1c 100644 --- a/examples/otel/services.py +++ b/examples/otel/services.py @@ -3,11 +3,13 @@ just demo otel # otel-gui, a broker, this, and the OTLP exporter Two simulated services, checkout and payments, handle orders as traced -requests and log through the standard `logging` module. OTel's SDK turns each -into log records and spans, and the exporters in `logs.py` and `spans.py` -publish them, as a client, to the broker's `logs` and `spans` streams. About -one order in five fails: payments declines the card, and the failed request -is one trace across both services. +requests, log through the standard `logging` module, and record metrics: an +order count by outcome and each request's duration. OTel's SDK turns each +into log records, spans and metric data points, and the exporters in +`logs.py`, `spans.py` and `metrics.py` publish them, as a client, to the +broker's `logs`, `spans` and `metrics` streams. About one order in five +fails: payments declines the card, and the failed request is one trace +across both services. """ from __future__ import annotations @@ -20,23 +22,30 @@ import random import signal import threading +import time from dataclasses import dataclass from typing import TYPE_CHECKING from opentelemetry.instrumentation.logging.handler import LoggingHandler from opentelemetry.sdk._logs import LoggerProvider from opentelemetry.sdk._logs.export import BatchLogRecordProcessor +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.trace import SpanKind, Status, StatusCode import streamcast -from examples.otel import logs, spans +from examples.otel import logs, metrics, spans if TYPE_CHECKING: + from opentelemetry.metrics import Counter, Histogram from opentelemetry.trace import Tracer +DURATION_BOUNDS = (0.0005, 0.001, 0.0025, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25) +"""Request-duration buckets, in seconds: these requests take milliseconds.""" + # -- the services --------------------------------------------------------------- @@ -44,6 +53,7 @@ class Exporters: logs: logs.StreamLogExporter spans: spans.StreamSpanExporter + metrics: metrics.StreamMetricExporter @dataclass @@ -52,7 +62,12 @@ class Service: tracer: Tracer log_provider: LoggerProvider tracer_provider: TracerProvider + meter_provider: MeterProvider handler: logging.Handler + orders: Counter + """Orders handled, by `outcome`: what a dashboard's error rate is.""" + duration: Histogram + """`http.server.request.duration`, OTel's semantic convention, in seconds.""" def detach(self) -> None: """Take the handler off the logger, which outlives this service. @@ -63,10 +78,12 @@ def detach(self) -> None: self.logger.removeHandler(self.handler) def shutdown(self) -> None: - """Detach, then flush and stop both providers.""" + """Detach, then flush and stop every provider.""" self.detach() self.tracer_provider.shutdown() self.log_provider.shutdown() + # Last: its final collection is the requests made before this call. + self.meter_provider.shutdown() def instrument(name: str, exporters: Exporters) -> Service: @@ -82,6 +99,12 @@ def instrument(name: str, exporters: Exporters) -> Service: tracer_provider.add_span_processor( BatchSpanProcessor(exporters.spans, schedule_delay_millis=500) ) + # Every second, so a chart moves; each collection is one row per series. + meter_provider = MeterProvider( + [PeriodicExportingMetricReader(exporters.metrics, export_interval_millis=1000)], + resource=resource, + ) + meter = meter_provider.get_meter("shop") logger = logging.getLogger(f"shop.{name}") logger.setLevel(logging.INFO) @@ -93,7 +116,17 @@ def instrument(name: str, exporters: Exporters) -> Service: tracer_provider.get_tracer("shop"), log_provider, tracer_provider, + meter_provider, handler, + meter.create_counter( + "shop.orders", unit="{order}", description="Orders handled" + ), + meter.create_histogram( + "http.server.request.duration", + unit="s", + description="Duration of HTTP server requests", + explicit_bucket_boundaries_advisory=DURATION_BOUNDS, + ), ) @@ -102,6 +135,7 @@ def order( ) -> str | None: """One traced request through both services. Returns its trace id if it failed.""" failed = Status(StatusCode.ERROR, "card declined") + started = time.perf_counter() with checkout.tracer.start_as_current_span( "POST /orders", kind=SpanKind.SERVER, attributes={"order.id": order_id} ) as request: @@ -118,6 +152,7 @@ def order( "POST /charge", kind=SpanKind.SERVER, attributes={"amount": amount} ) as charge, ): + charging = time.perf_counter() declined = amount > 100 if declined: charge.add_event("card declined", {"amount": amount}) @@ -136,13 +171,32 @@ def order( "card charged", extra={"order_id": order_id, "amount": amount} ) + # Inside the span, so the SDK can keep this measurement as an + # exemplar: a metric's link back to the trace it measured. + payments.duration.record( + time.perf_counter() - charging, + { + "http.route": "/charge", + "http.response.status_code": 402 if declined else 200, + }, + ) + if declined: request.set_status(failed) checkout.logger.warning("order failed", extra={"order_id": order_id}) - return format(request.get_span_context().trace_id, "032x") - - checkout.logger.info("order confirmed", extra={"order_id": order_id}) - return None + else: + checkout.logger.info("order confirmed", extra={"order_id": order_id}) + + outcome = "failed" if declined else "confirmed" + checkout.orders.add(1, {"outcome": outcome}) + checkout.duration.record( + time.perf_counter() - started, + { + "http.route": "/orders", + "http.response.status_code": 402 if declined else 201, + }, + ) + return format(request.get_span_context().trace_id, "032x") if declined else None def traffic(exporters: Exporters) -> tuple[str, int]: @@ -197,15 +251,17 @@ async def exporting(base: str): # noqa: ANN201 async with ( streamcast.publish(f"{base}/logs") as log_publication, streamcast.publish(f"{base}/spans") as span_publication, + streamcast.publish(f"{base}/metrics") as metric_publication, ): yield Exporters( logs.StreamLogExporter(log_publication, loop), spans.StreamSpanExporter(span_publication, loop), + metrics.StreamMetricExporter(metric_publication, loop), ) async def run(broker: str) -> None: - """Orders until cancelled, published to `broker`'s two streams.""" + """Orders until cancelled, published to `broker`'s three streams.""" async with exporting(broker) as exporters: stop = threading.Event() # A thread of its own, joined on the way out: the services flush their @@ -222,7 +278,9 @@ async def run(broker: str) -> None: if __name__ == "__main__": parser = argparse.ArgumentParser(description=__doc__.split("\n\n")[0]) parser.add_argument( - "--broker", default="ws://127.0.0.1:8766", help="serving /logs and /spans" + "--broker", + default="ws://127.0.0.1:8766", + help="serving /logs, /spans and /metrics", ) arguments = parser.parse_args() signal.signal(signal.SIGTERM, signal.default_int_handler) diff --git a/tests/test_otel_example.py b/tests/test_otel_example.py index 893c4de..2c05ed3 100644 --- a/tests/test_otel_example.py +++ b/tests/test_otel_example.py @@ -18,7 +18,15 @@ pytest.importorskip("opentelemetry.sdk", reason="the OTel example's dev dependency") import streamcast # noqa: E402 -from examples.otel import analytics, common, demo, export, logs, spans # noqa: E402 +from examples.otel import ( # noqa: E402 + analytics, + common, + demo, + export, + logs, + metrics, + spans, +) class TestTheDemo: @@ -49,6 +57,25 @@ async def test_it_runs_start_to_finish(self, tmp_path): assert seen["errors"] == [{"service": "payments", "errors": 1}] assert seen["lag_ms"] >= 0 + async def test_the_metrics_count_what_happened(self, tmp_path): + """Totals, not row counts: how many rows a run makes depends on how + many one-second collections it spans; what they add up to does not.""" + seen = await demo.main(tmp_path) + + assert seen["points"] > 0 + assert seen["outcomes"] == [ + {"outcome": "confirmed", "orders": 2}, + {"outcome": "failed", "orders": 1}, + ] + assert [(d["service"], d["requests"]) for d in seen["durations"]] == [ + ("checkout", 3), + ("payments", 3), + ] + assert all(d["mean_ms"] > 0 for d in seen["durations"]) + # The one failed order was counted inside its request's span: its + # exemplar is that request's trace. + assert seen["exemplars"] == [seen["failed"]] + async def test_the_failed_request_is_one_trace_across_both_services(self, tmp_path): seen = await demo.main(tmp_path) # Three spans a request: checkout's request, its call, payments' charge. @@ -229,12 +256,18 @@ def test_it_reaches_an_otlp_receiver_as_otel_encodes_it(self): import threading from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter + from opentelemetry.exporter.otlp.proto.http.metric_exporter import ( + OTLPMetricExporter, + ) from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( OTLPSpanExporter, ) from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ( ExportLogsServiceRequest, ) + from opentelemetry.proto.collector.metrics.v1.metrics_service_pb2 import ( + ExportMetricsServiceRequest, + ) from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ( ExportTraceServiceRequest, ) @@ -272,10 +305,16 @@ def log_message(self, format: str, *args: Any) -> None: # noqa: A002, ANN401 span_processor.on_end(export.span(spans.row(original_span))) assert span_processor.force_flush(timeout_millis=10_000) span_processor.shutdown() + + metric_exporter = OTLPMetricExporter(endpoint=f"{receiver}/v1/metrics") + original_metrics, trace = collected() + stored = metrics.rows(original_metrics) + metric_exporter.export(export.metrics_data(stored)) + metric_exporter.shutdown() finally: server.shutdown() - assert sorted(received) == ["/v1/logs", "/v1/traces"] + assert sorted(received) == ["/v1/logs", "/v1/metrics", "/v1/traces"] request = ExportLogsServiceRequest.FromString(received["/v1/logs"]) [resource] = request.resource_logs [scope] = resource.scope_logs @@ -295,6 +334,161 @@ def log_message(self, format: str, *args: Any) -> None: # noqa: A002, ANN401 assert span.status.code == 2 # STATUS_CODE_ERROR assert [e.name for e in span.events] == ["card declined"] + sent = ExportMetricsServiceRequest.FromString(received["/v1/metrics"]) + by_name = {m.name: m for m in sent.resource_metrics[0].scope_metrics[0].metrics} + orders = by_name["shop.orders"].sum + assert orders.aggregation_temporality == 1 # DELTA, as stored + assert orders.is_monotonic + [orders_point] = orders.data_points + assert orders_point.as_int == 2 + [exemplar] = orders_point.exemplars + assert exemplar.trace_id == trace.to_bytes(16, "big") + [duration] = by_name["http.server.request.duration"].histogram.data_points + assert list(duration.bucket_counts) == [1, 1, 0] + + +def collected(): + """Real SDK metrics of every kind, as a reader collects them. + + A counter measured inside a sampled span, so it carries an exemplar; an + up-down counter, a gauge, an explicit-bucket and an exponential histogram. + """ + from opentelemetry.sdk.metrics import MeterProvider + from opentelemetry.sdk.metrics.export import InMemoryMetricReader + from opentelemetry.sdk.metrics.view import ( + ExponentialBucketHistogramAggregation, + View, + ) + from opentelemetry.sdk.resources import Resource + from opentelemetry.sdk.trace import TracerProvider + + reader = InMemoryMetricReader(preferred_temporality=metrics.PREFERRED_TEMPORALITY) + provider = MeterProvider( + [reader], + resource=Resource({"service.name": "checkout"}), + views=[ + View( + instrument_name="shop.amount", + aggregation=ExponentialBucketHistogramAggregation(), + ) + ], + ) + meter = provider.get_meter("shop", "1.0") + with TracerProvider().get_tracer("t").start_as_current_span("request") as span: + meter.create_counter("shop.orders", unit="{order}").add( + 2, {"outcome": "failed", "items": ("book", "pen")} + ) + trace = span.get_span_context().trace_id + + meter.create_up_down_counter("shop.in_flight").add(-3) + meter.create_gauge("shop.temperature", unit="Cel").set(21.5, {"room": "a"}) + duration = meter.create_histogram( + "http.server.request.duration", + unit="s", + explicit_bucket_boundaries_advisory=[0.01, 0.1], + ) + # Two buckets of three filled, unevenly: counts that read the same + # backwards would hide a reversal. + duration.record(0.005) + duration.record(0.05) + amount = meter.create_histogram("shop.amount") + amount.record(0.3) + amount.record(0) + return reader.get_metrics_data(), trace + + +class TestTheMetrics: + """A metrics export as rows, one per data point, and back.""" + + def test_every_kind_survives_the_round_trip_as_otlp_encodes_it(self): + """Compared as OTLP, encoded by OTel's own encoder: what a receiver gets.""" + from opentelemetry.exporter.otlp.proto.common.metrics_encoder import ( + encode_metrics, + ) + + original, _trace = collected() + again = export.metrics_data(metrics.rows(original)) + + assert encode_metrics(again) == encode_metrics(original) + + def test_a_row_is_one_data_point_carrying_its_metric(self): + original, trace = collected() + by_name = {row["name"]: row for row in metrics.rows(original)} + + orders = by_name["shop.orders"] + assert (orders["type"], orders["temporality"], orders["is_monotonic"]) == ( + "sum", + metrics.DELTA, + True, + ) + assert (orders["value_int"], orders["value_double"]) == (2, None) + exemplars = orders["exemplars"] + assert isinstance(exemplars, list) + [exemplar] = exemplars + assert exemplar["trace_id"] == trace.to_bytes(16, "big") + + # The SDK's `delta` preference keeps an up-down counter cumulative. + assert by_name["shop.in_flight"]["temporality"] == metrics.CUMULATIVE + assert by_name["shop.in_flight"]["value_int"] == -3 + gauge = by_name["shop.temperature"] + assert (gauge["type"], gauge["temporality"], gauge["value_double"]) == ( + "gauge", + None, + 21.5, + ) + duration = by_name["http.server.request.duration"] + assert duration["bucket_counts"] == [1, 1, 0] + assert duration["explicit_bounds"] == [0.01, 0.1] + amount = by_name["shop.amount"] + assert (amount["type"], amount["count"], amount["zero_count"]) == ( + "exponential_histogram", + 2, + 1, + ) + + def test_a_non_finite_double_is_stored_as_null(self): + """A stream refuses NaN and ±inf; the point is kept, the value is not.""" + from opentelemetry.sdk.metrics.export import ( + Gauge, + Metric, + MetricsData, + NumberDataPoint, + ResourceMetrics, + ScopeMetrics, + ) + from opentelemetry.sdk.resources import Resource + from opentelemetry.sdk.util.instrumentation import InstrumentationScope + + point = NumberDataPoint({}, None, 1, math.inf) # ty: ignore[invalid-argument-type] + data = MetricsData( + [ + ResourceMetrics( + Resource({}), + [ + ScopeMetrics( + InstrumentationScope("s"), + [Metric("g", "", "", Gauge([point]))], + "", + ) + ], + "", + ) + ] + ) + [row] = metrics.rows(data) + assert (row["value_int"], row["value_double"]) == (None, None) + + def test_the_exporter_asks_for_delta_counters_and_histograms(self): + from opentelemetry.sdk.metrics import Counter, Histogram, UpDownCounter + from opentelemetry.sdk.metrics.export import AggregationTemporality + + exporter = metrics.StreamMetricExporter(None, None) # ty: ignore[invalid-argument-type] + preferred = exporter._preferred_temporality # noqa: SLF001 + assert preferred is not None + assert preferred[Counter] == AggregationTemporality.DELTA + assert preferred[Histogram] == AggregationTemporality.DELTA + assert preferred[UpDownCounter] == AggregationTemporality.CUMULATIVE + def span_row(i: int, service: str, *, kind: int = 2, failed: bool = False) -> dict: """The fields `spans.SCHEMA` requires, plus the service: enough to count.""" From 2eb95c4409d866a7de3a6144e9401a76544b68c0 Mon Sep 17 00:00:00 2001 From: nhobin219 <92616895+nhobin219@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:46:43 +0000 Subject: [PATCH 2/2] fix(examples): run the one-shot OTel demo without a maintainer As the migration demo does: a run this short never fills a log enough to seal it. Its five maintainer processes were still opening the logs when the run ended, so they did nothing but cost time and print tracebacks on SIGTERM. demo.main took 3-11 s with them and about 1 s without, and the OTel test file went from 14.5 s to 4 s. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi --- CHANGELOG.md | 7 +++++++ examples/otel/demo.py | 10 +++++++++- 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b3921c..968935c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,13 @@ there was nothing to have changed from. Everything above it is ordinary. shows them in otel-gui's Metrics tab.** The dashboard moves to otel-gui 3.0.0, the first release that accepts metrics. +### Changed + +- **The one-shot OTel demo runs without a maintainer**, as the migration demo + does. Its maintainer processes were still opening the logs when the run + ended, and cost it several seconds: `demo.main` took 3–11 s with them and + about 1 s without. + ## 0.13.1 — 2026-10-03 ### Changed diff --git a/examples/otel/demo.py b/examples/otel/demo.py index 815dfc8..6f7dfdb 100644 --- a/examples/otel/demo.py +++ b/examples/otel/demo.py @@ -82,8 +82,16 @@ async def one_trace( async def main(root: Path) -> dict[str, Any]: """Run the whole demo under `root`, print it, and return what it saw.""" log_stream, span_stream, metric_stream = streams(root) + # No maintainer, as in the migration demo: a run this short never fills a + # log enough to seal it. Its five processes were still opening the logs + # when the run ended, and cost the run several seconds. `just demo otel`'s + # broker keeps `serve`'s default, which starts them. server = await streamcast.serve( - [log_stream, span_stream, metric_stream], "127.0.0.1", 0, publish=True + [log_stream, span_stream, metric_stream], + "127.0.0.1", + 0, + publish=True, + maintain=False, ) base = f"ws://127.0.0.1:{server.sockets[0].getsockname()[1]}" try: