From e292b278910fd9647bdba096b133a5d19db63db3 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 22 Jul 2026 19:40:53 +0000 Subject: [PATCH 1/2] Add OpenTelemetry metrics monitor (faust[opentelemetry]) Add OpenTelemetryMonitor, a sensor that reports the same metric set as the Statsd and Datadog monitors through the OpenTelemetry metrics API, so Faust metrics can be exported to any OpenTelemetry-compatible backend (OTLP, Prometheus, console, ...). Faust depends only on opentelemetry-api via the new optional faust[opentelemetry] extra; every instrument is a cheap no-op until the application configures a global MeterProvider. Following OpenTelemetry conventions, each instrument is dimensioned by attributes (topic, partition, stream, table, ...) rather than baking those into the metric name as Statsd does. - faust/sensors/otel.py: OTelMetrics instrument container + OpenTelemetryMonitor. - requirements/extras/opentelemetry.txt + setup.py bundle + test requirements. - tests/unit/sensors/test_otel.py: 17 tests driving the hooks against an in-memory SDK metric reader. - docs: user-guide section, installation bundle entry, reference autodoc page. Assisted-by: Claude Opus 4.8 Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01LJ1voGxC8Nqs3AUPKnFVzZ --- docs/includes/installation.txt | 3 + docs/reference/faust.sensors.otel.rst | 11 + docs/reference/index.rst | 1 + docs/userguide/sensors.rst | 49 +++ faust/sensors/otel.py | 479 ++++++++++++++++++++++++++ requirements/extras/opentelemetry.txt | 2 + requirements/test.txt | 1 + setup.py | 1 + tests/unit/sensors/test_otel.py | 245 +++++++++++++ 9 files changed, 792 insertions(+) create mode 100644 docs/reference/faust.sensors.otel.rst create mode 100644 faust/sensors/otel.py create mode 100644 requirements/extras/opentelemetry.txt create mode 100644 tests/unit/sensors/test_otel.py diff --git a/docs/includes/installation.txt b/docs/includes/installation.txt index a44f485d6..16a4a239f 100644 --- a/docs/includes/installation.txt +++ b/docs/includes/installation.txt @@ -81,6 +81,9 @@ Sensors :``faust[prometheus]``: for using the Prometheus Faust monitor. +:``faust[opentelemetry]``: + for using the OpenTelemetry Faust monitor. + :``faust[sentry]``: for reporting worker errors to Sentry via :pypi:`sentry-sdk`. diff --git a/docs/reference/faust.sensors.otel.rst b/docs/reference/faust.sensors.otel.rst new file mode 100644 index 000000000..d237d6954 --- /dev/null +++ b/docs/reference/faust.sensors.otel.rst @@ -0,0 +1,11 @@ +===================================================== + ``faust.sensors.otel`` +===================================================== + +.. contents:: + :local: +.. currentmodule:: faust.sensors.otel + +.. automodule:: faust.sensors.otel + :members: + :undoc-members: diff --git a/docs/reference/index.rst b/docs/reference/index.rst index 979616bda..ef126dae7 100644 --- a/docs/reference/index.rst +++ b/docs/reference/index.rst @@ -105,6 +105,7 @@ Sensors faust.sensors.base faust.sensors.datadog faust.sensors.monitor + faust.sensors.otel faust.sensors.prometheus faust.sensors.statsd diff --git a/docs/userguide/sensors.rst b/docs/userguide/sensors.rst index 55ce68679..8b86bdd2b 100644 --- a/docs/userguide/sensors.rst +++ b/docs/userguide/sensors.rst @@ -72,6 +72,55 @@ in ``app.monitor``: # emit how many events are being processed every second. print(app.monitor.events_s) +.. _sensor-opentelemetry: + +Reporting metrics to OpenTelemetry +================================== + +The :class:`~faust.sensors.otel.OpenTelemetryMonitor` reports the same metrics +as the Statsd and Datadog monitors, but through the `OpenTelemetry`_ metrics +API, so they can be exported to any OpenTelemetry-compatible backend (OTLP, +Prometheus, the console, ...). Install the extra: + +.. sourcecode:: console + + $ pip install "faust[opentelemetry]" + +Faust depends only on ``opentelemetry-api``; every instrument is a cheap no-op +until *your application* configures a global ``MeterProvider`` with the exporter +of your choice. Unlike Statsd -- which bakes dimensions such as topic and +partition into the metric *name* -- the OpenTelemetry monitor follows the +OpenTelemetry conventions and dimensions each instrument by *attributes* +(``topic``, ``partition``, ``stream``, ``table``, ...). + +Configure a ``MeterProvider`` once at startup and set ``app.monitor``: + +.. sourcecode:: python + + import faust + from opentelemetry import metrics + from opentelemetry.sdk.metrics import MeterProvider + from opentelemetry.sdk.metrics.export import ( + ConsoleMetricExporter, + PeriodicExportingMetricReader, + ) + from faust.sensors.otel import OpenTelemetryMonitor + + # Swap ConsoleMetricExporter for an OTLP/Prometheus exporter. + reader = PeriodicExportingMetricReader(ConsoleMetricExporter()) + metrics.set_meter_provider(MeterProvider(metric_readers=[reader])) + + app = faust.App('example', broker='kafka://localhost:9092') + app.monitor = OpenTelemetryMonitor() + +The monitor emits counters (e.g. ``faust.messages.received``, +``faust.events.total``), up/down counters for in-flight work +(``faust.messages.active``, ``faust.events.active``), millisecond latency +histograms (``faust.send.latency``, ``faust.events.runtime``, ...) and gauges +for the latest read/committed/end offsets per topic-partition. + +.. _`OpenTelemetry`: https://opentelemetry.io/ + .. _monitor-reference: Monitor API Reference diff --git a/faust/sensors/otel.py b/faust/sensors/otel.py new file mode 100644 index 000000000..6a5332790 --- /dev/null +++ b/faust/sensors/otel.py @@ -0,0 +1,479 @@ +"""Monitor reporting metrics via OpenTelemetry. + +This sensor mirrors the metric set that :class:`~faust.sensors.statsd.StatsdMonitor` +reports, but records it through the `OpenTelemetry`_ metrics API instead of a +statsd socket. Faust depends only on ``opentelemetry-api``: every instrument is +a cheap no-op until the *application* configures a global ``MeterProvider`` with +the exporter of its choice (OTLP, Prometheus, console, ...), so importing this +module never hard-fails Faust core. + +Unlike statsd -- which encodes dimensions such as topic and partition into the +metric *name* -- the OpenTelemetry conventions favour a single instrument +dimensioned by *attributes*. This sensor follows that convention: e.g. a single +``faust.messages.received`` counter carries ``topic``/``partition`` attributes +rather than one metric per topic. + +.. _OpenTelemetry: https://opentelemetry.io/ +""" + +import typing +from typing import Any, Dict, Optional, cast + +from mode.utils.objects import cached_property + +import faust +from faust import web +from faust.exceptions import ImproperlyConfigured +from faust.types import ( + TP, + AppT, + CollectionT, + EventT, + Message, + RecordMetadata, + StreamT, +) +from faust.types.assignor import PartitionAssignorT +from faust.types.transports import ConsumerT, ProducerT + +from .monitor import Monitor, TPOffsetMapping + +try: + from opentelemetry import metrics as otel_metrics +except ImportError: # pragma: no cover + otel_metrics = None + +if typing.TYPE_CHECKING: # pragma: no cover + from opentelemetry.metrics import ( + Counter, + Gauge, + Histogram, + Meter, + UpDownCounter, + ) +else: + + class Meter: ... # noqa + + class Counter: ... # noqa + + class UpDownCounter: ... # noqa + + class Histogram: ... # noqa + + class Gauge: ... # noqa + + +__all__ = ["OTelMetrics", "OpenTelemetryMonitor"] + +#: Attributes = OpenTelemetry's term for the key/value dimensions on a metric +#: point (statsd bakes these into the metric name instead). +Attributes = Dict[str, Any] + + +class OTelMetrics: + """Container for the OpenTelemetry instruments used by the monitor. + + Instruments are created once, up front, from a single ``Meter``; recording + happens by supplying attributes at call time. Latency histograms use the + ``ms`` unit because Faust observes every latency in milliseconds (see + :meth:`~faust.sensors.monitor.Monitor.ms_since`). + """ + + #: Monotonic counters (``.add``). + messages_received: Counter + events_total: Counter + messages_sent: Counter + messages_send_errors: Counter + table_operations: Counter + assignments: Counter + http_requests: Counter + custom_counts: Counter + + #: Up/down counters for values that rise and fall (``.add`` +/-). + messages_active: UpDownCounter + events_active: UpDownCounter + rebalances_active: UpDownCounter + rebalances_recovering: UpDownCounter + + #: Latency histograms in milliseconds (``.record``). + events_runtime: Histogram + commit_latency: Histogram + send_latency: Histogram + send_error_latency: Histogram + assignment_latency: Histogram + rebalance_return_latency: Histogram + rebalance_end_latency: Histogram + http_latency: Histogram + + #: Synchronous gauges for last-known values (``.set``). + offset_read: Gauge + offset_committed: Gauge + offset_end: Gauge + producer_buffer: Gauge + + def __init__(self, meter: Meter) -> None: + self.meter = meter + + self.messages_received = meter.create_counter( + "faust.messages.received", + unit="1", + description="Messages received from Kafka.", + ) + self.events_total = meter.create_counter( + "faust.events.total", + unit="1", + description="Stream events processed by agents.", + ) + self.messages_sent = meter.create_counter( + "faust.messages.sent", + unit="1", + description="Messages successfully sent to Kafka.", + ) + self.messages_send_errors = meter.create_counter( + "faust.messages.send_errors", + unit="1", + description="Messages that failed to send.", + ) + self.table_operations = meter.create_counter( + "faust.table.operations", + unit="1", + description="Table key operations (get/set/del).", + ) + self.assignments = meter.create_counter( + "faust.assignments", + unit="1", + description="Partition assignor runs (completed/error).", + ) + self.http_requests = meter.create_counter( + "faust.http.requests", + unit="1", + description="Web requests served, by status code.", + ) + self.custom_counts = meter.create_counter( + "faust.custom.count", + unit="1", + description="Application counters recorded via Monitor.count().", + ) + + self.messages_active = meter.create_up_down_counter( + "faust.messages.active", + unit="1", + description="Messages currently being processed.", + ) + self.events_active = meter.create_up_down_counter( + "faust.events.active", + unit="1", + description="Stream events currently being processed.", + ) + self.rebalances_active = meter.create_up_down_counter( + "faust.rebalances.active", + unit="1", + description="Cluster rebalances currently in progress.", + ) + self.rebalances_recovering = meter.create_up_down_counter( + "faust.rebalances.recovering", + unit="1", + description="Rebalances currently in the recovery phase.", + ) + + self.events_runtime = meter.create_histogram( + "faust.events.runtime", + unit="ms", + description="Time spent processing a stream event.", + ) + self.commit_latency = meter.create_histogram( + "faust.commit.latency", + unit="ms", + description="Consumer offset-commit latency.", + ) + self.send_latency = meter.create_histogram( + "faust.send.latency", + unit="ms", + description="Producer send latency (success).", + ) + self.send_error_latency = meter.create_histogram( + "faust.send.error_latency", + unit="ms", + description="Producer send latency (failure).", + ) + self.assignment_latency = meter.create_histogram( + "faust.assignment.latency", + unit="ms", + description="Partition-assignment latency.", + ) + self.rebalance_return_latency = meter.create_histogram( + "faust.rebalance.return_latency", + unit="ms", + description="Time until the consumer returned its assignment.", + ) + self.rebalance_end_latency = meter.create_histogram( + "faust.rebalance.end_latency", + unit="ms", + description="Total rebalance time, including recovery.", + ) + self.http_latency = meter.create_histogram( + "faust.http.latency", + unit="ms", + description="Web request handling latency.", + ) + + self.offset_read = meter.create_gauge( + "faust.offset.read", + unit="1", + description="Last read offset per topic/partition.", + ) + self.offset_committed = meter.create_gauge( + "faust.offset.committed", + unit="1", + description="Last committed offset per topic/partition.", + ) + self.offset_end = meter.create_gauge( + "faust.offset.end", + unit="1", + description="Last known end offset per topic/partition.", + ) + self.producer_buffer = meter.create_gauge( + "faust.producer.buffer", + unit="1", + description="Size of the threaded producer send buffer.", + ) + + +class OpenTelemetryMonitor(Monitor): + """OpenTelemetry Faust Sensor. + + Records the same metrics as :class:`~faust.sensors.statsd.StatsdMonitor` + through the OpenTelemetry metrics API, dimensioned by attributes. + + The application is responsible for configuring a global ``MeterProvider`` + (and an exporter) before starting the app; until then the instruments are + cheap no-ops. A ``meter`` may be supplied explicitly, otherwise one is + resolved from the globally-registered provider. + """ + + def __init__( + self, + *, + meter: Optional[Meter] = None, + service_name: str = "faust", + metrics: Optional[OTelMetrics] = None, + **kwargs: Any, + ) -> None: + if otel_metrics is None: + raise ImproperlyConfigured( + f"{type(self).__name__} requires " + '`pip install "faust[opentelemetry]"`.' + ) + self.service_name = service_name + self._meter = meter + self._metrics = metrics + super().__init__(**kwargs) + + @cached_property + def meter(self) -> Meter: + """Return the OpenTelemetry meter for this monitor.""" + if self._meter is not None: + return self._meter + return otel_metrics.get_meter(self.service_name, faust.__version__) + + @cached_property + def metrics(self) -> OTelMetrics: + """Return the container of OpenTelemetry instruments.""" + if self._metrics is not None: + return self._metrics + return OTelMetrics(self.meter) + + # -- Messages ----------------------------------------------------------- + + def on_message_in(self, tp: TP, offset: int, message: Message) -> None: + """Call before message is delegated to streams.""" + super().on_message_in(tp, offset, message) + attrs = self._tp_attrs(tp) + self.metrics.messages_received.add(1, attrs) + self.metrics.messages_active.add(1, attrs) + self.metrics.offset_read.set(offset, attrs) + + def on_message_out(self, tp: TP, offset: int, message: Message) -> None: + """Call when message is fully acknowledged and can be committed.""" + super().on_message_out(tp, offset, message) + self.metrics.messages_active.add(-1, self._tp_attrs(tp)) + + # -- Stream events ------------------------------------------------------ + + def on_stream_event_in( + self, tp: TP, offset: int, stream: StreamT, event: EventT + ) -> Optional[Dict]: + """Call when stream starts processing an event.""" + state = super().on_stream_event_in(tp, offset, stream, event) + attrs = self._tp_attrs(tp, stream) + self.metrics.events_total.add(1, attrs) + self.metrics.events_active.add(1, attrs) + return state + + def on_stream_event_out( + self, + tp: TP, + offset: int, + stream: StreamT, + event: EventT, + state: Dict = None, + ) -> None: + """Call when stream is done processing an event.""" + super().on_stream_event_out(tp, offset, stream, event, state) + attrs = self._tp_attrs(tp, stream) + self.metrics.events_active.add(-1, attrs) + if state is not None: + self.metrics.events_runtime.record( + self.secs_to_ms(self.events_runtime[-1]), attrs + ) + + # -- Tables ------------------------------------------------------------- + + def on_table_get(self, table: CollectionT, key: Any) -> None: + """Call when value in table is retrieved.""" + super().on_table_get(table, key) + self.metrics.table_operations.add(1, self._table_attrs(table, "get")) + + def on_table_set(self, table: CollectionT, key: Any, value: Any) -> None: + """Call when new value for key in table is set.""" + super().on_table_set(table, key, value) + self.metrics.table_operations.add(1, self._table_attrs(table, "set")) + + def on_table_del(self, table: CollectionT, key: Any) -> None: + """Call when key in a table is deleted.""" + super().on_table_del(table, key) + self.metrics.table_operations.add(1, self._table_attrs(table, "del")) + + # -- Producer ----------------------------------------------------------- + + def on_send_completed( + self, producer: ProducerT, state: Any, metadata: RecordMetadata + ) -> None: + """Call when producer finished sending message.""" + super().on_send_completed(producer, state, metadata) + attrs = {"topic": metadata.topic} if metadata.topic is not None else {} + self.metrics.messages_sent.add(1, attrs) + self.metrics.send_latency.record(self.ms_since(cast(float, state))) + + def on_send_error( + self, producer: ProducerT, exc: BaseException, state: Any + ) -> None: + """Call when producer was unable to publish message.""" + super().on_send_error(producer, exc, state) + self.metrics.messages_send_errors.add(1) + self.metrics.send_error_latency.record(self.ms_since(cast(float, state))) + + # -- Assignments -------------------------------------------------------- + + def on_assignment_error( + self, assignor: PartitionAssignorT, state: Dict, exc: BaseException + ) -> None: + """Partition assignor did not complete assignor due to error.""" + super().on_assignment_error(assignor, state, exc) + attrs = {"result": "error"} + self.metrics.assignments.add(1, attrs) + self.metrics.assignment_latency.record( + self.ms_since(state["time_start"]), attrs + ) + + def on_assignment_completed( + self, assignor: PartitionAssignorT, state: Dict + ) -> None: + """Partition assignor completed assignment.""" + super().on_assignment_completed(assignor, state) + attrs = {"result": "completed"} + self.metrics.assignments.add(1, attrs) + self.metrics.assignment_latency.record( + self.ms_since(state["time_start"]), attrs + ) + + # -- Rebalances --------------------------------------------------------- + + def on_rebalance_start(self, app: AppT) -> Dict: + """Cluster rebalance in progress.""" + state = super().on_rebalance_start(app) + self.metrics.rebalances_active.add(1) + return state + + def on_rebalance_return(self, app: AppT, state: Dict) -> None: + """Consumer replied assignment is done to broker.""" + super().on_rebalance_return(app, state) + self.metrics.rebalances_active.add(-1) + self.metrics.rebalances_recovering.add(1) + self.metrics.rebalance_return_latency.record( + self.ms_since(state["time_return"]) + ) + + def on_rebalance_end(self, app: AppT, state: Dict) -> None: + """Cluster rebalance fully completed (including recovery).""" + super().on_rebalance_end(app, state) + self.metrics.rebalances_recovering.add(-1) + self.metrics.rebalance_end_latency.record(self.ms_since(state["time_end"])) + + # -- Offsets / commit --------------------------------------------------- + + def on_commit_completed(self, consumer: ConsumerT, state: Any) -> None: + """Call when consumer commit offset operation completed.""" + super().on_commit_completed(consumer, state) + self.metrics.commit_latency.record(self.ms_since(cast(float, state))) + + def on_tp_commit(self, tp_offsets: TPOffsetMapping) -> None: + """Call when offset in topic partition is committed.""" + super().on_tp_commit(tp_offsets) + for tp, offset in tp_offsets.items(): + self.metrics.offset_committed.set(offset, self._tp_attrs(tp)) + + def track_tp_end_offset(self, tp: TP, offset: int) -> None: + """Track new topic partition end offset for monitoring lags.""" + super().track_tp_end_offset(tp, offset) + self.metrics.offset_end.set(offset, self._tp_attrs(tp)) + + # -- Web ---------------------------------------------------------------- + + def on_web_request_end( + self, + app: AppT, + request: web.Request, + response: Optional[web.Response], + state: Dict, + *, + view: web.View = None, + ) -> None: + """Web server finished working on request.""" + super().on_web_request_end(app, request, response, state, view=view) + status_code = int(state["status_code"]) + self.metrics.http_requests.add(1, {"status_code": status_code}) + self.metrics.http_latency.record(self.ms_since(state["time_end"])) + + def on_threaded_producer_buffer_processed(self, app: AppT, size: int) -> None: + """Call when the threaded producer flushed its buffer.""" + super().on_threaded_producer_buffer_processed(app, size) + self.metrics.producer_buffer.set(size) + + # -- Custom counters ---------------------------------------------------- + + def count(self, metric_name: str, count: int = 1) -> None: + """Count metric by name.""" + super().count(metric_name, count=count) + self.metrics.custom_counts.add(count, {"name": metric_name}) + + # -- Attribute builders ------------------------------------------------- + + def _tp_attrs(self, tp: TP, stream: Optional[StreamT] = None) -> Attributes: + attrs: Attributes = {"topic": tp.topic, "partition": tp.partition} + if stream is not None: + attrs["stream"] = self._stream_label(stream) + return attrs + + def _table_attrs(self, table: CollectionT, operation: str) -> Attributes: + return {"table": table.name, "operation": operation} + + def _stream_label(self, stream: StreamT) -> str: + return ( + self._normalize( + stream.shortlabel.lstrip("Stream:"), + ) + .strip("_") + .lower() + ) diff --git a/requirements/extras/opentelemetry.txt b/requirements/extras/opentelemetry.txt new file mode 100644 index 000000000..95ceff058 --- /dev/null +++ b/requirements/extras/opentelemetry.txt @@ -0,0 +1,2 @@ +opentelemetry-api>=1.23.0 +opentelemetry-sdk>=1.23.0 diff --git a/requirements/test.txt b/requirements/test.txt index b06081cb4..59476a837 100644 --- a/requirements/test.txt +++ b/requirements/test.txt @@ -28,6 +28,7 @@ wheel intervaltree -r requirements.txt -r extras/datadog.txt +-r extras/opentelemetry.txt -r extras/opentracing.txt -r extras/redis.txt -r extras/statsd.txt diff --git a/setup.py b/setup.py index 7a3190b70..bf5909c84 100644 --- a/setup.py +++ b/setup.py @@ -30,6 +30,7 @@ "datadog", "debug", "fast", + "opentelemetry", "opentracing", "orjson", "prometheus", diff --git a/tests/unit/sensors/test_otel.py b/tests/unit/sensors/test_otel.py new file mode 100644 index 000000000..43017f384 --- /dev/null +++ b/tests/unit/sensors/test_otel.py @@ -0,0 +1,245 @@ +from typing import Any, Dict, List, Optional +from unittest.mock import Mock + +import pytest + +from faust.exceptions import ImproperlyConfigured +from faust.sensors.otel import OpenTelemetryMonitor, OTelMetrics +from faust.types import TP + +opentelemetry = pytest.importorskip("opentelemetry") + +from opentelemetry.sdk.metrics import MeterProvider # noqa: E402 +from opentelemetry.sdk.metrics.export import InMemoryMetricReader # noqa: E402 + +TP1 = TP("foo", 3) +TP2 = TP("bar", 4) + + +def _time() -> float: + return 101.1 + + +class Collected: + """Flattened view of the metrics a reader collected.""" + + def __init__(self, reader: InMemoryMetricReader) -> None: + self.points: Dict[str, List[Any]] = {} + data = reader.get_metrics_data() + if data is None: + return + for rm in data.resource_metrics: + for sm in rm.scope_metrics: + self.scope = sm.scope + for metric in sm.metrics: + self.points.setdefault(metric.name, []) + for point in metric.data.data_points: + self.points[metric.name].append(point) + + def names(self) -> List[str]: + return sorted(self.points) + + def point(self, name: str, attributes: Optional[Dict] = None) -> Any: + candidates = self.points.get(name, []) + if attributes is None: + assert len(candidates) == 1, f"{name}: {candidates}" + return candidates[0] + for point in candidates: + if dict(point.attributes) == attributes: + return point + raise AssertionError( + f"no point for {name} with attributes {attributes}: " + f"{[dict(p.attributes) for p in candidates]}" + ) + + def value(self, name: str, attributes: Optional[Dict] = None) -> Any: + return self.point(name, attributes).value + + def sum(self, name: str, attributes: Optional[Dict] = None) -> Any: + return self.point(name, attributes).sum + + def count(self, name: str, attributes: Optional[Dict] = None) -> int: + return self.point(name, attributes).count + + +class TestOpenTelemetryMonitor: + @pytest.fixture() + def reader(self) -> InMemoryMetricReader: + return InMemoryMetricReader() + + @pytest.fixture() + def meter(self, reader: InMemoryMetricReader): + # An isolated provider per test -- never touch the global provider. + provider = MeterProvider(metric_readers=[reader]) + return provider.get_meter("faust-test") + + @pytest.fixture() + def mon(self, meter) -> OpenTelemetryMonitor: + return OpenTelemetryMonitor(meter=meter, time=_time) + + @pytest.fixture() + def stream(self): + stream = Mock(name="stream") + stream.shortlabel = "Stream: Topic: foo" + return stream + + @pytest.fixture() + def event(self): + return Mock(name="event") + + @pytest.fixture() + def table(self): + table = Mock(name="table") + table.name = "table1" + return table + + def test_raises_if_opentelemetry_not_installed(self, *, monkeypatch): + monkeypatch.setattr("faust.sensors.otel.otel_metrics", None) + with pytest.raises(ImproperlyConfigured): + OpenTelemetryMonitor() + + def test_resolves_global_meter_by_default(self, *, monkeypatch): + get_meter = Mock(name="get_meter") + monkeypatch.setattr("faust.sensors.otel.otel_metrics.get_meter", get_meter) + mon = OpenTelemetryMonitor(service_name="svc") + assert mon.meter is get_meter.return_value + get_meter.assert_called_once() + assert get_meter.call_args.args[0] == "svc" + + def test_metrics_container_built_lazily(self, *, mon): + assert isinstance(mon.metrics, OTelMetrics) + + def test_on_message_in(self, *, mon, reader): + mon.on_message_in(TP1, 400, Mock(name="message")) + + c = Collected(reader) + attrs = {"topic": "foo", "partition": 3} + assert c.value("faust.messages.received", attrs) == 1 + assert c.value("faust.messages.active", attrs) == 1 + assert c.value("faust.offset.read", attrs) == 400 + + def test_on_message_out_decrements_active(self, *, mon, reader): + mon.on_message_in(TP1, 400, Mock(name="message", time_in=100.0)) + mon.on_message_out(TP1, 400, Mock(name="message", time_in=100.0)) + + c = Collected(reader) + attrs = {"topic": "foo", "partition": 3} + assert c.value("faust.messages.active", attrs) == 0 + + def test_on_stream_event_in_out(self, *, mon, reader, stream, event): + state = mon.on_stream_event_in(TP1, 401, stream, event) + attrs = {"topic": "foo", "partition": 3, "stream": "topic_foo"} + + c = Collected(reader) + assert c.value("faust.events.total", attrs) == 1 + assert c.value("faust.events.active", attrs) == 1 + + mon.on_stream_event_out(TP1, 401, stream, event, state) + c = Collected(reader) + assert c.value("faust.events.active", attrs) == 0 + assert c.count("faust.events.runtime", attrs) == 1 + + def test_on_stream_event_out_no_runtime_without_state( + self, *, mon, reader, stream, event + ): + mon.on_stream_event_out(TP1, 401, stream, event, state=None) + c = Collected(reader) + assert "faust.events.runtime" not in c.names() + + def test_table_operations(self, *, mon, reader, table): + mon.on_table_get(table, "k") + mon.on_table_set(table, "k", "v") + mon.on_table_del(table, "k") + + c = Collected(reader) + base = {"table": "table1"} + assert c.value("faust.table.operations", {**base, "operation": "get"}) == 1 + assert c.value("faust.table.operations", {**base, "operation": "set"}) == 1 + assert c.value("faust.table.operations", {**base, "operation": "del"}) == 1 + + def test_on_send_completed(self, *, mon, reader): + metadata = Mock(name="metadata") + metadata.topic = "foo" + mon.on_send_completed(Mock(name="producer"), 100.1, metadata) + + c = Collected(reader) + assert c.value("faust.messages.sent", {"topic": "foo"}) == 1 + assert c.count("faust.send.latency") == 1 + # ms_since(100.1) with time()==101.1 -> 1000ms + assert c.sum("faust.send.latency") == pytest.approx(1000.0) + + def test_on_send_error(self, *, mon, reader): + mon.on_send_error(Mock(name="producer"), RuntimeError("boom"), 100.1) + + c = Collected(reader) + assert c.value("faust.messages.send_errors") == 1 + assert c.count("faust.send.error_latency") == 1 + + def test_assignments(self, *, mon, reader): + assignor = Mock(name="assignor") + mon.on_assignment_completed(assignor, {"time_start": 100.1}) + mon.on_assignment_error(assignor, {"time_start": 100.1}, RuntimeError()) + + c = Collected(reader) + assert c.value("faust.assignments", {"result": "completed"}) == 1 + assert c.value("faust.assignments", {"result": "error"}) == 1 + assert c.count("faust.assignment.latency", {"result": "completed"}) == 1 + assert c.count("faust.assignment.latency", {"result": "error"}) == 1 + + def test_rebalance_lifecycle(self, *, mon, reader): + app = Mock(name="app") + state = mon.on_rebalance_start(app) + c = Collected(reader) + assert c.value("faust.rebalances.active") == 1 + + mon.on_rebalance_return(app, state) + c = Collected(reader) + assert c.value("faust.rebalances.active") == 0 + assert c.value("faust.rebalances.recovering") == 1 + assert c.count("faust.rebalance.return_latency") == 1 + + mon.on_rebalance_end(app, state) + c = Collected(reader) + assert c.value("faust.rebalances.recovering") == 0 + assert c.count("faust.rebalance.end_latency") == 1 + + def test_on_commit_completed(self, *, mon, reader): + mon.on_commit_completed(Mock(name="consumer"), 100.1) + c = Collected(reader) + assert c.count("faust.commit.latency") == 1 + assert c.sum("faust.commit.latency") == pytest.approx(1000.0) + + def test_offsets(self, *, mon, reader): + mon.on_tp_commit({TP1: 10, TP2: 20}) + mon.track_tp_end_offset(TP1, 99) + + c = Collected(reader) + assert c.value("faust.offset.committed", {"topic": "foo", "partition": 3}) == 10 + assert c.value("faust.offset.committed", {"topic": "bar", "partition": 4}) == 20 + assert c.value("faust.offset.end", {"topic": "foo", "partition": 3}) == 99 + + def test_on_web_request_end(self, *, mon, reader): + response = Mock(name="response") + response.status = 200 + state = {"time_start": 100.1} + mon.on_web_request_end( + Mock(name="app"), + Mock(name="request"), + response, + state, + view=Mock(name="view"), + ) + c = Collected(reader) + assert c.value("faust.http.requests", {"status_code": 200}) == 1 + assert c.count("faust.http.latency") == 1 + + def test_threaded_producer_buffer(self, *, mon, reader): + mon.on_threaded_producer_buffer_processed(Mock(name="app"), 17) + c = Collected(reader) + assert c.value("faust.producer.buffer") == 17 + + def test_count(self, *, mon, reader): + mon.count("my_metric", 5) + mon.count("my_metric", 3) + c = Collected(reader) + assert c.value("faust.custom.count", {"name": "my_metric"}) == 8 From 20e563ef2d058be8d41bdd39dd6461ef849581b7 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 23 Jul 2026 03:28:24 +0000 Subject: [PATCH 2/2] test(otel): cover injected OTelMetrics container path Adds a test that passes a pre-built OTelMetrics container via the `metrics=` kwarg, covering the last uncovered line/branch in OpenTelemetryMonitor.metrics (patch coverage 98.9% -> 100%). Assisted-by: Claude Opus 4.8 Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01LJ1voGxC8Nqs3AUPKnFVzZ --- tests/unit/sensors/test_otel.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/tests/unit/sensors/test_otel.py b/tests/unit/sensors/test_otel.py index 43017f384..f6e600869 100644 --- a/tests/unit/sensors/test_otel.py +++ b/tests/unit/sensors/test_otel.py @@ -109,6 +109,11 @@ def test_resolves_global_meter_by_default(self, *, monkeypatch): def test_metrics_container_built_lazily(self, *, mon): assert isinstance(mon.metrics, OTelMetrics) + def test_metrics_container_can_be_injected(self, *, meter): + metrics = OTelMetrics(meter) + mon = OpenTelemetryMonitor(meter=meter, metrics=metrics, time=_time) + assert mon.metrics is metrics + def test_on_message_in(self, *, mon, reader): mon.on_message_in(TP1, 400, Mock(name="message"))