From d116d9b0868e1d29b94c226ef49bf3587be55b99 Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 00:03:39 +0500 Subject: [PATCH 1/8] feat(otel): inert event trace points in reflex-base and the reflex-otel instrumentor --- .github/scripts/dispatch_release/detect.sh | 3 +- .github/workflows/dispatch_release.yml | 5 + .github/workflows/publish.yml | 1 + packages/reflex-base/news/6227.feature.md | 1 + packages/reflex-base/pyproject.toml | 1 + .../src/reflex_base/event/context.py | 6 + .../event/processor/event_processor.py | 11 +- packages/reflex-base/src/reflex_base/otel.py | 108 +++++++++++++ packages/reflex-otel/CHANGELOG.md | 1 + packages/reflex-otel/README.md | 13 ++ packages/reflex-otel/news/6227.feature.md | 1 + packages/reflex-otel/pyproject.toml | 31 ++++ .../reflex-otel/src/reflex_otel/__init__.py | 47 ++++++ pyproject.toml | 4 + tests/units/reflex_base/conftest.py | 26 ++++ .../event/processor/test_event_processor.py | 43 ++++++ tests/units/reflex_base/test_otel.py | 95 ++++++++++++ tests/units/reflex_otel/__init__.py | 0 tests/units/reflex_otel/test_init.py | 35 +++++ uv.lock | 146 ++++++++++++++---- 20 files changed, 545 insertions(+), 33 deletions(-) create mode 100644 packages/reflex-base/news/6227.feature.md create mode 100644 packages/reflex-base/src/reflex_base/otel.py create mode 100644 packages/reflex-otel/CHANGELOG.md create mode 100644 packages/reflex-otel/README.md create mode 100644 packages/reflex-otel/news/6227.feature.md create mode 100644 packages/reflex-otel/pyproject.toml create mode 100644 packages/reflex-otel/src/reflex_otel/__init__.py create mode 100644 tests/units/reflex_base/conftest.py create mode 100644 tests/units/reflex_base/test_otel.py create mode 100644 tests/units/reflex_otel/__init__.py create mode 100644 tests/units/reflex_otel/test_init.py diff --git a/.github/scripts/dispatch_release/detect.sh b/.github/scripts/dispatch_release/detect.sh index 9a57a7d48dc..4199c0f85b8 100755 --- a/.github/scripts/dispatch_release/detect.sh +++ b/.github/scripts/dispatch_release/detect.sh @@ -18,9 +18,10 @@ declare -A MAP=( [reflex_components_sonner]=reflex-components-sonner [reflex_docgen]=reflex-docgen [reflex_hosting_cli]=reflex-hosting-cli + [reflex_otel]=reflex-otel [reflex_release]=reflex-release ) -ORDER=(hatch_reflex_pyi reflex_base reflex_components_code reflex_components_core reflex_components_dataeditor reflex_components_gridjs reflex_components_lucide reflex_components_markdown reflex_components_moment reflex_components_plotly reflex_components_radix reflex_components_react_player reflex_components_recharts reflex_components_sonner reflex_docgen reflex_hosting_cli reflex_release) +ORDER=(hatch_reflex_pyi reflex_base reflex_components_code reflex_components_core reflex_components_dataeditor reflex_components_gridjs reflex_components_lucide reflex_components_markdown reflex_components_moment reflex_components_plotly reflex_components_radix reflex_components_react_player reflex_components_recharts reflex_components_sonner reflex_docgen reflex_hosting_cli reflex_otel reflex_release) PACKAGES=() for key in "${ORDER[@]}"; do diff --git a/.github/workflows/dispatch_release.yml b/.github/workflows/dispatch_release.yml index f73008151d8..b48cfd3d139 100644 --- a/.github/workflows/dispatch_release.yml +++ b/.github/workflows/dispatch_release.yml @@ -129,6 +129,10 @@ on: description: "reflex-hosting-cli" type: boolean default: false + reflex_otel: + description: "reflex-otel" + type: boolean + default: false reflex_release: description: "reflex-release" type: boolean @@ -171,6 +175,7 @@ jobs: reflex_docgen: ${{ inputs.reflex_docgen }} reflex_hosting_cli: ${{ inputs.reflex_hosting_cli }} reflex_release: ${{ inputs.reflex_release }} + reflex_otel: ${{ inputs.reflex_otel }} run: bash .github/scripts/dispatch_release/detect.sh materialize: diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml index 6c9528bd7b5..69ff909ccbc 100644 --- a/.github/workflows/publish.yml +++ b/.github/workflows/publish.yml @@ -68,6 +68,7 @@ on: - reflex-components-sonner - reflex-docgen - reflex-hosting-cli + - reflex-otel - reflex-release - reflex-site-shared version: diff --git a/packages/reflex-base/news/6227.feature.md b/packages/reflex-base/news/6227.feature.md new file mode 100644 index 00000000000..86c31cc08c5 --- /dev/null +++ b/packages/reflex-base/news/6227.feature.md @@ -0,0 +1 @@ +Add inert OpenTelemetry trace points around event handler execution (`reflex_base.otel`); they cost one boolean check until the `reflex-otel` package enables them. diff --git a/packages/reflex-base/pyproject.toml b/packages/reflex-base/pyproject.toml index 7e744790a18..68bec385afb 100644 --- a/packages/reflex-base/pyproject.toml +++ b/packages/reflex-base/pyproject.toml @@ -8,6 +8,7 @@ authors = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] maintainers = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] requires-python = ">=3.10" dependencies = [ + "opentelemetry-api >=1.24.0,<2.0", "packaging >=24.2,<27", "rich >=13,<16", "typing_extensions >=4.13.0", diff --git a/packages/reflex-base/src/reflex_base/event/context.py b/packages/reflex-base/src/reflex_base/event/context.py index df3f0200e62..fba066556d0 100644 --- a/packages/reflex-base/src/reflex_base/event/context.py +++ b/packages/reflex-base/src/reflex_base/event/context.py @@ -8,10 +8,13 @@ from collections.abc import Callable, Mapping from typing import TYPE_CHECKING, Any, Protocol +from reflex_base import otel from reflex_base.context.base import BaseContext from reflex_base.utils.format import to_snake_case if TYPE_CHECKING: + from opentelemetry.context import Context + from reflex.istate.manager import StateManager from reflex_base.event import Event @@ -98,6 +101,8 @@ class EventContext(BaseContext): cached_states: dict[type, Any] = dataclasses.field( default_factory=dict, init=False, repr=False ) + # OpenTelemetry context active when this event was enqueued (None when tracing is off). + otel_context: Context | None = dataclasses.field(default=None, repr=False) def fork(self, token: str | None = None) -> EventContext: """Return a new EventContext with the specified fields replaced. @@ -115,6 +120,7 @@ def fork(self, token: str | None = None) -> EventContext: enqueue_impl=self.enqueue_impl, emit_delta_impl=self.emit_delta_impl, emit_event_impl=self.emit_event_impl, + otel_context=otel.capture_context(), ) async def emit_delta(self, delta: Mapping[str, Mapping[str, Any]]) -> None: diff --git a/packages/reflex-base/src/reflex_base/event/processor/event_processor.py b/packages/reflex-base/src/reflex_base/event/processor/event_processor.py index 29e35ed2367..15e2b58d299 100644 --- a/packages/reflex-base/src/reflex_base/event/processor/event_processor.py +++ b/packages/reflex-base/src/reflex_base/event/processor/event_processor.py @@ -20,6 +20,7 @@ from reflex.app_mixins.middleware import MiddlewareMixin from reflex.istate.manager import StateManager from reflex.utils import console +from reflex_base import otel from reflex_base.event.context import EventContext from reflex_base.event.processor.future import EventFuture from reflex_base.event.processor.timeout import DrainTimeoutManager @@ -576,7 +577,15 @@ async def _process_event_queue_entry( """ # Set up the event context for this task. EventContext.set(entry.ctx) - await self._execute_event(entry=entry, registered_handler=registered_handler) + if not otel.enabled: + await self._execute_event( + entry=entry, registered_handler=registered_handler + ) + return + with otel.event_span(entry.event, entry.ctx, registered_handler): + await self._execute_event( + entry=entry, registered_handler=registered_handler + ) def _create_event_task( self, diff --git a/packages/reflex-base/src/reflex_base/otel.py b/packages/reflex-base/src/reflex_base/otel.py new file mode 100644 index 00000000000..074a4df52da --- /dev/null +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -0,0 +1,108 @@ +"""OpenTelemetry trace points for the Reflex runtime. + +The framework calls into this module at a small number of fixed points (event +dispatch, event context forks). Every entry point checks the module-level +``enabled`` flag first, so with no instrumentation installed the cost is one +attribute read and no ``opentelemetry`` object is ever created. + +The ``reflex-otel`` package flips the flag via :func:`enable` once a tracer +provider is available. +""" + +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager +from typing import TYPE_CHECKING + +from opentelemetry import context as otel_context +from opentelemetry import trace +from opentelemetry.trace import SpanKind + +from reflex_base.constants.base import Reflex + +if TYPE_CHECKING: + from opentelemetry.context import Context + + from reflex_base.event import Event + from reflex_base.event.context import EventContext + from reflex_base.registry import RegisteredEventHandler + +INSTRUMENTATION_NAME = "reflex" + +# Attribute keys emitted on Reflex spans. +ATTR_EVENT_NAME = "reflex.event.name" +ATTR_EVENT_TXID = "reflex.event.txid" +ATTR_EVENT_PARENT_TXID = "reflex.event.parent_txid" +ATTR_EVENT_BACKGROUND = "reflex.event.background" +ATTR_SESSION_ID = "session.id" +ATTR_CODE_FUNCTION_NAME = "code.function.name" + +# Read at every trace point; True only after enable() ran. +enabled: bool = False +_tracer: trace.Tracer = trace.NoOpTracer() + + +def enable(tracer_provider: trace.TracerProvider | None = None) -> None: + """Turn the trace points on. + + Args: + tracer_provider: The provider to obtain the tracer from. Defaults to + the global provider, which may be configured later; the returned + proxy tracer picks it up automatically. + """ + global _tracer, enabled + _tracer = trace.get_tracer( + INSTRUMENTATION_NAME, Reflex.VERSION, tracer_provider=tracer_provider + ) + enabled = True + + +def disable() -> None: + """Turn the trace points off and drop the tracer.""" + global _tracer, enabled + enabled = False + _tracer = trace.NoOpTracer() + + +def capture_context() -> Context | None: + """Snapshot the current OpenTelemetry context for a forked event context. + + Returns: + The current context when tracing is enabled, otherwise None. + """ + return otel_context.get_current() if enabled else None + + +@contextmanager +def event_span( + event: Event, ctx: EventContext, registered_handler: RegisteredEventHandler +) -> Iterator[trace.Span]: + """Open the span for one event handler execution. + + Chained events are parented under the span that enqueued them; events + that arrive from the frontend start a new trace. + + Args: + event: The event being processed. + ctx: The event context for this execution. + registered_handler: The handler resolved for the event. + + Yields: + The active span. + """ + handler = registered_handler.handler + with _tracer.start_as_current_span( + event.name, + context=ctx.otel_context, + kind=SpanKind.SERVER, + attributes={ + ATTR_EVENT_NAME: event.name, + ATTR_EVENT_TXID: ctx.txid, + ATTR_EVENT_BACKGROUND: handler.is_background, + ATTR_SESSION_ID: ctx.token, + ATTR_CODE_FUNCTION_NAME: getattr(handler.fn, "__qualname__", event.name), + } + | ({ATTR_EVENT_PARENT_TXID: ctx.parent_txid} if ctx.parent_txid else {}), + ) as span: + yield span diff --git a/packages/reflex-otel/CHANGELOG.md b/packages/reflex-otel/CHANGELOG.md new file mode 100644 index 00000000000..825c32f0d03 --- /dev/null +++ b/packages/reflex-otel/CHANGELOG.md @@ -0,0 +1 @@ +# Changelog diff --git a/packages/reflex-otel/README.md b/packages/reflex-otel/README.md new file mode 100644 index 00000000000..298f1b33a71 --- /dev/null +++ b/packages/reflex-otel/README.md @@ -0,0 +1,13 @@ +# reflex-otel + +OpenTelemetry instrumentation for the Reflex framework. + +```python +from reflex_otel import ReflexInstrumentor + +ReflexInstrumentor().instrument() +``` + +The package registers an `opentelemetry_instrumentor` entry point, so +`opentelemetry-instrument reflex run` enables it automatically. Configure a +tracer provider (for example with `opentelemetry-sdk`) to export the spans. diff --git a/packages/reflex-otel/news/6227.feature.md b/packages/reflex-otel/news/6227.feature.md new file mode 100644 index 00000000000..19fa9421ea3 --- /dev/null +++ b/packages/reflex-otel/news/6227.feature.md @@ -0,0 +1 @@ +Add the `reflex-otel` package: an OpenTelemetry instrumentor that turns on the framework's built-in trace points (one span per event handler run, chained events parented under the enqueuing span). diff --git a/packages/reflex-otel/pyproject.toml b/packages/reflex-otel/pyproject.toml new file mode 100644 index 00000000000..1eb2a10c36a --- /dev/null +++ b/packages/reflex-otel/pyproject.toml @@ -0,0 +1,31 @@ +[project] +name = "reflex-otel" +dynamic = ["version"] +description = "OpenTelemetry instrumentation for the Reflex framework." +license.text = "Apache-2.0" +readme = "README.md" +authors = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] +maintainers = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] +requires-python = ">=3.10" +dependencies = [ + "opentelemetry-api >=1.24.0,<2.0", + "opentelemetry-instrumentation >=0.45b0,<1.0", + "reflex-base >= 0.9.7.post45.dev0", +] + +[project.optional-dependencies] +instruments = ["reflex-base >= 0.9.7.post45.dev0"] + +[project.entry-points.opentelemetry_instrumentor] +reflex = "reflex_otel:ReflexInstrumentor" + +[tool.hatch.version] +source = "uv-dynamic-versioning" + +[tool.uv-dynamic-versioning] +pattern-prefix = "reflex-otel-" +fallback-version = "0.0.0dev0" + +[build-system] +requires = ["hatchling", "uv-dynamic-versioning"] +build-backend = "hatchling.build" diff --git a/packages/reflex-otel/src/reflex_otel/__init__.py b/packages/reflex-otel/src/reflex_otel/__init__.py new file mode 100644 index 00000000000..bbc530ad431 --- /dev/null +++ b/packages/reflex-otel/src/reflex_otel/__init__.py @@ -0,0 +1,47 @@ +"""OpenTelemetry instrumentation for the Reflex framework.""" + +from __future__ import annotations + +from collections.abc import Collection +from typing import Any + +from opentelemetry.instrumentation.instrumentor import BaseInstrumentor +from reflex_base import otel + +_instruments = ("reflex-base >= 0.9.7.post45.dev0",) + + +class ReflexInstrumentor(BaseInstrumentor): + """Enable the trace points built into the Reflex runtime. + + Usage:: + + ReflexInstrumentor().instrument(tracer_provider=provider) + + or let ``opentelemetry-instrument`` load it through the + ``opentelemetry_instrumentor`` entry point. + """ + + def instrumentation_dependencies(self) -> Collection[str]: + """Return the packages this instrumentor targets. + + Returns: + The dependency specifiers for the instrumented package. + """ + return _instruments + + def _instrument(self, **kwargs: Any) -> None: + """Turn on the Reflex trace points. + + Args: + **kwargs: ``tracer_provider`` selects the provider; defaults to the global one. + """ + otel.enable(tracer_provider=kwargs.get("tracer_provider")) + + def _uninstrument(self, **kwargs: Any) -> None: + """Turn off the Reflex trace points. + + Args: + **kwargs: Ignored. + """ + otel.disable() diff --git a/pyproject.toml b/pyproject.toml index 4a3759b3775..7abfac5bfbd 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -86,6 +86,7 @@ dev = [ "hatchling", "libsass", "numpy", + "opentelemetry-sdk", "pandas", "pillow", "playwright", @@ -107,6 +108,7 @@ dev = [ "python-dotenv", "pyyaml", "reflex-docgen", + "reflex-otel", "reflex-release", "reflex-site-shared", "ruff", @@ -159,6 +161,7 @@ extraPaths = [ "packages/reflex-recharts/src", "packages/reflex-sonner/src", "packages/reflex-components-internal/src", + "packages/reflex-otel/src", "packages/reflex-site-shared/src", "packages/integrations-docs/src", ] @@ -427,6 +430,7 @@ reflex-components-sonner.workspace = true reflex-docgen.workspace = true reflex-release.workspace = true reflex-hosting-cli.workspace = true +reflex-otel.workspace = true reflex-components-internal.workspace = true reflex-site-shared.workspace = true reflex-integrations-docs.workspace = true diff --git a/tests/units/reflex_base/conftest.py b/tests/units/reflex_base/conftest.py new file mode 100644 index 00000000000..5ae9153f10e --- /dev/null +++ b/tests/units/reflex_base/conftest.py @@ -0,0 +1,26 @@ +"""Shared fixtures for reflex_base unit tests.""" + +from collections.abc import Generator + +import pytest +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from reflex_base import otel + + +@pytest.fixture +def otel_exporter() -> Generator[InMemorySpanExporter, None, None]: + """Enable the reflex_base.otel trace points against an in-memory exporter. + + Yields: + The exporter collecting finished spans. + """ + exporter = InMemorySpanExporter() + provider = TracerProvider() + provider.add_span_processor(SimpleSpanProcessor(exporter)) + otel.enable(tracer_provider=provider) + try: + yield exporter + finally: + otel.disable() diff --git a/tests/units/reflex_base/event/processor/test_event_processor.py b/tests/units/reflex_base/event/processor/test_event_processor.py index d5dda19dca3..fba1e5d7333 100644 --- a/tests/units/reflex_base/event/processor/test_event_processor.py +++ b/tests/units/reflex_base/event/processor/test_event_processor.py @@ -792,3 +792,46 @@ async def _watcher(): # noqa: RUF029 collected = [v async for v in _stream_queue_until_done(queue, _watcher())] assert collected == [99] + + +async def test_no_spans_when_otel_disabled( + mock_event_processor: EventProcessor, token: str +): + """With tracing off the processor never touches the tracer. + + Args: + mock_event_processor: The event processor with mock root context. + token: The client token. + """ + from reflex_base import otel + + assert otel.enabled is False + async with mock_event_processor as ep: + await ep.enqueue(token, Event.from_event_type(noop_event())[0]) + + +async def test_event_spans_chain_parent_child(token: str, otel_exporter): + """Each event gets a span; chained events are children of the enqueuing span. + + Args: + token: The client token. + otel_exporter: In-memory span exporter with tracing enabled. + """ + from reflex_base import otel + + ep = EventProcessor(graceful_shutdown_timeout=2) + ep.configure() + async with ep: + await ep.enqueue(token, Event.from_event_type(chaining_event())[0]) + assert _CALL_LOG == [{"value": "chained"}] + spans = {s.name.rsplit(".", 1)[-1]: s for s in otel_exporter.get_finished_spans()} + parent = spans["_chaining_handler"] + child = spans["_logging_handler"] + assert parent.parent is None + assert child.parent is not None + assert child.parent.span_id == parent.context.span_id + assert ( + child.attributes[otel.ATTR_EVENT_PARENT_TXID] + == parent.attributes[otel.ATTR_EVENT_TXID] + ) + assert child.attributes[otel.ATTR_SESSION_ID] == token diff --git a/tests/units/reflex_base/test_otel.py b/tests/units/reflex_base/test_otel.py new file mode 100644 index 00000000000..93674e1269d --- /dev/null +++ b/tests/units/reflex_base/test_otel.py @@ -0,0 +1,95 @@ +"""Tests for the reflex_base.otel trace points.""" + +import pytest +from opentelemetry import context as otel_context +from opentelemetry import trace +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from opentelemetry.trace import SpanKind, StatusCode +from reflex_base import otel +from reflex_base.event.context import EventContext +from reflex_base.registry import RegisteredEventHandler + +from reflex.event import Event, EventHandler + + +def _ctx(token: str = "tok", parent_txid: str | None = None) -> EventContext: + return EventContext( + token=token, + state_manager=None, # type: ignore[arg-type] + enqueue_impl=None, # type: ignore[arg-type] + parent_txid=parent_txid, + ) + + +async def _handler(): + """A no-op handler.""" + + +def test_disabled_by_default(): + assert otel.enabled is False + assert otel.capture_context() is None + + +def test_enable_disable_toggle(): + otel.enable() + assert otel.enabled is True + assert otel.capture_context() is not None + otel.disable() + assert otel.enabled is False + assert isinstance(otel._tracer, trace.NoOpTracer) + + +def test_capture_context_returns_current(otel_exporter: InMemorySpanExporter): + with otel._tracer.start_as_current_span("outer") as span: + captured = otel.capture_context() + assert captured is not None + assert trace.get_current_span(captured) is span + assert otel_context.get_current() is not captured + + +def test_event_span_attributes(otel_exporter: InMemorySpanExporter): + ctx = _ctx(parent_txid="parent123") + event = Event(name="state.sub.handler") + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + with otel.event_span(event, ctx, registered) as span: + assert trace.get_current_span() is span + (finished,) = otel_exporter.get_finished_spans() + assert finished.name == "state.sub.handler" + assert finished.kind == SpanKind.SERVER + assert finished.parent is None + assert finished.attributes == { + otel.ATTR_EVENT_NAME: "state.sub.handler", + otel.ATTR_EVENT_TXID: ctx.txid, + otel.ATTR_EVENT_PARENT_TXID: "parent123", + otel.ATTR_EVENT_BACKGROUND: False, + otel.ATTR_SESSION_ID: "tok", + otel.ATTR_CODE_FUNCTION_NAME: "_handler", + } + assert finished.status.status_code == StatusCode.UNSET + + +def test_event_span_records_exception(otel_exporter: InMemorySpanExporter): + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + with ( + pytest.raises(RuntimeError, match="boom"), + otel.event_span(Event(name="e"), _ctx(), registered), + ): + msg = "boom" + raise RuntimeError(msg) + (finished,) = otel_exporter.get_finished_spans() + assert finished.status.status_code == StatusCode.ERROR + assert finished.events[0].name == "exception" + + +def test_event_span_parents_under_captured_context( + otel_exporter: InMemorySpanExporter, +): + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + with otel._tracer.start_as_current_span("root") as root: + child_ctx = _ctx().fork() + assert child_ctx.otel_context is not None + with otel.event_span(Event(name="child"), child_ctx, registered): + pass + _root, child = otel_exporter.get_finished_spans() + assert child.parent is not None + assert child.parent.span_id == root.get_span_context().span_id diff --git a/tests/units/reflex_otel/__init__.py b/tests/units/reflex_otel/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/tests/units/reflex_otel/test_init.py b/tests/units/reflex_otel/test_init.py new file mode 100644 index 00000000000..c1143f2404c --- /dev/null +++ b/tests/units/reflex_otel/test_init.py @@ -0,0 +1,35 @@ +"""Tests for the reflex_otel instrumentor.""" + +from collections.abc import Generator + +import pytest +from opentelemetry.sdk.trace import TracerProvider +from reflex_base import otel +from reflex_otel import ReflexInstrumentor + + +@pytest.fixture +def instrumentor() -> Generator[ReflexInstrumentor, None, None]: + inst = ReflexInstrumentor() + yield inst + if inst._is_instrumented_by_opentelemetry: + inst.uninstrument() + + +def test_instrument_toggles_trace_points(instrumentor: ReflexInstrumentor): + assert otel.enabled is False + instrumentor.instrument(tracer_provider=TracerProvider()) + assert otel.enabled is True + instrumentor.uninstrument() + assert otel.enabled is False + + +def test_instrument_is_idempotent(instrumentor: ReflexInstrumentor): + instrumentor.instrument() + instrumentor.instrument() + assert otel.enabled is True + + +def test_dependencies_target_reflex_base(instrumentor: ReflexInstrumentor): + (dep,) = instrumentor.instrumentation_dependencies() + assert dep.startswith("reflex-base") diff --git a/uv.lock b/uv.lock index 75a88ad310d..aebd8a13630 100644 --- a/uv.lock +++ b/uv.lock @@ -63,6 +63,7 @@ members = [ "reflex-docs-bundle", "reflex-hosting-cli", "reflex-integrations-docs", + "reflex-otel", "reflex-release", "reflex-site-shared", ] @@ -585,7 +586,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/66/54/eb9bfc647b19f2009dd5c7f5ec51c4e6ca831725f1aea7a993034f483147/contourpy-1.3.2.tar.gz", hash = "sha256:b6945942715a034c671b7fc54f9588126b0b8bf23db2696e3ca8328f3ff0ab54", size = 13466130, upload-time = "2025-04-15T17:47:53.79Z" } wheels = [ @@ -663,7 +664,7 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/58/01/1253e6698a07380cd31a736d248a3f2a50a7c88779a1813da27503cadc2a/contourpy-1.3.3.tar.gz", hash = "sha256:083e12155b210502d0bca491432bb04d56dc3432f95a979b429f2848c3dbe880", size = 13466174, upload-time = "2025-07-26T12:03:12.549Z" } @@ -1003,7 +1004,7 @@ name = "exceptiongroup" version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "typing-extensions" }, + { name = "typing-extensions", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" } wheels = [ @@ -1879,15 +1880,15 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "contourpy", version = "1.3.2", source = { registry = "https://pypi.org/simple" } }, - { name = "cycler" }, - { name = "fonttools" }, - { name = "kiwisolver" }, - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, - { name = "packaging" }, - { name = "pillow" }, - { name = "pyparsing" }, - { name = "python-dateutil" }, + { name = "contourpy", version = "1.3.2", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "cycler", marker = "python_full_version < '3.11'" }, + { name = "fonttools", marker = "python_full_version < '3.11'" }, + { name = "kiwisolver", marker = "python_full_version < '3.11'" }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "packaging", marker = "python_full_version < '3.11'" }, + { name = "pillow", marker = "python_full_version < '3.11'" }, + { name = "pyparsing", marker = "python_full_version < '3.11'" }, + { name = "python-dateutil", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/63/1b/4be5be87d43d327a0cf4de1a56e86f7f84c89312452406cf122efe2839e6/matplotlib-3.10.9.tar.gz", hash = "sha256:fd66508e8c6877d98e586654b608a0456db8d7e8a546eb1e2600efd957302358", size = 34811233, upload-time = "2026-04-24T00:14:13.539Z" } wheels = [ @@ -1963,16 +1964,16 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "contourpy", version = "1.3.3", source = { registry = "https://pypi.org/simple" } }, - { name = "cycler" }, - { name = "fonttools" }, - { name = "kiwisolver" }, - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "contourpy", version = "1.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "cycler", marker = "python_full_version >= '3.11'" }, + { name = "fonttools", marker = "python_full_version >= '3.11'" }, + { name = "kiwisolver", marker = "python_full_version >= '3.11'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, - { name = "packaging" }, - { name = "pillow" }, - { name = "pyparsing" }, - { name = "python-dateutil" }, + { name = "packaging", marker = "python_full_version >= '3.11'" }, + { name = "pillow", marker = "python_full_version >= '3.11'" }, + { name = "pyparsing", marker = "python_full_version >= '3.11'" }, + { name = "python-dateutil", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/49/64/f9a391af28f518b11ad45a8a712353c94a0aefce09d3703200e5c54b610a/matplotlib-3.11.1.tar.gz", hash = "sha256:69647db5746941c793d6e445a4cd349323ffb87d9cc958c2ad84a659b4832d30", size = 32612045, upload-time = "2026-07-18T03:39:46.63Z" } wheels = [ @@ -2424,6 +2425,60 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/78/0f/cc6afea3542a5142c5d8fc8211c5e059a8375105d004a41dfa2c7948dbb0/openai-2.53.0-py3-none-any.whl", hash = "sha256:c694ffc747a3c4d1663ef2b07b811315a476164ee5efa3a993967349ebca7618", size = 1659829, upload-time = "2026-08-03T21:41:59.581Z" }, ] +[[package]] +name = "opentelemetry-api" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/ee/8b/aa9e2d8b8dfa7c946f7dec5d1f8f6ba8eca062f43509a06bdb5ce93d26c0/opentelemetry_api-1.44.0.tar.gz", hash = "sha256:67647e5e9566edcf421166fdf022b3537f818635daa852b289e34604dc6fb33a", size = 72406, upload-time = "2026-07-16T15:25:32.678Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/ca/6f/a04e900f465ff3221ccc395522503e2d10e79fa21f2723c8e177aae1e0d1/opentelemetry_api-1.44.0-py3-none-any.whl", hash = "sha256:94b98c893a91b88657eaac1e3ba89618cdb85be6918196705354f34728b2cdef", size = 60018, upload-time = "2026-07-16T15:25:11.657Z" }, +] + +[[package]] +name = "opentelemetry-instrumentation" +version = "0.65b0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "opentelemetry-semantic-conventions" }, + { name = "packaging" }, + { name = "wrapt" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/13/91/3c58961cb0360cd60509064734f0be4275383c8681d73c580a40ca83ddce/opentelemetry_instrumentation-0.65b0.tar.gz", hash = "sha256:071d9d9eced9bd6460444ec3b0c77229870ed05a881c22c84fdede58e4eed09b", size = 42689, upload-time = "2026-07-16T15:25:50.275Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/40/7b/85eab1215f72adf0e68d3dc4a679b9bff993fa679ff34cd8dd378e2659fd/opentelemetry_instrumentation-0.65b0-py3-none-any.whl", hash = "sha256:ea967a72b9939b5fcfdad572753b4306c59dcb99e3f382d95dae04286805e137", size = 36717, upload-time = "2026-07-16T15:24:51.424Z" }, +] + +[[package]] +name = "opentelemetry-sdk" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "opentelemetry-semantic-conventions" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/5d/77/a6592cbc7c8d9bcc9d6757a9df45e04a7c585e3e6e7a13456da522b21109/opentelemetry_sdk-1.44.0.tar.gz", hash = "sha256:cebe7f65dc12f26ead75c6064de12fd2a9052e5060c0272d402cfa203aae123b", size = 208624, upload-time = "2026-07-16T15:25:46.078Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/e7/23/ff077e61886ee020a17ce9c8b6fa11c601c8d8345b09ea24f605445df62a/opentelemetry_sdk-1.44.0-py3-none-any.whl", hash = "sha256:df081c4c6bcfdb1211e3e86140376792643128a25f8d72d1d27675936e7e96ad", size = 137221, upload-time = "2026-07-16T15:25:29.534Z" }, +] + +[[package]] +name = "opentelemetry-semantic-conventions" +version = "0.65b0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/8f/73/0cbdebcb4cf545fdd328da14f5137e37d0770c3f26185e478b0d15d94f50/opentelemetry_semantic_conventions-0.65b0.tar.gz", hash = "sha256:f9b2b81e9d5b64f11bc952075e7e9c7fb0aab075c7fd1c46d597f1b919852d60", size = 148774, upload-time = "2026-07-16T15:25:46.902Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/a6/0e/49df70d9b81fb5cbae4bbf2a49d865b09bcbcbc4eb53f5851b1027738d78/opentelemetry_semantic_conventions-0.65b0-py3-none-any.whl", hash = "sha256:1cacde7b0ad306f84c5ef08c3dbe1bbaf20165bba6f8bff43b670e555a086bcb", size = 204645, upload-time = "2026-07-16T15:25:30.688Z" }, +] + [[package]] name = "orjson" version = "3.11.9" @@ -2534,10 +2589,10 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, - { name = "python-dateutil" }, - { name = "pytz" }, - { name = "tzdata" }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "python-dateutil", marker = "python_full_version < '3.11'" }, + { name = "pytz", marker = "python_full_version < '3.11'" }, + { name = "tzdata", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/33/01/d40b85317f86cf08d853a4f495195c73815fdf205eef3993821720274518/pandas-2.3.3.tar.gz", hash = "sha256:e05e1af93b977f7eafa636d043f9f94c7ee3ac81af99c13508215942e64c993b", size = 4495223, upload-time = "2025-09-29T23:34:51.853Z" } wheels = [ @@ -2606,10 +2661,10 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, - { name = "python-dateutil" }, - { name = "tzdata", marker = "sys_platform == 'emscripten' or sys_platform == 'win32'" }, + { name = "python-dateutil", marker = "python_full_version >= '3.11'" }, + { name = "tzdata", marker = "(python_full_version >= '3.11' and sys_platform == 'emscripten') or (python_full_version >= '3.11' and sys_platform == 'win32')" }, ] sdist = { url = "https://files.pythonhosted.org/packages/be/4f/5f3422a2afec5ffc46308b79e53291365a93748b498ac2e58bead0197916/pandas-3.0.5.tar.gz", hash = "sha256:dca3734d6ab7c906e6730f0788b0a1dbb9f2467731f9711f77995c8e9d62d712", size = 4658219, upload-time = "2026-07-22T22:19:28.819Z" } wheels = [ @@ -3690,6 +3745,7 @@ dev = [ { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, + { name = "opentelemetry-sdk" }, { name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, { name = "pandas", version = "3.0.5", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, { name = "pillow" }, @@ -3712,6 +3768,7 @@ dev = [ { name = "python-dotenv" }, { name = "pyyaml" }, { name = "reflex-docgen" }, + { name = "reflex-otel" }, { name = "reflex-release" }, { name = "reflex-site-shared" }, { name = "ruff" }, @@ -3772,6 +3829,7 @@ dev = [ { name = "hatchling" }, { name = "libsass" }, { name = "numpy" }, + { name = "opentelemetry-sdk" }, { name = "pandas" }, { name = "pillow" }, { name = "playwright" }, @@ -3793,6 +3851,7 @@ dev = [ { name = "python-dotenv" }, { name = "pyyaml" }, { name = "reflex-docgen", editable = "packages/reflex-docgen" }, + { name = "reflex-otel", editable = "packages/reflex-otel" }, { name = "reflex-release", editable = "packages/reflex-release" }, { name = "reflex-site-shared", editable = "packages/reflex-site-shared" }, { name = "ruff" }, @@ -3811,6 +3870,7 @@ dev = [ name = "reflex-base" source = { editable = "packages/reflex-base" } dependencies = [ + { name = "opentelemetry-api" }, { name = "packaging" }, { name = "platformdirs" }, { name = "rich" }, @@ -3824,6 +3884,7 @@ pydantic = [ [package.metadata] requires-dist = [ + { name = "opentelemetry-api", specifier = ">=1.24.0,<2.0" }, { name = "packaging", specifier = ">=24.2,<27" }, { name = "platformdirs", specifier = ">=4.3.7,<5.0" }, { name = "pydantic", marker = "extra == 'pydantic'", specifier = ">=2.12.0,<3.0" }, @@ -4149,6 +4210,29 @@ requires-dist = [ name = "reflex-integrations-docs" source = { editable = "packages/integrations-docs" } +[[package]] +name = "reflex-otel" +source = { editable = "packages/reflex-otel" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "opentelemetry-instrumentation" }, + { name = "reflex-base" }, +] + +[package.optional-dependencies] +instruments = [ + { name = "reflex-base" }, +] + +[package.metadata] +requires-dist = [ + { name = "opentelemetry-api", specifier = ">=1.24.0,<2.0" }, + { name = "opentelemetry-instrumentation", specifier = ">=0.45b0,<1.0" }, + { name = "reflex-base", editable = "packages/reflex-base" }, + { name = "reflex-base", marker = "extra == 'instruments'", editable = "packages/reflex-base" }, +] +provides-extras = ["instruments"] + [[package]] name = "reflex-pyplot" version = "0.2.1" @@ -4304,7 +4388,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/0f/37/6964b830433e654ec7485e45a00fc9a27cf868d622838f6b6d9c5ec0d532/scipy-1.15.3.tar.gz", hash = "sha256:eae3cf522bc7df64b42cad3925c876e1b0b6c35c1337c93e12c0f366f55b0eaf", size = 59419214, upload-time = "2025-05-08T16:13:05.955Z" } wheels = [ @@ -4365,7 +4449,7 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/7a/97/5a3609c4f8d58b039179648e62dd220f89864f56f7357f5d4f45c29eb2cc/scipy-1.17.1.tar.gz", hash = "sha256:95d8e012d8cb8816c226aef832200b1d45109ed4464303e997c5b13122b297c0", size = 30573822, upload-time = "2026-02-23T00:26:24.851Z" } wheels = [ @@ -4444,7 +4528,7 @@ resolution-markers = [ "python_full_version >= '3.12' and python_full_version < '3.14' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/a7/25/c2700dfaf6442b4effaa91af24ebce5dc9d31bb4a69706313aae70d72cd0/scipy-1.18.0.tar.gz", hash = "sha256:67b2ad2ad54c72ca6d04975a9b2df8c3638c34ddd5b28738e94fc2b57929d378", size = 30774447, upload-time = "2026-06-19T15:01:43.456Z" } wheels = [ From aa82464975196bf1e021cd4cbc29ba4cc7e5d61c Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 00:22:30 +0500 Subject: [PATCH 2/8] feat(otel): traceparent propagation, ASGI middleware and runtime metrics --- news/6227.feature.md | 1 + packages/reflex-base/news/6227.feature.md | 2 +- .../event/processor/base_state_processor.py | 5 + packages/reflex-base/src/reflex_base/otel.py | 214 +++++++++++++++--- packages/reflex-otel/README.md | 32 ++- packages/reflex-otel/news/6227.feature.md | 2 +- packages/reflex-otel/pyproject.toml | 4 +- .../reflex-otel/src/reflex_otel/__init__.py | 47 +++- reflex/app.py | 46 +++- tests/units/conftest.py | 38 ++++ tests/units/reflex_base/conftest.py | 26 --- .../processor/test_base_state_processor.py | 37 +++ tests/units/reflex_base/test_otel.py | 99 ++++++++ tests/units/reflex_otel/test_init.py | 14 ++ tests/units/test_app.py | 99 ++++++++ uv.lock | 43 +++- 16 files changed, 640 insertions(+), 69 deletions(-) create mode 100644 news/6227.feature.md delete mode 100644 tests/units/reflex_base/conftest.py diff --git a/news/6227.feature.md b/news/6227.feature.md new file mode 100644 index 00000000000..95d7d92747f --- /dev/null +++ b/news/6227.feature.md @@ -0,0 +1 @@ +Propagate a frontend `traceparent` into event spans, count websocket connections and message sizes, and wrap the ASGI app when the `reflex-otel` instrumentor is active. diff --git a/packages/reflex-base/news/6227.feature.md b/packages/reflex-base/news/6227.feature.md index 86c31cc08c5..0ed2e91a1e3 100644 --- a/packages/reflex-base/news/6227.feature.md +++ b/packages/reflex-base/news/6227.feature.md @@ -1 +1 @@ -Add inert OpenTelemetry trace points around event handler execution (`reflex_base.otel`); they cost one boolean check until the `reflex-otel` package enables them. +Add inert OpenTelemetry trace points and metrics around event handler execution, state acquisition and socket messages (`reflex_base.otel`); they cost one boolean check until the `reflex-otel` package enables them. diff --git a/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py index 9411d62927f..8f0449b42a4 100644 --- a/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py +++ b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py @@ -9,12 +9,14 @@ from collections.abc import Mapping, Sequence from enum import Enum from importlib.util import find_spec +from time import perf_counter from typing import TYPE_CHECKING, Any from reflex.istate.data import RouterData from reflex.istate.manager.token import BaseStateToken from reflex.istate.proxy import StateProxy from reflex.utils import console, types +from reflex_base import otel from reflex_base.event.context import EventContext from reflex_base.event.processor.event_processor import EventProcessor, EventQueueEntry from reflex_base.registry import RegisteredEventHandler @@ -340,6 +342,7 @@ async def _execute_event( ctx = entry.ctx event = entry.event router_data = event.router_data or {} + acquire_start = perf_counter() if otel.enabled else 0.0 # Get the state for the session exclusively. async with ctx.state_manager.modify_state_with_links( BaseStateToken( @@ -348,6 +351,8 @@ async def _execute_event( ), event=entry.event, ) as state: + if acquire_start: + otel.record_state_acquired(acquire_start, event) # Compatibility hack rehydrate the state before processing this event. needs_to_rehydrate = bool( not state.router_data and event.name != _hydrate_event_name() diff --git a/packages/reflex-base/src/reflex_base/otel.py b/packages/reflex-base/src/reflex_base/otel.py index 074a4df52da..e971ae9a70a 100644 --- a/packages/reflex-base/src/reflex_base/otel.py +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -1,9 +1,10 @@ -"""OpenTelemetry trace points for the Reflex runtime. +"""OpenTelemetry trace points and metrics for the Reflex runtime. The framework calls into this module at a small number of fixed points (event -dispatch, event context forks). Every entry point checks the module-level -``enabled`` flag first, so with no instrumentation installed the cost is one -attribute read and no ``opentelemetry`` object is ever created. +dispatch, event context forks, state acquisition, socket messages). Every entry +point checks the module-level ``enabled`` flag first, so with no +instrumentation installed the cost is one attribute read and no +``opentelemetry`` object is ever created. The ``reflex-otel`` package flips the flag via :func:`enable` once a tracer provider is available. @@ -11,18 +12,20 @@ from __future__ import annotations -from collections.abc import Iterator -from contextlib import contextmanager -from typing import TYPE_CHECKING +from collections.abc import Callable, Iterator, Mapping +from contextlib import contextmanager, nullcontext +from time import perf_counter +from typing import TYPE_CHECKING, Any from opentelemetry import context as otel_context -from opentelemetry import trace +from opentelemetry import metrics, propagate, trace +from opentelemetry.context import Context from opentelemetry.trace import SpanKind from reflex_base.constants.base import Reflex if TYPE_CHECKING: - from opentelemetry.context import Context + from contextlib import AbstractContextManager from reflex_base.event import Event from reflex_base.event.context import EventContext @@ -30,39 +33,104 @@ INSTRUMENTATION_NAME = "reflex" -# Attribute keys emitted on Reflex spans. +# Attribute keys emitted on Reflex spans and metrics. ATTR_EVENT_NAME = "reflex.event.name" ATTR_EVENT_TXID = "reflex.event.txid" ATTR_EVENT_PARENT_TXID = "reflex.event.parent_txid" ATTR_EVENT_BACKGROUND = "reflex.event.background" ATTR_SESSION_ID = "session.id" ATTR_CODE_FUNCTION_NAME = "code.function.name" +ATTR_ERROR_TYPE = "error.type" +ATTR_NETWORK_IO_DIRECTION = "network.io.direction" + +# Metric instrument names. +METRIC_EVENT_DURATION = "reflex.event.duration" +METRIC_STATE_ACQUIRE_DURATION = "reflex.state.acquire.duration" +METRIC_WEBSOCKET_MESSAGE_SIZE = "reflex.websocket.message.size" +METRIC_WEBSOCKET_CONNECTIONS = "reflex.websocket.connections" + +# Key of the W3C trace context carried in an event payload sent by the frontend. +TRACEPARENT_FIELD = "traceparent" # Read at every trace point; True only after enable() ran. enabled: bool = False +# Wraps the app's ASGI callable when set (installed by enable()). +asgi_middleware: Callable[[Any], Any] | None = None + _tracer: trace.Tracer = trace.NoOpTracer() +_noop_meter = metrics.NoOpMeter(INSTRUMENTATION_NAME) +_event_duration = _noop_meter.create_histogram(METRIC_EVENT_DURATION) +_state_acquire_duration = _noop_meter.create_histogram(METRIC_STATE_ACQUIRE_DURATION) +_message_size = _noop_meter.create_histogram(METRIC_WEBSOCKET_MESSAGE_SIZE) +_ws_connections = _noop_meter.create_up_down_counter(METRIC_WEBSOCKET_CONNECTIONS) + + +def _create_instruments(meter: metrics.Meter) -> None: + """Create the metric instruments on the given meter. + + Args: + meter: The meter to create the instruments on. + """ + global _event_duration, _state_acquire_duration, _message_size, _ws_connections + _event_duration = meter.create_histogram( + METRIC_EVENT_DURATION, + unit="s", + description="Duration of event handler executions.", + ) + _state_acquire_duration = meter.create_histogram( + METRIC_STATE_ACQUIRE_DURATION, + unit="s", + description="Time an event waited to acquire and load its session state.", + ) + _message_size = meter.create_histogram( + METRIC_WEBSOCKET_MESSAGE_SIZE, + unit="By", + description="Serialized size of socket messages exchanged with the client.", + ) + _ws_connections = meter.create_up_down_counter( + METRIC_WEBSOCKET_CONNECTIONS, + unit="{connection}", + description="Number of open client socket connections.", + ) -def enable(tracer_provider: trace.TracerProvider | None = None) -> None: - """Turn the trace points on. +def enable( + tracer_provider: trace.TracerProvider | None = None, + meter_provider: metrics.MeterProvider | None = None, + asgi_middleware_factory: Callable[[Any], Any] | None = None, +) -> None: + """Turn the trace points and metrics on. Args: tracer_provider: The provider to obtain the tracer from. Defaults to the global provider, which may be configured later; the returned proxy tracer picks it up automatically. + meter_provider: The provider to obtain the meter from. Defaults to the + global provider. + asgi_middleware_factory: Callable that wraps the app's ASGI callable, + e.g. the OpenTelemetry ASGI middleware. Applied by the app when it + builds its ASGI app. """ - global _tracer, enabled + global _tracer, enabled, asgi_middleware _tracer = trace.get_tracer( INSTRUMENTATION_NAME, Reflex.VERSION, tracer_provider=tracer_provider ) + _create_instruments( + metrics.get_meter( + INSTRUMENTATION_NAME, Reflex.VERSION, meter_provider=meter_provider + ) + ) + asgi_middleware = asgi_middleware_factory enabled = True def disable() -> None: - """Turn the trace points off and drop the tracer.""" - global _tracer, enabled + """Turn the trace points off and drop the tracer and instruments.""" + global _tracer, enabled, asgi_middleware enabled = False + asgi_middleware = None _tracer = trace.NoOpTracer() + _create_instruments(_noop_meter) def capture_context() -> Context | None: @@ -74,11 +142,54 @@ def capture_context() -> Context | None: return otel_context.get_current() if enabled else None +class _AttachedContext: + """Attach a context on enter and detach it on exit.""" + + __slots__ = ("_context", "_token") + + def __init__(self, context: Context): + """Store the context to attach. + + Args: + context: The context to make current. + """ + self._context = context + + def __enter__(self) -> None: + """Attach the context.""" + self._token = otel_context.attach(self._context) + + def __exit__(self, *exc_info: object) -> None: + """Detach the context. + + Args: + *exc_info: Ignored exception details. + """ + otel_context.detach(self._token) + + +def remote_context(carrier: Mapping[str, Any]) -> AbstractContextManager[None]: + """Make the trace context carried by a frontend event the current context. + + Events that carry no ``traceparent`` start a new trace: the websocket + connection span (if any) is deliberately not used as their parent. + + Args: + carrier: The raw event fields received from the frontend. + + Returns: + A context manager to run the enqueue under; a no-op when tracing is off. + """ + if not enabled: + return nullcontext() + return _AttachedContext(propagate.extract(carrier, context=Context())) + + @contextmanager def event_span( event: Event, ctx: EventContext, registered_handler: RegisteredEventHandler ) -> Iterator[trace.Span]: - """Open the span for one event handler execution. + """Open the span for one event handler execution and record its duration. Chained events are parented under the span that enqueued them; events that arrive from the frontend start a new trace. @@ -92,17 +203,60 @@ def event_span( The active span. """ handler = registered_handler.handler - with _tracer.start_as_current_span( - event.name, - context=ctx.otel_context, - kind=SpanKind.SERVER, - attributes={ - ATTR_EVENT_NAME: event.name, - ATTR_EVENT_TXID: ctx.txid, - ATTR_EVENT_BACKGROUND: handler.is_background, - ATTR_SESSION_ID: ctx.token, - ATTR_CODE_FUNCTION_NAME: getattr(handler.fn, "__qualname__", event.name), - } - | ({ATTR_EVENT_PARENT_TXID: ctx.parent_txid} if ctx.parent_txid else {}), - ) as span: - yield span + metric_attributes = { + ATTR_EVENT_NAME: event.name, + ATTR_EVENT_BACKGROUND: handler.is_background, + } + start = perf_counter() + try: + with _tracer.start_as_current_span( + event.name, + context=ctx.otel_context, + kind=SpanKind.SERVER, + attributes=metric_attributes + | { + ATTR_EVENT_TXID: ctx.txid, + ATTR_SESSION_ID: ctx.token, + ATTR_CODE_FUNCTION_NAME: getattr( + handler.fn, "__qualname__", event.name + ), + } + | ({ATTR_EVENT_PARENT_TXID: ctx.parent_txid} if ctx.parent_txid else {}), + ) as span: + yield span + except BaseException as ex: + metric_attributes[ATTR_ERROR_TYPE] = type(ex).__qualname__ + raise + finally: + _event_duration.record(perf_counter() - start, metric_attributes) + + +def record_state_acquired(start: float, event: Event) -> None: + """Record how long an event waited for its session state. + + Args: + start: ``perf_counter()`` value taken before the state was requested. + event: The event that acquired the state. + """ + _state_acquire_duration.record( + perf_counter() - start, {ATTR_EVENT_NAME: event.name} + ) + + +def record_message_size(size: int, direction: str) -> None: + """Record the serialized size of one socket message. + + Args: + size: The message size in bytes. + direction: ``"transmit"`` for server-to-client, ``"receive"`` otherwise. + """ + _message_size.record(size, {ATTR_NETWORK_IO_DIRECTION: direction}) + + +def record_connection(delta: int) -> None: + """Adjust the open socket connection count. + + Args: + delta: ``1`` on connect, ``-1`` on disconnect. + """ + _ws_connections.add(delta) diff --git a/packages/reflex-otel/README.md b/packages/reflex-otel/README.md index 298f1b33a71..874f8a8f506 100644 --- a/packages/reflex-otel/README.md +++ b/packages/reflex-otel/README.md @@ -10,4 +10,34 @@ ReflexInstrumentor().instrument() The package registers an `opentelemetry_instrumentor` entry point, so `opentelemetry-instrument reflex run` enables it automatically. Configure a -tracer provider (for example with `opentelemetry-sdk`) to export the spans. +tracer provider and a meter provider (for example with `opentelemetry-sdk`) +to export the data. + +## What you get + +Traces: + +- One `SERVER` span per event handler run, named after the event. Chained + events are children of the span that enqueued them. An event sent by the + frontend with a `traceparent` field continues that trace; otherwise it + starts a new one. +- HTTP requests and the websocket connection are wrapped in the standard + OpenTelemetry ASGI middleware (per-message websocket spans are off). + +Metrics: + +| Instrument | Type | Unit | Attributes | +| --- | --- | --- | --- | +| `reflex.event.duration` | histogram | s | `reflex.event.name`, `reflex.event.background`, `error.type` | +| `reflex.state.acquire.duration` | histogram | s | `reflex.event.name` | +| `reflex.websocket.message.size` | histogram | By | `network.io.direction` (`transmit`/`receive`) | +| `reflex.websocket.connections` | up-down counter | `{connection}` | | + +Plus the ASGI middleware's `http.server.*` metrics. + +## Options + +`instrument()` accepts `tracer_provider`, `meter_provider`, `excluded_urls` +(comma-separated URL patterns skipped by the ASGI middleware; defaults to the +`OTEL_PYTHON_REFLEX_EXCLUDED_URLS` environment variable) and the ASGI hooks +`server_request_hook`, `client_request_hook`, `client_response_hook`. diff --git a/packages/reflex-otel/news/6227.feature.md b/packages/reflex-otel/news/6227.feature.md index 19fa9421ea3..fdc6ea8d17b 100644 --- a/packages/reflex-otel/news/6227.feature.md +++ b/packages/reflex-otel/news/6227.feature.md @@ -1 +1 @@ -Add the `reflex-otel` package: an OpenTelemetry instrumentor that turns on the framework's built-in trace points (one span per event handler run, chained events parented under the enqueuing span). +Add the `reflex-otel` package: an OpenTelemetry instrumentor that turns on the framework's built-in trace points and metrics (one span per event handler run, chained events parented under the enqueuing span, frontend `traceparent` propagation, event/state/websocket metrics) and wraps the ASGI app in the OpenTelemetry ASGI middleware. diff --git a/packages/reflex-otel/pyproject.toml b/packages/reflex-otel/pyproject.toml index 1eb2a10c36a..4c693ffa3d8 100644 --- a/packages/reflex-otel/pyproject.toml +++ b/packages/reflex-otel/pyproject.toml @@ -9,7 +9,9 @@ maintainers = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] requires-python = ">=3.10" dependencies = [ "opentelemetry-api >=1.24.0,<2.0", - "opentelemetry-instrumentation >=0.45b0,<1.0", + "opentelemetry-instrumentation >=0.49b0,<1.0", + "opentelemetry-instrumentation-asgi >=0.49b0,<1.0", + "opentelemetry-util-http >=0.49b0,<1.0", "reflex-base >= 0.9.7.post45.dev0", ] diff --git a/packages/reflex-otel/src/reflex_otel/__init__.py b/packages/reflex-otel/src/reflex_otel/__init__.py index bbc530ad431..95aef1e90e9 100644 --- a/packages/reflex-otel/src/reflex_otel/__init__.py +++ b/packages/reflex-otel/src/reflex_otel/__init__.py @@ -3,16 +3,21 @@ from __future__ import annotations from collections.abc import Collection -from typing import Any +from typing import Any, Literal +from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware from opentelemetry.instrumentation.instrumentor import BaseInstrumentor +from opentelemetry.util.http import get_excluded_urls from reflex_base import otel _instruments = ("reflex-base >= 0.9.7.post45.dev0",) +# Per-message websocket spans are noise; Reflex emits one span per event instead. +_ASGI_EXCLUDED_SPANS: list[Literal["receive", "send"]] = ["receive", "send"] + class ReflexInstrumentor(BaseInstrumentor): - """Enable the trace points built into the Reflex runtime. + """Enable the trace points and metrics built into the Reflex runtime. Usage:: @@ -20,6 +25,11 @@ class ReflexInstrumentor(BaseInstrumentor): or let ``opentelemetry-instrument`` load it through the ``opentelemetry_instrumentor`` entry point. + + Besides the per-event spans and metrics, the app's ASGI callable is + wrapped in the OpenTelemetry ASGI middleware, so HTTP requests (uploads, + custom API routes) and the websocket connection get server spans and + HTTP metrics as well. """ def instrumentation_dependencies(self) -> Collection[str]: @@ -34,9 +44,38 @@ def _instrument(self, **kwargs: Any) -> None: """Turn on the Reflex trace points. Args: - **kwargs: ``tracer_provider`` selects the provider; defaults to the global one. + **kwargs: ``tracer_provider`` and ``meter_provider`` select the + providers (default: the global ones). ``excluded_urls`` is a + comma-separated list of URL patterns the ASGI middleware skips + (default: ``OTEL_PYTHON_REFLEX_EXCLUDED_URLS``). + ``server_request_hook``, ``client_request_hook`` and + ``client_response_hook`` are forwarded to the ASGI middleware. """ - otel.enable(tracer_provider=kwargs.get("tracer_provider")) + tracer_provider = kwargs.get("tracer_provider") + meter_provider = kwargs.get("meter_provider") + excluded_urls = kwargs.get("excluded_urls") + + def asgi_middleware(app: Any) -> Any: + return OpenTelemetryMiddleware( + app, + excluded_urls=( + get_excluded_urls("REFLEX") + if excluded_urls is None + else excluded_urls + ), + server_request_hook=kwargs.get("server_request_hook"), + client_request_hook=kwargs.get("client_request_hook"), + client_response_hook=kwargs.get("client_response_hook"), + tracer_provider=tracer_provider, + meter_provider=meter_provider, + exclude_spans=_ASGI_EXCLUDED_SPANS, + ) + + otel.enable( + tracer_provider=tracer_provider, + meter_provider=meter_provider, + asgi_middleware_factory=asgi_middleware, + ) def _uninstrument(self, **kwargs: Any) -> None: """Turn off the Reflex trace points. diff --git a/reflex/app.py b/reflex/app.py index 258410fe9f7..fc82ca50a27 100644 --- a/reflex/app.py +++ b/reflex/app.py @@ -27,7 +27,7 @@ from types import SimpleNamespace from typing import TYPE_CHECKING, Any, overload -from reflex_base import constants +from reflex_base import constants, otel from reflex_base.components.component import Component, ComponentStyle from reflex_base.config import get_config, reload_config from reflex_base.context.base import BaseContext @@ -577,8 +577,8 @@ def _setup_state(self) -> None: ping_interval=environment.REFLEX_SOCKET_INTERVAL.get(), ping_timeout=environment.REFLEX_SOCKET_TIMEOUT.get(), json=SimpleNamespace( - dumps=staticmethod(format.json_dumps), - loads=staticmethod(json.loads), + dumps=staticmethod(_sio_dumps), + loads=staticmethod(_sio_loads), ), allow_upgrades=False, transports=[config.transport], @@ -809,6 +809,8 @@ def __call__(self) -> ASGIApp: self._context_middleware(asgi_app), ) App._add_cors(top_asgi_app) + if otel.asgi_middleware is not None: + return otel.asgi_middleware(top_asgi_app) return top_asgi_app def _add_default_endpoints(self): @@ -1917,6 +1919,37 @@ async def health(_request: Request) -> JSONResponse: return JSONResponse(content=health_status, status_code=status_code) +def _sio_dumps(obj: Any, **kwargs: Any) -> str: + """Serialize an outgoing Socket.IO packet, recording its size when telemetry is on. + + Args: + obj: The packet payload. + **kwargs: Options forwarded to the JSON encoder. + + Returns: + The JSON string. + """ + data = format.json_dumps(obj, **kwargs) + if otel.enabled: + otel.record_message_size(len(data), "transmit") + return data + + +def _sio_loads(data: str | bytes, **kwargs: Any) -> Any: + """Deserialize an incoming Socket.IO packet, recording its size when telemetry is on. + + Args: + data: The JSON string. + **kwargs: Options forwarded to the JSON decoder. + + Returns: + The decoded payload. + """ + if otel.enabled: + otel.record_message_size(len(data), "receive") + return json.loads(data, **kwargs) + + class EventNamespace(AsyncNamespace): """The event namespace.""" @@ -1997,6 +2030,8 @@ async def on_connect(self, sid: str, environ: dict): console.warn( f"Frontend version {subprotocol} for session {sid} does not match the backend version {constants.Reflex.VERSION}." ) + if otel.enabled: + otel.record_connection(1) def on_disconnect(self, sid: str) -> asyncio.Task | None: """Event for when the websocket disconnects. @@ -2007,6 +2042,8 @@ def on_disconnect(self, sid: str) -> asyncio.Task | None: Returns: An asyncio Task for cleaning up the token, or None. """ + if otel.enabled: + otel.record_connection(-1) self._client_error_counts.pop(sid, None) # Get token before cleaning up disconnect_token = self.sid_to_token.get(sid) @@ -2141,7 +2178,8 @@ async def on_event(self, sid: str, data: Any): if (path := router_data.get(constants.RouteVar.PATH)) else "404" ).removeprefix("/") - await self.app.event_processor.enqueue(token, event) + with otel.remote_context(fields): + await self.app.event_processor.enqueue(token, event) async def on_ping(self, sid: str): """Event for testing the API endpoint. diff --git a/tests/units/conftest.py b/tests/units/conftest.py index 8ef4232fbc7..c2d215cc358 100644 --- a/tests/units/conftest.py +++ b/tests/units/conftest.py @@ -9,6 +9,12 @@ import pytest import pytest_asyncio +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import InMemoryMetricReader +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from reflex_base import otel from reflex_base.components.memo import MEMOS from reflex_base.event import Event, EventSpec from reflex_base.event.context import EventContext @@ -517,3 +523,35 @@ def preserve_memo_registries(): finally: MEMOS.clear() MEMOS.update(memos) + + +@pytest.fixture +def otel_exporter() -> Generator[InMemorySpanExporter, None, None]: + """Enable the reflex_base.otel trace points against an in-memory exporter. + + Yields: + The exporter collecting finished spans. + """ + exporter = InMemorySpanExporter() + provider = TracerProvider() + provider.add_span_processor(SimpleSpanProcessor(exporter)) + otel.enable(tracer_provider=provider) + try: + yield exporter + finally: + otel.disable() + + +@pytest.fixture +def otel_metrics() -> Generator[InMemoryMetricReader, None, None]: + """Enable the reflex_base.otel metrics against an in-memory reader. + + Yields: + The reader collecting recorded metrics. + """ + reader = InMemoryMetricReader() + otel.enable(meter_provider=MeterProvider(metric_readers=[reader])) + try: + yield reader + finally: + otel.disable() diff --git a/tests/units/reflex_base/conftest.py b/tests/units/reflex_base/conftest.py deleted file mode 100644 index 5ae9153f10e..00000000000 --- a/tests/units/reflex_base/conftest.py +++ /dev/null @@ -1,26 +0,0 @@ -"""Shared fixtures for reflex_base unit tests.""" - -from collections.abc import Generator - -import pytest -from opentelemetry.sdk.trace import TracerProvider -from opentelemetry.sdk.trace.export import SimpleSpanProcessor -from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter -from reflex_base import otel - - -@pytest.fixture -def otel_exporter() -> Generator[InMemorySpanExporter, None, None]: - """Enable the reflex_base.otel trace points against an in-memory exporter. - - Yields: - The exporter collecting finished spans. - """ - exporter = InMemorySpanExporter() - provider = TracerProvider() - provider.add_span_processor(SimpleSpanProcessor(exporter)) - otel.enable(tracer_provider=provider) - try: - yield exporter - finally: - otel.disable() diff --git a/tests/units/reflex_base/event/processor/test_base_state_processor.py b/tests/units/reflex_base/event/processor/test_base_state_processor.py index 6f89d9adf49..1243f6df6dd 100644 --- a/tests/units/reflex_base/event/processor/test_base_state_processor.py +++ b/tests/units/reflex_base/event/processor/test_base_state_processor.py @@ -214,3 +214,40 @@ async def preprocess(self, app, state, event) -> StateUpdate: client_event_names = {e.name for _, events in emitted_events for e in events} assert "_call_function" in client_event_names assert "_redirect" in client_event_names + + +async def test_execute_event_records_state_acquire_duration( + wired_app: App, + real_base_state_processor: BaseStateEventProcessor, + token: str, + otel_metrics, +): + """Executing an event records how long acquiring the session state took. + + Args: + wired_app: The App wired to the processor's state manager. + real_base_state_processor: The unmocked BaseStateEventProcessor. + token: The client token. + otel_metrics: In-memory metric reader with metrics enabled. + """ + from reflex_base import otel + + class AcquireState(State): + @event + def noop(self): + pass + + async with real_base_state_processor as processor: + await processor.enqueue(token, Event.from_event_type(AcquireState.noop())[0]) + await processor.join(1) + + metrics = otel_metrics.get_metrics_data() + (metric,) = [ + m + for rm in metrics.resource_metrics + for sm in rm.scope_metrics + for m in sm.metrics + if m.name == otel.METRIC_STATE_ACQUIRE_DURATION + ] + names = {p.attributes[otel.ATTR_EVENT_NAME] for p in metric.data.data_points} + assert Event.from_event_type(AcquireState.noop())[0].name in names diff --git a/tests/units/reflex_base/test_otel.py b/tests/units/reflex_base/test_otel.py index 93674e1269d..6d95cda6f59 100644 --- a/tests/units/reflex_base/test_otel.py +++ b/tests/units/reflex_base/test_otel.py @@ -1,8 +1,12 @@ """Tests for the reflex_base.otel trace points.""" +from contextlib import nullcontext +from time import perf_counter + import pytest from opentelemetry import context as otel_context from opentelemetry import trace +from opentelemetry.sdk.metrics.export import InMemoryMetricReader from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter from opentelemetry.trace import SpanKind, StatusCode from reflex_base import otel @@ -93,3 +97,98 @@ def test_event_span_parents_under_captured_context( _root, child = otel_exporter.get_finished_spans() assert child.parent is not None assert child.parent.span_id == root.get_span_context().span_id + + +def _metric_points(reader: InMemoryMetricReader, name: str) -> list: + data = reader.get_metrics_data() + assert data is not None + return [ + point + for rm in data.resource_metrics + for sm in rm.scope_metrics + for metric in sm.metrics + if metric.name == name + for point in metric.data.data_points + ] + + +def test_event_span_records_duration_metric(otel_metrics: InMemoryMetricReader): + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + with otel.event_span(Event(name="ok"), _ctx(), registered): + pass + with ( + pytest.raises(RuntimeError), + otel.event_span(Event(name="bad"), _ctx(), registered), + ): + raise RuntimeError + bad, ok = sorted( + _metric_points(otel_metrics, otel.METRIC_EVENT_DURATION), + key=lambda p: p.attributes[otel.ATTR_EVENT_NAME], + ) + assert bad.attributes == { + otel.ATTR_EVENT_NAME: "bad", + otel.ATTR_EVENT_BACKGROUND: False, + otel.ATTR_ERROR_TYPE: "RuntimeError", + } + assert ok.attributes == { + otel.ATTR_EVENT_NAME: "ok", + otel.ATTR_EVENT_BACKGROUND: False, + } + assert ok.count == bad.count == 1 + assert ok.sum >= 0 + + +def test_metric_helpers_record(otel_metrics: InMemoryMetricReader): + otel.record_state_acquired(perf_counter(), Event(name="e")) + otel.record_message_size(42, "transmit") + otel.record_connection(1) + otel.record_connection(1) + otel.record_connection(-1) + (acquire,) = _metric_points(otel_metrics, otel.METRIC_STATE_ACQUIRE_DURATION) + assert acquire.attributes == {otel.ATTR_EVENT_NAME: "e"} + (size,) = _metric_points(otel_metrics, otel.METRIC_WEBSOCKET_MESSAGE_SIZE) + assert size.sum == 42 + assert size.attributes == {otel.ATTR_NETWORK_IO_DIRECTION: "transmit"} + (conns,) = _metric_points(otel_metrics, otel.METRIC_WEBSOCKET_CONNECTIONS) + assert conns.value == 1 + + +def test_remote_context_disabled_is_noop(): + with otel._tracer.start_as_current_span("outer"), otel.remote_context({}): + pass + assert isinstance(otel.remote_context({"traceparent": "x"}), nullcontext) + + +def test_remote_context_uses_traceparent(otel_exporter: InMemorySpanExporter): + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + traceparent = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01" + with otel.remote_context({"traceparent": traceparent}): + ctx = _ctx().fork() + with otel.event_span(Event(name="e"), ctx, registered): + pass + (span,) = otel_exporter.get_finished_spans() + assert span.parent is not None + assert span.parent.is_remote + assert format(span.parent.trace_id, "032x") == "0af7651916cd43dd8448eb211c80319c" + assert format(span.parent.span_id, "016x") == "b7ad6b7169203331" + + +def test_remote_context_without_traceparent_starts_new_trace( + otel_exporter: InMemorySpanExporter, +): + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + with otel._tracer.start_as_current_span("websocket"), otel.remote_context({}): + ctx = _ctx().fork() + with otel.event_span(Event(name="e"), ctx, registered): + pass + _ws, span = otel_exporter.get_finished_spans() + assert span.parent is None + + +def test_asgi_middleware_hook_toggles(): + assert otel.asgi_middleware is None + factory = lambda app: app # noqa: E731 + otel.enable(asgi_middleware_factory=factory) + assert otel.asgi_middleware is factory + otel.disable() + assert otel.asgi_middleware is None diff --git a/tests/units/reflex_otel/test_init.py b/tests/units/reflex_otel/test_init.py index c1143f2404c..dfba88b0ac4 100644 --- a/tests/units/reflex_otel/test_init.py +++ b/tests/units/reflex_otel/test_init.py @@ -3,6 +3,7 @@ from collections.abc import Generator import pytest +from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware from opentelemetry.sdk.trace import TracerProvider from reflex_base import otel from reflex_otel import ReflexInstrumentor @@ -33,3 +34,16 @@ def test_instrument_is_idempotent(instrumentor: ReflexInstrumentor): def test_dependencies_target_reflex_base(instrumentor: ReflexInstrumentor): (dep,) = instrumentor.instrumentation_dependencies() assert dep.startswith("reflex-base") + + +def test_instrument_installs_asgi_middleware(instrumentor: ReflexInstrumentor): + instrumentor.instrument(tracer_provider=TracerProvider()) + assert otel.asgi_middleware is not None + + async def app(scope, receive, send): ... + + wrapped = otel.asgi_middleware(app) + assert isinstance(wrapped, OpenTelemetryMiddleware) + assert wrapped.app is app + instrumentor.uninstrument() + assert otel.asgi_middleware is None diff --git a/tests/units/test_app.py b/tests/units/test_app.py index 60e3122d70a..41ee6a5871f 100644 --- a/tests/units/test_app.py +++ b/tests/units/test_app.py @@ -4257,3 +4257,102 @@ def test_client_error_constants_match_frontend(): f'const ERROR_TYPE_STATE_UPDATE = "{constants.ClientErrorType.STATE_UPDATE}"' in state_js ) + + +def test_call_app_wraps_with_otel_asgi_middleware(): + """The app's ASGI callable is wrapped when instrumentation installs a middleware.""" + from reflex_base import otel + + app = App() + app._compile = unittest.mock.Mock() + wrapped = [] + otel.enable(asgi_middleware_factory=lambda asgi: wrapped.append(asgi) or asgi) + try: + api = app() + finally: + otel.disable() + assert wrapped == [api] + + +def test_sio_json_records_message_sizes(otel_metrics): + """Socket.IO packet serialization records sizes in both directions.""" + from reflex_base import otel + + from reflex.app import _sio_dumps, _sio_loads + + data = _sio_dumps({"a": 1}, separators=(",", ":")) + assert data == '{"a":1}' + assert _sio_loads(data) == {"a": 1} + metrics = otel_metrics.get_metrics_data() + (metric,) = [ + m + for rm in metrics.resource_metrics + for sm in rm.scope_metrics + for m in sm.metrics + if m.name == otel.METRIC_WEBSOCKET_MESSAGE_SIZE + ] + points = { + p.attributes[otel.ATTR_NETWORK_IO_DIRECTION]: p.sum + for p in metric.data.data_points + } + assert points == {"transmit": len(data), "receive": len(data)} + + +@pytest.mark.asyncio +async def test_on_event_uses_frontend_traceparent(otel_exporter): + """A traceparent in the event payload becomes the parent of the event span.""" + from opentelemetry import trace + from reflex_base import otel + + from reflex.app import EventNamespace + + mock_app = unittest.mock.Mock() + mock_app.router.return_value = "/" + mock_app.sio.get_environ.return_value = { + "asgi.scope": {"headers": [], "client": ("127.0.0.1", 1)} + } + seen: list = [] + + async def enqueue(token, event): + await asyncio.sleep(0) + seen.append(trace.get_current_span().get_span_context()) + + mock_app.event_processor.enqueue = enqueue + ns = EventNamespace(namespace="/", app=mock_app) + ns._token_manager.sid_to_token["sid"] = "tok" + + traceparent = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01" + with otel._tracer.start_as_current_span("websocket"): + await ns.on_event("sid", {"name": "state.h", "traceparent": traceparent}) + await ns.on_event("sid", {"name": "state.h"}) + remote, fresh = seen + assert f"{remote.trace_id:032x}" == "0af7651916cd43dd8448eb211c80319c" + assert not fresh.is_valid + + +@pytest.mark.asyncio +async def test_connect_disconnect_counts_connections(otel_metrics): + """Connect and disconnect adjust the open connection gauge.""" + from reflex_base import otel + + from reflex.app import EventNamespace + + mock_app = unittest.mock.Mock() + mock_app._state = None + ns = EventNamespace(namespace="/", app=mock_app) + ns.emit = unittest.mock.AsyncMock() + await ns.on_connect("sid1", {"QUERY_STRING": "token=t1"}) + await ns.on_connect("sid2", {"QUERY_STRING": "token=t2"}) + task = ns.on_disconnect("sid1") + if task is not None: + await task + metrics = otel_metrics.get_metrics_data() + (metric,) = [ + m + for rm in metrics.resource_metrics + for sm in rm.scope_metrics + for m in sm.metrics + if m.name == otel.METRIC_WEBSOCKET_CONNECTIONS + ] + (point,) = metric.data.data_points + assert point.value == 1 diff --git a/uv.lock b/uv.lock index aebd8a13630..b1a0cbd61bc 100644 --- a/uv.lock +++ b/uv.lock @@ -288,6 +288,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/f1/a9/ceb963a6f5be13024f850816db1ab9b53efbcf4be688c64a4c1ecccda810/asgiproxy-0.2.0-py3-none-any.whl", hash = "sha256:b952917b2a3318c558f2cba6a91329b060af5621e5f4e301948d922fe2fdcc5b", size = 8839, upload-time = "2025-02-07T10:18:49.692Z" }, ] +[[package]] +name = "asgiref" +version = "3.12.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "typing-extensions", marker = "python_full_version < '3.11'" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/e6/26/3b59f2bdae5f640389becb1f673cded775287f5fc4f816309d9ca9a3f93d/asgiref-3.12.1.tar.gz", hash = "sha256:59dcb51c272ad209d59bed5708a64a333083e86017d7fcdd67498eeab7784340", size = 42378, upload-time = "2026-07-14T09:56:18.087Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/c0/1b/54f4ad77cd8a584fa70746c47df988e002cf1ee1eba43364d46f87803647/asgiref-3.12.1-py3-none-any.whl", hash = "sha256:fe386d1c2bff7259ea95929266d12a8cf9a8b5a1c2598402967d8792e7a7c094", size = 25478, upload-time = "2026-07-14T09:56:16.926Z" }, +] + [[package]] name = "async-timeout" version = "5.0.1" @@ -2452,6 +2464,22 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/40/7b/85eab1215f72adf0e68d3dc4a679b9bff993fa679ff34cd8dd378e2659fd/opentelemetry_instrumentation-0.65b0-py3-none-any.whl", hash = "sha256:ea967a72b9939b5fcfdad572753b4306c59dcb99e3f382d95dae04286805e137", size = 36717, upload-time = "2026-07-16T15:24:51.424Z" }, ] +[[package]] +name = "opentelemetry-instrumentation-asgi" +version = "0.65b0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "asgiref" }, + { name = "opentelemetry-api" }, + { name = "opentelemetry-instrumentation" }, + { name = "opentelemetry-semantic-conventions" }, + { name = "opentelemetry-util-http" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/17/83/8e8e83b7ac285281687c7be2fd305213ccccbb8c0a2dd4fb45a8ccaf12c7/opentelemetry_instrumentation_asgi-0.65b0.tar.gz", hash = "sha256:892bca67c56522ffa85a8a83cf934d7b50b3be2132e45cbee705825f0a5ba426", size = 26140, upload-time = "2026-07-16T15:25:54.544Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/0b/9c/376962840b619d2d55fe8ee2285f8c70971c090e5fff614516fc654a6f3a/opentelemetry_instrumentation_asgi-0.65b0-py3-none-any.whl", hash = "sha256:3a845a8ebd1c4ef0d8263401e6545f5b219b2feee612090d50f578a87e71fd65", size = 15903, upload-time = "2026-07-16T15:24:57.198Z" }, +] + [[package]] name = "opentelemetry-sdk" version = "1.44.0" @@ -2479,6 +2507,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/a6/0e/49df70d9b81fb5cbae4bbf2a49d865b09bcbcbc4eb53f5851b1027738d78/opentelemetry_semantic_conventions-0.65b0-py3-none-any.whl", hash = "sha256:1cacde7b0ad306f84c5ef08c3dbe1bbaf20165bba6f8bff43b670e555a086bcb", size = 204645, upload-time = "2026-07-16T15:25:30.688Z" }, ] +[[package]] +name = "opentelemetry-util-http" +version = "0.65b0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/32/a9/d7525a59fdd240e69b5af4a6338e78fafa1b4203394122cbd6701fb5f84a/opentelemetry_util_http-0.65b0.tar.gz", hash = "sha256:84f82d826978bba416ab453460ff6a7391cdc3534c93a786595e4068680016b7", size = 11243, upload-time = "2026-07-16T15:26:27.898Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/23/3f/ab8d29df207ce5f470a07fa96ebb48af4e95b7fab7e7635311b9a32f2fab/opentelemetry_util_http-0.65b0-py3-none-any.whl", hash = "sha256:7553b606f963097cb190536dc30556cce85090692e471a422fff30ca29b04348", size = 8245, upload-time = "2026-07-16T15:25:46.482Z" }, +] + [[package]] name = "orjson" version = "3.11.9" @@ -4216,6 +4253,8 @@ source = { editable = "packages/reflex-otel" } dependencies = [ { name = "opentelemetry-api" }, { name = "opentelemetry-instrumentation" }, + { name = "opentelemetry-instrumentation-asgi" }, + { name = "opentelemetry-util-http" }, { name = "reflex-base" }, ] @@ -4227,7 +4266,9 @@ instruments = [ [package.metadata] requires-dist = [ { name = "opentelemetry-api", specifier = ">=1.24.0,<2.0" }, - { name = "opentelemetry-instrumentation", specifier = ">=0.45b0,<1.0" }, + { name = "opentelemetry-instrumentation", specifier = ">=0.49b0,<1.0" }, + { name = "opentelemetry-instrumentation-asgi", specifier = ">=0.49b0,<1.0" }, + { name = "opentelemetry-util-http", specifier = ">=0.49b0,<1.0" }, { name = "reflex-base", editable = "packages/reflex-base" }, { name = "reflex-base", marker = "extra == 'instruments'", editable = "packages/reflex-base" }, ] From fc39a8b9125c3d8eacd5f4cde0accbe002696b93 Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 00:45:08 +0500 Subject: [PATCH 3/8] fix(otel): seconds/bytes bucket advisories on the histograms (opentelemetry-api >=1.30) --- packages/reflex-base/pyproject.toml | 2 +- packages/reflex-base/src/reflex_base/otel.py | 22 ++++++++++++++++++++ packages/reflex-otel/pyproject.toml | 2 +- uv.lock | 4 ++-- 4 files changed, 26 insertions(+), 4 deletions(-) diff --git a/packages/reflex-base/pyproject.toml b/packages/reflex-base/pyproject.toml index 68bec385afb..310626adce4 100644 --- a/packages/reflex-base/pyproject.toml +++ b/packages/reflex-base/pyproject.toml @@ -8,7 +8,7 @@ authors = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] maintainers = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] requires-python = ">=3.10" dependencies = [ - "opentelemetry-api >=1.24.0,<2.0", + "opentelemetry-api >=1.30.0,<2.0", "packaging >=24.2,<27", "rich >=13,<16", "typing_extensions >=4.13.0", diff --git a/packages/reflex-base/src/reflex_base/otel.py b/packages/reflex-base/src/reflex_base/otel.py index e971ae9a70a..2809bd44baa 100644 --- a/packages/reflex-base/src/reflex_base/otel.py +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -52,6 +52,25 @@ # Key of the W3C trace context carried in an event payload sent by the frontend. TRACEPARENT_FIELD = "traceparent" +# Histogram bucket advisories: seconds (semconv http.server.request.duration) and bytes. +_DURATION_BUCKETS = ( + 0.005, + 0.01, + 0.025, + 0.05, + 0.075, + 0.1, + 0.25, + 0.5, + 0.75, + 1, + 2.5, + 5, + 7.5, + 10, +) +_SIZE_BUCKETS = (128, 512, 1024, 4096, 16384, 65536, 262144, 1048576, 4194304) + # Read at every trace point; True only after enable() ran. enabled: bool = False # Wraps the app's ASGI callable when set (installed by enable()). @@ -76,16 +95,19 @@ def _create_instruments(meter: metrics.Meter) -> None: METRIC_EVENT_DURATION, unit="s", description="Duration of event handler executions.", + explicit_bucket_boundaries_advisory=_DURATION_BUCKETS, ) _state_acquire_duration = meter.create_histogram( METRIC_STATE_ACQUIRE_DURATION, unit="s", description="Time an event waited to acquire and load its session state.", + explicit_bucket_boundaries_advisory=_DURATION_BUCKETS, ) _message_size = meter.create_histogram( METRIC_WEBSOCKET_MESSAGE_SIZE, unit="By", description="Serialized size of socket messages exchanged with the client.", + explicit_bucket_boundaries_advisory=_SIZE_BUCKETS, ) _ws_connections = meter.create_up_down_counter( METRIC_WEBSOCKET_CONNECTIONS, diff --git a/packages/reflex-otel/pyproject.toml b/packages/reflex-otel/pyproject.toml index 4c693ffa3d8..9a8821f7345 100644 --- a/packages/reflex-otel/pyproject.toml +++ b/packages/reflex-otel/pyproject.toml @@ -8,7 +8,7 @@ authors = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] maintainers = [{ name = "Khaleel Al-Adhami", email = "khaleel@reflex.dev" }] requires-python = ">=3.10" dependencies = [ - "opentelemetry-api >=1.24.0,<2.0", + "opentelemetry-api >=1.30.0,<2.0", "opentelemetry-instrumentation >=0.49b0,<1.0", "opentelemetry-instrumentation-asgi >=0.49b0,<1.0", "opentelemetry-util-http >=0.49b0,<1.0", diff --git a/uv.lock b/uv.lock index b1a0cbd61bc..ef71d88d6c3 100644 --- a/uv.lock +++ b/uv.lock @@ -3921,7 +3921,7 @@ pydantic = [ [package.metadata] requires-dist = [ - { name = "opentelemetry-api", specifier = ">=1.24.0,<2.0" }, + { name = "opentelemetry-api", specifier = ">=1.30.0,<2.0" }, { name = "packaging", specifier = ">=24.2,<27" }, { name = "platformdirs", specifier = ">=4.3.7,<5.0" }, { name = "pydantic", marker = "extra == 'pydantic'", specifier = ">=2.12.0,<3.0" }, @@ -4265,7 +4265,7 @@ instruments = [ [package.metadata] requires-dist = [ - { name = "opentelemetry-api", specifier = ">=1.24.0,<2.0" }, + { name = "opentelemetry-api", specifier = ">=1.30.0,<2.0" }, { name = "opentelemetry-instrumentation", specifier = ">=0.49b0,<1.0" }, { name = "opentelemetry-instrumentation-asgi", specifier = ">=0.49b0,<1.0" }, { name = "opentelemetry-util-http", specifier = ">=0.49b0,<1.0" }, From 54ec281c65897672b11d42fcbdf7436fddddb320 Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 01:18:51 +0500 Subject: [PATCH 4/8] =?UTF-8?q?refactor(otel):=20review=20fixes=20?= =?UTF-8?q?=E2=80=94=20span=20kinds,=20byte=20sizes,=20lazy=20ASGI=20impor?= =?UTF-8?q?t,=20/ping=20excluded,=20shared=20test=20fixtures?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../event/processor/base_state_processor.py | 5 +- packages/reflex-base/src/reflex_base/otel.py | 70 +++++++++++------- packages/reflex-otel/README.md | 31 +++++--- packages/reflex-otel/pyproject.toml | 1 - .../reflex-otel/src/reflex_otel/__init__.py | 25 ++++--- reflex/app.py | 9 ++- tests/units/conftest.py | 73 ++++++++++++++----- .../processor/test_base_state_processor.py | 16 ++-- .../event/processor/test_event_processor.py | 15 ++-- tests/units/reflex_base/test_otel.py | 53 +++++++------- tests/units/test_app.py | 59 +++++---------- uv.lock | 2 - 12 files changed, 210 insertions(+), 149 deletions(-) diff --git a/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py index 8f0449b42a4..efe6d60fa7f 100644 --- a/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py +++ b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py @@ -351,7 +351,7 @@ async def _execute_event( ), event=entry.event, ) as state: - if acquire_start: + if otel.enabled: otel.record_state_acquired(acquire_start, event) # Compatibility hack rehydrate the state before processing this event. needs_to_rehydrate = bool( @@ -416,6 +416,9 @@ async def _handle_backend_exception( if ev_ctx is not None: # Ensure the event context is set for the exception handler. EventContext.set(ev_ctx) + if otel.enabled: + # Keep events chained by the handler in the failed event's trace. + otel.attach_context(ev_ctx.otel_context) if events := self.backend_exception_handler(ex): await chain_updates( events=events, diff --git a/packages/reflex-base/src/reflex_base/otel.py b/packages/reflex-base/src/reflex_base/otel.py index 2809bd44baa..bbb8d083249 100644 --- a/packages/reflex-base/src/reflex_base/otel.py +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -12,8 +12,8 @@ from __future__ import annotations -from collections.abc import Callable, Iterator, Mapping -from contextlib import contextmanager, nullcontext +from collections.abc import Awaitable, Callable, Iterator, Mapping +from contextlib import contextmanager from time import perf_counter from typing import TYPE_CHECKING, Any @@ -25,8 +25,6 @@ from reflex_base.constants.base import Reflex if TYPE_CHECKING: - from contextlib import AbstractContextManager - from reflex_base.event import Event from reflex_base.event.context import EventContext from reflex_base.registry import RegisteredEventHandler @@ -71,10 +69,13 @@ ) _SIZE_BUCKETS = (128, 512, 1024, 4096, 16384, 65536, 262144, 1048576, 4194304) +# An ASGI callable: (scope, receive, send) -> awaitable. +ASGIApp = Callable[..., Awaitable[None]] + # Read at every trace point; True only after enable() ran. enabled: bool = False # Wraps the app's ASGI callable when set (installed by enable()). -asgi_middleware: Callable[[Any], Any] | None = None +asgi_middleware: Callable[[ASGIApp], ASGIApp] | None = None _tracer: trace.Tracer = trace.NoOpTracer() _noop_meter = metrics.NoOpMeter(INSTRUMENTATION_NAME) @@ -119,7 +120,7 @@ def _create_instruments(meter: metrics.Meter) -> None: def enable( tracer_provider: trace.TracerProvider | None = None, meter_provider: metrics.MeterProvider | None = None, - asgi_middleware_factory: Callable[[Any], Any] | None = None, + asgi_middleware_factory: Callable[[ASGIApp], ASGIApp] | None = None, ) -> None: """Turn the trace points and metrics on. @@ -164,6 +165,16 @@ def capture_context() -> Context | None: return otel_context.get_current() if enabled else None +def attach_context(context: Context | None) -> None: + """Make a captured context current for the rest of the running task. + + Args: + context: The context to attach; None is ignored. + """ + if context is not None: + otel_context.attach(context) + + class _AttachedContext: """Attach a context on enter and detach it on exit.""" @@ -190,20 +201,19 @@ def __exit__(self, *exc_info: object) -> None: otel_context.detach(self._token) -def remote_context(carrier: Mapping[str, Any]) -> AbstractContextManager[None]: +def remote_context(carrier: Mapping[str, Any]) -> _AttachedContext: """Make the trace context carried by a frontend event the current context. - Events that carry no ``traceparent`` start a new trace: the websocket - connection span (if any) is deliberately not used as their parent. + Only call when ``enabled``. Events that carry no ``traceparent`` start a + new trace: the websocket connection span (if any) is deliberately not + used as their parent. Args: carrier: The raw event fields received from the frontend. Returns: - A context manager to run the enqueue under; a no-op when tracing is off. + A context manager to run the enqueue under. """ - if not enabled: - return nullcontext() return _AttachedContext(propagate.extract(carrier, context=Context())) @@ -213,8 +223,9 @@ def event_span( ) -> Iterator[trace.Span]: """Open the span for one event handler execution and record its duration. - Chained events are parented under the span that enqueued them; events - that arrive from the frontend start a new trace. + Chained events are INTERNAL children of the span that enqueued them; + events that arrive from the frontend are SERVER spans (a new trace, or a + child of the browser's remote span). Args: event: The event being processed. @@ -229,24 +240,29 @@ def event_span( ATTR_EVENT_NAME: event.name, ATTR_EVENT_BACKGROUND: handler.is_background, } + attributes: dict[str, Any] = { + **metric_attributes, + ATTR_EVENT_TXID: ctx.txid, + ATTR_SESSION_ID: ctx.token, + ATTR_CODE_FUNCTION_NAME: getattr(handler.fn, "__qualname__", event.name), + } + if ctx.parent_txid: + attributes[ATTR_EVENT_PARENT_TXID] = ctx.parent_txid + # An event whose parent span is local (a chained event) is an internal + # step of that request; anything else is a new inbound request. + parent = trace.get_current_span(ctx.otel_context).get_span_context() + kind = ( + SpanKind.INTERNAL + if parent.is_valid and not parent.is_remote + else SpanKind.SERVER + ) start = perf_counter() try: with _tracer.start_as_current_span( - event.name, - context=ctx.otel_context, - kind=SpanKind.SERVER, - attributes=metric_attributes - | { - ATTR_EVENT_TXID: ctx.txid, - ATTR_SESSION_ID: ctx.token, - ATTR_CODE_FUNCTION_NAME: getattr( - handler.fn, "__qualname__", event.name - ), - } - | ({ATTR_EVENT_PARENT_TXID: ctx.parent_txid} if ctx.parent_txid else {}), + event.name, context=ctx.otel_context, kind=kind, attributes=attributes ) as span: yield span - except BaseException as ex: + except Exception as ex: metric_attributes[ATTR_ERROR_TYPE] = type(ex).__qualname__ raise finally: diff --git a/packages/reflex-otel/README.md b/packages/reflex-otel/README.md index 874f8a8f506..7dcdcc5bb1a 100644 --- a/packages/reflex-otel/README.md +++ b/packages/reflex-otel/README.md @@ -9,18 +9,29 @@ ReflexInstrumentor().instrument() ``` The package registers an `opentelemetry_instrumentor` entry point, so -`opentelemetry-instrument reflex run` enables it automatically. Configure a -tracer provider and a meter provider (for example with `opentelemetry-sdk`) -to export the data. +`opentelemetry-instrument reflex run` enables it too (the auto-instrumentation +`sitecustomize` reaches the backend worker through the inherited `PYTHONPATH`). +Configure a tracer provider and a meter provider (for example with +`opentelemetry-sdk`) to export the data. + +Reflex hot reload re-imports the app module in-process, so guard one-time SDK +setup: + +```python +if not ReflexInstrumentor().is_instrumented_by_opentelemetry: + trace.set_tracer_provider(provider) + ReflexInstrumentor().instrument(tracer_provider=provider) +``` ## What you get Traces: -- One `SERVER` span per event handler run, named after the event. Chained - events are children of the span that enqueued them. An event sent by the - frontend with a `traceparent` field continues that trace; otherwise it - starts a new one. +- One span per event handler run, named after the event: `SERVER` for events + sent by the frontend (a new trace, or a child of the browser span when the + event carries a `traceparent` field), `INTERNAL` for chained events, which + are children of the span that enqueued them. `traceparent`/`tracestate` + are consumed and never reach the handler. - HTTP requests and the websocket connection are wrapped in the standard OpenTelemetry ASGI middleware (per-message websocket spans are off). @@ -30,7 +41,7 @@ Metrics: | --- | --- | --- | --- | | `reflex.event.duration` | histogram | s | `reflex.event.name`, `reflex.event.background`, `error.type` | | `reflex.state.acquire.duration` | histogram | s | `reflex.event.name` | -| `reflex.websocket.message.size` | histogram | By | `network.io.direction` (`transmit`/`receive`) | +| `reflex.websocket.message.size` | histogram | By | `network.io.direction` (`transmit`/`receive`); default `sio` only | | `reflex.websocket.connections` | up-down counter | `{connection}` | | Plus the ASGI middleware's `http.server.*` metrics. @@ -38,6 +49,6 @@ Plus the ASGI middleware's `http.server.*` metrics. ## Options `instrument()` accepts `tracer_provider`, `meter_provider`, `excluded_urls` -(comma-separated URL patterns skipped by the ASGI middleware; defaults to the -`OTEL_PYTHON_REFLEX_EXCLUDED_URLS` environment variable) and the ASGI hooks +(comma-separated URL patterns skipped by the ASGI middleware; defaults to +`OTEL_PYTHON_REFLEX_EXCLUDED_URLS`, else `/ping`) and the ASGI hooks `server_request_hook`, `client_request_hook`, `client_response_hook`. diff --git a/packages/reflex-otel/pyproject.toml b/packages/reflex-otel/pyproject.toml index 9a8821f7345..df219b2bddc 100644 --- a/packages/reflex-otel/pyproject.toml +++ b/packages/reflex-otel/pyproject.toml @@ -11,7 +11,6 @@ dependencies = [ "opentelemetry-api >=1.30.0,<2.0", "opentelemetry-instrumentation >=0.49b0,<1.0", "opentelemetry-instrumentation-asgi >=0.49b0,<1.0", - "opentelemetry-util-http >=0.49b0,<1.0", "reflex-base >= 0.9.7.post45.dev0", ] diff --git a/packages/reflex-otel/src/reflex_otel/__init__.py b/packages/reflex-otel/src/reflex_otel/__init__.py index 95aef1e90e9..316ab480c29 100644 --- a/packages/reflex-otel/src/reflex_otel/__init__.py +++ b/packages/reflex-otel/src/reflex_otel/__init__.py @@ -2,18 +2,19 @@ from __future__ import annotations +import os from collections.abc import Collection from typing import Any, Literal -from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware from opentelemetry.instrumentation.instrumentor import BaseInstrumentor -from opentelemetry.util.http import get_excluded_urls from reflex_base import otel _instruments = ("reflex-base >= 0.9.7.post45.dev0",) # Per-message websocket spans are noise; Reflex emits one span per event instead. _ASGI_EXCLUDED_SPANS: list[Literal["receive", "send"]] = ["receive", "send"] +# Frontend health polling; override with excluded_urls / OTEL_PYTHON_REFLEX_EXCLUDED_URLS. +_DEFAULT_EXCLUDED_URLS = "/ping" class ReflexInstrumentor(BaseInstrumentor): @@ -47,22 +48,26 @@ def _instrument(self, **kwargs: Any) -> None: **kwargs: ``tracer_provider`` and ``meter_provider`` select the providers (default: the global ones). ``excluded_urls`` is a comma-separated list of URL patterns the ASGI middleware skips - (default: ``OTEL_PYTHON_REFLEX_EXCLUDED_URLS``). + (default: ``OTEL_PYTHON_REFLEX_EXCLUDED_URLS``, else ``/ping``). ``server_request_hook``, ``client_request_hook`` and ``client_response_hook`` are forwarded to the ASGI middleware. """ + # Imported here so `from reflex_otel import OtelPlugin` in rxconfig.py + # stays cheap for CLI processes that never instrument. + from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware + tracer_provider = kwargs.get("tracer_provider") meter_provider = kwargs.get("meter_provider") - excluded_urls = kwargs.get("excluded_urls") + excluded_urls = kwargs.get("excluded_urls") or ( + os.environ.get("OTEL_PYTHON_REFLEX_EXCLUDED_URLS") + or os.environ.get("OTEL_PYTHON_EXCLUDED_URLS") + or _DEFAULT_EXCLUDED_URLS + ) - def asgi_middleware(app: Any) -> Any: + def asgi_middleware(app: otel.ASGIApp) -> otel.ASGIApp: return OpenTelemetryMiddleware( app, - excluded_urls=( - get_excluded_urls("REFLEX") - if excluded_urls is None - else excluded_urls - ), + excluded_urls=excluded_urls, server_request_hook=kwargs.get("server_request_hook"), client_request_hook=kwargs.get("client_request_hook"), client_response_hook=kwargs.get("client_response_hook"), diff --git a/reflex/app.py b/reflex/app.py index fc82ca50a27..aa21898bde6 100644 --- a/reflex/app.py +++ b/reflex/app.py @@ -1931,7 +1931,7 @@ def _sio_dumps(obj: Any, **kwargs: Any) -> str: """ data = format.json_dumps(obj, **kwargs) if otel.enabled: - otel.record_message_size(len(data), "transmit") + otel.record_message_size(len(data.encode()), "transmit") return data @@ -1946,7 +1946,9 @@ def _sio_loads(data: str | bytes, **kwargs: Any) -> Any: The decoded payload. """ if otel.enabled: - otel.record_message_size(len(data), "receive") + otel.record_message_size( + len(data.encode()) if isinstance(data, str) else len(data), "receive" + ) return json.loads(data, **kwargs) @@ -2178,6 +2180,9 @@ async def on_event(self, sid: str, data: Any): if (path := router_data.get(constants.RouteVar.PATH)) else "404" ).removeprefix("/") + if not otel.enabled: + await self.app.event_processor.enqueue(token, event) + return with otel.remote_context(fields): await self.app.event_processor.enqueue(token, event) diff --git a/tests/units/conftest.py b/tests/units/conftest.py index c2d215cc358..60012f2fed4 100644 --- a/tests/units/conftest.py +++ b/tests/units/conftest.py @@ -526,32 +526,71 @@ def preserve_memo_registries(): @pytest.fixture -def otel_exporter() -> Generator[InMemorySpanExporter, None, None]: - """Enable the reflex_base.otel trace points against an in-memory exporter. +def otel_sdk() -> Generator[ + tuple[InMemorySpanExporter, InMemoryMetricReader], None, None +]: + """Enable the reflex_base.otel trace points and metrics against in-memory sinks. Yields: - The exporter collecting finished spans. + The span exporter and the metric reader. """ exporter = InMemorySpanExporter() - provider = TracerProvider() - provider.add_span_processor(SimpleSpanProcessor(exporter)) - otel.enable(tracer_provider=provider) + tracer_provider = TracerProvider() + tracer_provider.add_span_processor(SimpleSpanProcessor(exporter)) + reader = InMemoryMetricReader() + otel.enable( + tracer_provider=tracer_provider, + meter_provider=MeterProvider(metric_readers=[reader]), + ) try: - yield exporter + yield exporter, reader finally: otel.disable() @pytest.fixture -def otel_metrics() -> Generator[InMemoryMetricReader, None, None]: - """Enable the reflex_base.otel metrics against an in-memory reader. +def otel_exporter(otel_sdk) -> InMemorySpanExporter: + """The in-memory span exporter of the enabled otel_sdk. - Yields: - The reader collecting recorded metrics. + Args: + otel_sdk: The enabled sinks. + + Returns: + The span exporter. """ - reader = InMemoryMetricReader() - otel.enable(meter_provider=MeterProvider(metric_readers=[reader])) - try: - yield reader - finally: - otel.disable() + return otel_sdk[0] + + +@pytest.fixture +def otel_metrics(otel_sdk) -> InMemoryMetricReader: + """The in-memory metric reader of the enabled otel_sdk. + + Args: + otel_sdk: The enabled sinks. + + Returns: + The metric reader. + """ + return otel_sdk[1] + + +def metric_points(reader: InMemoryMetricReader, name: str) -> list: + """Collect the data points recorded for one metric. + + Args: + reader: The in-memory reader to collect from. + name: The metric name. + + Returns: + The data points, in recording order. + """ + data = reader.get_metrics_data() + assert data is not None + return [ + point + for rm in data.resource_metrics + for sm in rm.scope_metrics + for metric in sm.metrics + if metric.name == name + for point in metric.data.data_points + ] diff --git a/tests/units/reflex_base/event/processor/test_base_state_processor.py b/tests/units/reflex_base/event/processor/test_base_state_processor.py index 1243f6df6dd..37aa0b53066 100644 --- a/tests/units/reflex_base/event/processor/test_base_state_processor.py +++ b/tests/units/reflex_base/event/processor/test_base_state_processor.py @@ -6,6 +6,7 @@ import pytest import pytest_asyncio +from reflex_base import otel from reflex_base.constants import CompileVars from reflex_base.constants.state import FIELD_MARKER from reflex_base.event.context import EventContext @@ -19,6 +20,7 @@ from reflex.istate.manager.memory import StateManagerMemory from reflex.middleware.middleware import Middleware from reflex.state import OnLoadInternalState, State, StateUpdate +from tests.units.conftest import metric_points @pytest.fixture @@ -230,7 +232,6 @@ async def test_execute_event_records_state_acquire_duration( token: The client token. otel_metrics: In-memory metric reader with metrics enabled. """ - from reflex_base import otel class AcquireState(State): @event @@ -241,13 +242,8 @@ def noop(self): await processor.enqueue(token, Event.from_event_type(AcquireState.noop())[0]) await processor.join(1) - metrics = otel_metrics.get_metrics_data() - (metric,) = [ - m - for rm in metrics.resource_metrics - for sm in rm.scope_metrics - for m in sm.metrics - if m.name == otel.METRIC_STATE_ACQUIRE_DURATION - ] - names = {p.attributes[otel.ATTR_EVENT_NAME] for p in metric.data.data_points} + names = { + p.attributes[otel.ATTR_EVENT_NAME] + for p in metric_points(otel_metrics, otel.METRIC_STATE_ACQUIRE_DURATION) + } assert Event.from_event_type(AcquireState.noop())[0].name in names diff --git a/tests/units/reflex_base/event/processor/test_event_processor.py b/tests/units/reflex_base/event/processor/test_event_processor.py index fba1e5d7333..acc3f827e04 100644 --- a/tests/units/reflex_base/event/processor/test_event_processor.py +++ b/tests/units/reflex_base/event/processor/test_event_processor.py @@ -3,9 +3,12 @@ import asyncio import contextlib from typing import Any +from unittest.mock import Mock import pytest +from opentelemetry.trace import SpanKind from pytest_mock import MockerFixture +from reflex_base import otel from reflex_base.event.context import EventContext from reflex_base.event.processor.event_processor import ( EventProcessor, @@ -795,19 +798,21 @@ async def _watcher(): # noqa: RUF029 async def test_no_spans_when_otel_disabled( - mock_event_processor: EventProcessor, token: str + mock_event_processor: EventProcessor, token: str, monkeypatch ): """With tracing off the processor never touches the tracer. Args: mock_event_processor: The event processor with mock root context. token: The client token. + monkeypatch: Pytest monkeypatch fixture. """ - from reflex_base import otel - assert otel.enabled is False + tracer = Mock() + monkeypatch.setattr(otel, "_tracer", tracer) async with mock_event_processor as ep: await ep.enqueue(token, Event.from_event_type(noop_event())[0]) + tracer.start_as_current_span.assert_not_called() async def test_event_spans_chain_parent_child(token: str, otel_exporter): @@ -817,8 +822,6 @@ async def test_event_spans_chain_parent_child(token: str, otel_exporter): token: The client token. otel_exporter: In-memory span exporter with tracing enabled. """ - from reflex_base import otel - ep = EventProcessor(graceful_shutdown_timeout=2) ep.configure() async with ep: @@ -828,8 +831,10 @@ async def test_event_spans_chain_parent_child(token: str, otel_exporter): parent = spans["_chaining_handler"] child = spans["_logging_handler"] assert parent.parent is None + assert parent.kind == SpanKind.SERVER assert child.parent is not None assert child.parent.span_id == parent.context.span_id + assert child.kind == SpanKind.INTERNAL assert ( child.attributes[otel.ATTR_EVENT_PARENT_TXID] == parent.attributes[otel.ATTR_EVENT_TXID] diff --git a/tests/units/reflex_base/test_otel.py b/tests/units/reflex_base/test_otel.py index 6d95cda6f59..fce75140d9a 100644 --- a/tests/units/reflex_base/test_otel.py +++ b/tests/units/reflex_base/test_otel.py @@ -1,6 +1,6 @@ """Tests for the reflex_base.otel trace points.""" -from contextlib import nullcontext +import asyncio from time import perf_counter import pytest @@ -14,6 +14,7 @@ from reflex_base.registry import RegisteredEventHandler from reflex.event import Event, EventHandler +from tests.units.conftest import metric_points def _ctx(token: str = "tok", parent_txid: str | None = None) -> EventContext: @@ -97,19 +98,7 @@ def test_event_span_parents_under_captured_context( _root, child = otel_exporter.get_finished_spans() assert child.parent is not None assert child.parent.span_id == root.get_span_context().span_id - - -def _metric_points(reader: InMemoryMetricReader, name: str) -> list: - data = reader.get_metrics_data() - assert data is not None - return [ - point - for rm in data.resource_metrics - for sm in rm.scope_metrics - for metric in sm.metrics - if metric.name == name - for point in metric.data.data_points - ] + assert child.kind == SpanKind.INTERNAL def test_event_span_records_duration_metric(otel_metrics: InMemoryMetricReader): @@ -121,10 +110,16 @@ def test_event_span_records_duration_metric(otel_metrics: InMemoryMetricReader): otel.event_span(Event(name="bad"), _ctx(), registered), ): raise RuntimeError - bad, ok = sorted( - _metric_points(otel_metrics, otel.METRIC_EVENT_DURATION), + with ( + pytest.raises(asyncio.CancelledError), + otel.event_span(Event(name="cancelled"), _ctx(), registered), + ): + raise asyncio.CancelledError + bad, cancelled, ok = sorted( + metric_points(otel_metrics, otel.METRIC_EVENT_DURATION), key=lambda p: p.attributes[otel.ATTR_EVENT_NAME], ) + assert otel.ATTR_ERROR_TYPE not in cancelled.attributes assert bad.attributes == { otel.ATTR_EVENT_NAME: "bad", otel.ATTR_EVENT_BACKGROUND: False, @@ -144,21 +139,15 @@ def test_metric_helpers_record(otel_metrics: InMemoryMetricReader): otel.record_connection(1) otel.record_connection(1) otel.record_connection(-1) - (acquire,) = _metric_points(otel_metrics, otel.METRIC_STATE_ACQUIRE_DURATION) + (acquire,) = metric_points(otel_metrics, otel.METRIC_STATE_ACQUIRE_DURATION) assert acquire.attributes == {otel.ATTR_EVENT_NAME: "e"} - (size,) = _metric_points(otel_metrics, otel.METRIC_WEBSOCKET_MESSAGE_SIZE) + (size,) = metric_points(otel_metrics, otel.METRIC_WEBSOCKET_MESSAGE_SIZE) assert size.sum == 42 assert size.attributes == {otel.ATTR_NETWORK_IO_DIRECTION: "transmit"} - (conns,) = _metric_points(otel_metrics, otel.METRIC_WEBSOCKET_CONNECTIONS) + (conns,) = metric_points(otel_metrics, otel.METRIC_WEBSOCKET_CONNECTIONS) assert conns.value == 1 -def test_remote_context_disabled_is_noop(): - with otel._tracer.start_as_current_span("outer"), otel.remote_context({}): - pass - assert isinstance(otel.remote_context({"traceparent": "x"}), nullcontext) - - def test_remote_context_uses_traceparent(otel_exporter: InMemorySpanExporter): registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) traceparent = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01" @@ -167,6 +156,7 @@ def test_remote_context_uses_traceparent(otel_exporter: InMemorySpanExporter): with otel.event_span(Event(name="e"), ctx, registered): pass (span,) = otel_exporter.get_finished_spans() + assert span.kind == SpanKind.SERVER assert span.parent is not None assert span.parent.is_remote assert format(span.parent.trace_id, "032x") == "0af7651916cd43dd8448eb211c80319c" @@ -192,3 +182,16 @@ def test_asgi_middleware_hook_toggles(): assert otel.asgi_middleware is factory otel.disable() assert otel.asgi_middleware is None + + +def test_attach_context(otel_exporter: InMemorySpanExporter): + with otel._tracer.start_as_current_span("outer") as outer: + captured = otel.capture_context() + otel.attach_context(None) + assert not trace.get_current_span().get_span_context().is_valid + token = otel_context.attach(otel_context.get_current()) + try: + otel.attach_context(captured) + assert trace.get_current_span() is outer + finally: + otel_context.detach(token) diff --git a/tests/units/test_app.py b/tests/units/test_app.py index 41ee6a5871f..96d43a15c63 100644 --- a/tests/units/test_app.py +++ b/tests/units/test_app.py @@ -18,7 +18,9 @@ import pytest import reflex_base +from opentelemetry import trace from pytest_mock import MockerFixture +from reflex_base import otel from reflex_base.components.component import Component from reflex_base.constants.state import FIELD_MARKER from reflex_base.event import Event @@ -43,7 +45,14 @@ import reflex as rx from reflex import AdminDash, constants from reflex._upload import upload -from reflex.app import App, ComponentCallable, EventNamespace, default_overlay_component +from reflex.app import ( + App, + ComponentCallable, + EventNamespace, + _sio_dumps, + _sio_loads, + default_overlay_component, +) from reflex.compiler.compiler import ( _compile_app, _memoize_stateful_app_wraps, @@ -61,7 +70,7 @@ from reflex.state import BaseState, OnLoadInternalState, State, reload_state_module from reflex.utils import exec as exec_utils -from .conftest import chdir +from .conftest import chdir, metric_points from .states import GenState from .states.upload import ( ChildFileUploadState, @@ -4261,8 +4270,6 @@ def test_client_error_constants_match_frontend(): def test_call_app_wraps_with_otel_asgi_middleware(): """The app's ASGI callable is wrapped when instrumentation installs a middleware.""" - from reflex_base import otel - app = App() app._compile = unittest.mock.Mock() wrapped = [] @@ -4276,36 +4283,22 @@ def test_call_app_wraps_with_otel_asgi_middleware(): def test_sio_json_records_message_sizes(otel_metrics): """Socket.IO packet serialization records sizes in both directions.""" - from reflex_base import otel - - from reflex.app import _sio_dumps, _sio_loads - - data = _sio_dumps({"a": 1}, separators=(",", ":")) - assert data == '{"a":1}' - assert _sio_loads(data) == {"a": 1} - metrics = otel_metrics.get_metrics_data() - (metric,) = [ - m - for rm in metrics.resource_metrics - for sm in rm.scope_metrics - for m in sm.metrics - if m.name == otel.METRIC_WEBSOCKET_MESSAGE_SIZE - ] + data = _sio_dumps({"a": "é"}, separators=(",", ":")) + assert data == '{"a":"é"}' + assert _sio_loads(data) == {"a": "é"} + assert _sio_loads(data.encode()) == {"a": "é"} points = { p.attributes[otel.ATTR_NETWORK_IO_DIRECTION]: p.sum - for p in metric.data.data_points + for p in metric_points(otel_metrics, otel.METRIC_WEBSOCKET_MESSAGE_SIZE) } - assert points == {"transmit": len(data), "receive": len(data)} + size = len(data.encode()) + assert size == len(data) + 1 + assert points == {"transmit": size, "receive": 2 * size} @pytest.mark.asyncio async def test_on_event_uses_frontend_traceparent(otel_exporter): """A traceparent in the event payload becomes the parent of the event span.""" - from opentelemetry import trace - from reflex_base import otel - - from reflex.app import EventNamespace - mock_app = unittest.mock.Mock() mock_app.router.return_value = "/" mock_app.sio.get_environ.return_value = { @@ -4333,10 +4326,6 @@ async def enqueue(token, event): @pytest.mark.asyncio async def test_connect_disconnect_counts_connections(otel_metrics): """Connect and disconnect adjust the open connection gauge.""" - from reflex_base import otel - - from reflex.app import EventNamespace - mock_app = unittest.mock.Mock() mock_app._state = None ns = EventNamespace(namespace="/", app=mock_app) @@ -4346,13 +4335,5 @@ async def test_connect_disconnect_counts_connections(otel_metrics): task = ns.on_disconnect("sid1") if task is not None: await task - metrics = otel_metrics.get_metrics_data() - (metric,) = [ - m - for rm in metrics.resource_metrics - for sm in rm.scope_metrics - for m in sm.metrics - if m.name == otel.METRIC_WEBSOCKET_CONNECTIONS - ] - (point,) = metric.data.data_points + (point,) = metric_points(otel_metrics, otel.METRIC_WEBSOCKET_CONNECTIONS) assert point.value == 1 diff --git a/uv.lock b/uv.lock index ef71d88d6c3..1f5421157fd 100644 --- a/uv.lock +++ b/uv.lock @@ -4254,7 +4254,6 @@ dependencies = [ { name = "opentelemetry-api" }, { name = "opentelemetry-instrumentation" }, { name = "opentelemetry-instrumentation-asgi" }, - { name = "opentelemetry-util-http" }, { name = "reflex-base" }, ] @@ -4268,7 +4267,6 @@ requires-dist = [ { name = "opentelemetry-api", specifier = ">=1.30.0,<2.0" }, { name = "opentelemetry-instrumentation", specifier = ">=0.49b0,<1.0" }, { name = "opentelemetry-instrumentation-asgi", specifier = ">=0.49b0,<1.0" }, - { name = "opentelemetry-util-http", specifier = ">=0.49b0,<1.0" }, { name = "reflex-base", editable = "packages/reflex-base" }, { name = "reflex-base", marker = "extra == 'instruments'", editable = "packages/reflex-base" }, ] From 7b44ff8e27b9f42370f827966df2819d1c6346e6 Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 02:34:32 +0500 Subject: [PATCH 5/8] test(otel): release websocket tokens after connection-count test --- tests/units/test_app.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tests/units/test_app.py b/tests/units/test_app.py index 96d43a15c63..773b463044f 100644 --- a/tests/units/test_app.py +++ b/tests/units/test_app.py @@ -4337,3 +4337,5 @@ async def test_connect_disconnect_counts_connections(otel_metrics): await task (point,) = metric_points(otel_metrics, otel.METRIC_WEBSOCKET_CONNECTIONS) assert point.value == 1 + # Release t2 so a shared token store (redis) does not leak into other tests. + await ns._token_manager.disconnect_all() From 5af4670a43b42764a75034ca4333e7bbc4b36a2a Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 02:43:18 +0500 Subject: [PATCH 6/8] fix(otel): honour excluded_urls="", cover reflex_otel, tidy docs and changelog --- .../event/processor/base_state_processor.py | 3 ++- packages/reflex-base/src/reflex_base/otel.py | 4 ++-- packages/reflex-otel/CHANGELOG.md | 2 +- packages/reflex-otel/README.md | 3 ++- .../reflex-otel/src/reflex_otel/__init__.py | 16 ++++++++----- pyproject.toml | 1 + tests/units/reflex_otel/test_init.py | 23 +++++++++++++++++++ 7 files changed, 41 insertions(+), 11 deletions(-) diff --git a/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py index efe6d60fa7f..85f76a13b6d 100644 --- a/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py +++ b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py @@ -417,7 +417,8 @@ async def _handle_backend_exception( # Ensure the event context is set for the exception handler. EventContext.set(ev_ctx) if otel.enabled: - # Keep events chained by the handler in the failed event's trace. + # Chain the handler's events under the failed event's parent context + # (its remote or enqueuing span); parent_txid links them to it. otel.attach_context(ev_ctx.otel_context) if events := self.backend_exception_handler(ex): await chain_updates( diff --git a/packages/reflex-base/src/reflex_base/otel.py b/packages/reflex-base/src/reflex_base/otel.py index bbb8d083249..2474c9518fd 100644 --- a/packages/reflex-base/src/reflex_base/otel.py +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -3,8 +3,8 @@ The framework calls into this module at a small number of fixed points (event dispatch, event context forks, state acquisition, socket messages). Every entry point checks the module-level ``enabled`` flag first, so with no -instrumentation installed the cost is one attribute read and no -``opentelemetry`` object is ever created. +instrumentation installed the cost is one attribute read: no span is started +and nothing is recorded (only no-op API objects exist, created at import). The ``reflex-otel`` package flips the flag via :func:`enable` once a tracer provider is available. diff --git a/packages/reflex-otel/CHANGELOG.md b/packages/reflex-otel/CHANGELOG.md index 825c32f0d03..8cf66aa33c9 100644 --- a/packages/reflex-otel/CHANGELOG.md +++ b/packages/reflex-otel/CHANGELOG.md @@ -1 +1 @@ -# Changelog + diff --git a/packages/reflex-otel/README.md b/packages/reflex-otel/README.md index 7dcdcc5bb1a..feb06948483 100644 --- a/packages/reflex-otel/README.md +++ b/packages/reflex-otel/README.md @@ -50,5 +50,6 @@ Plus the ASGI middleware's `http.server.*` metrics. `instrument()` accepts `tracer_provider`, `meter_provider`, `excluded_urls` (comma-separated URL patterns skipped by the ASGI middleware; defaults to -`OTEL_PYTHON_REFLEX_EXCLUDED_URLS`, else `/ping`) and the ASGI hooks +`OTEL_PYTHON_REFLEX_EXCLUDED_URLS`, else `OTEL_PYTHON_EXCLUDED_URLS`, else +`/ping`; pass `""` to exclude nothing) and the ASGI hooks `server_request_hook`, `client_request_hook`, `client_response_hook`. diff --git a/packages/reflex-otel/src/reflex_otel/__init__.py b/packages/reflex-otel/src/reflex_otel/__init__.py index 316ab480c29..8034b25b512 100644 --- a/packages/reflex-otel/src/reflex_otel/__init__.py +++ b/packages/reflex-otel/src/reflex_otel/__init__.py @@ -48,7 +48,9 @@ def _instrument(self, **kwargs: Any) -> None: **kwargs: ``tracer_provider`` and ``meter_provider`` select the providers (default: the global ones). ``excluded_urls`` is a comma-separated list of URL patterns the ASGI middleware skips - (default: ``OTEL_PYTHON_REFLEX_EXCLUDED_URLS``, else ``/ping``). + (default: ``OTEL_PYTHON_REFLEX_EXCLUDED_URLS``, else + ``OTEL_PYTHON_EXCLUDED_URLS``, else ``/ping``; pass ``""`` to + exclude nothing). ``server_request_hook``, ``client_request_hook`` and ``client_response_hook`` are forwarded to the ASGI middleware. """ @@ -58,11 +60,13 @@ def _instrument(self, **kwargs: Any) -> None: tracer_provider = kwargs.get("tracer_provider") meter_provider = kwargs.get("meter_provider") - excluded_urls = kwargs.get("excluded_urls") or ( - os.environ.get("OTEL_PYTHON_REFLEX_EXCLUDED_URLS") - or os.environ.get("OTEL_PYTHON_EXCLUDED_URLS") - or _DEFAULT_EXCLUDED_URLS - ) + excluded_urls = kwargs.get("excluded_urls") + if excluded_urls is None: + excluded_urls = ( + os.environ.get("OTEL_PYTHON_REFLEX_EXCLUDED_URLS") + or os.environ.get("OTEL_PYTHON_EXCLUDED_URLS") + or _DEFAULT_EXCLUDED_URLS + ) def asgi_middleware(app: otel.ASGIApp) -> otel.ASGIApp: return OpenTelemetryMiddleware( diff --git a/pyproject.toml b/pyproject.toml index 7abfac5bfbd..be5e1f0681c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -288,6 +288,7 @@ source = [ "reflex_components_sonner", "reflex_base", "reflex_docgen", + "reflex_otel", "reflex_release", "reflex_site_shared", "reflex_components_internal", diff --git a/tests/units/reflex_otel/test_init.py b/tests/units/reflex_otel/test_init.py index dfba88b0ac4..df9f8ee4ae2 100644 --- a/tests/units/reflex_otel/test_init.py +++ b/tests/units/reflex_otel/test_init.py @@ -47,3 +47,26 @@ async def app(scope, receive, send): ... assert wrapped.app is app instrumentor.uninstrument() assert otel.asgi_middleware is None + + +@pytest.mark.parametrize( + ("excluded_urls", "ping_disabled"), + [(None, True), ("", False), ("/health", False)], +) +def test_excluded_urls_only_defaults_when_omitted( + instrumentor: ReflexInstrumentor, + excluded_urls: str | None, + ping_disabled: bool, +): + kwargs = {} if excluded_urls is None else {"excluded_urls": excluded_urls} + instrumentor.instrument(tracer_provider=TracerProvider(), **kwargs) + assert otel.asgi_middleware is not None + + async def app(scope, receive, send): ... + + wrapped = otel.asgi_middleware(app) + assert isinstance(wrapped, OpenTelemetryMiddleware) + assert ( + bool(wrapped.excluded_urls and wrapped.excluded_urls.url_disabled("/ping")) + is ping_disabled + ) From b0d7618dc0eb558fee471b1a06ec687fefaa8504 Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 02:51:28 +0500 Subject: [PATCH 7/8] docs(otel): note the ASGI middleware lifetime for uninstrument() --- packages/reflex-otel/README.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/reflex-otel/README.md b/packages/reflex-otel/README.md index feb06948483..fbcea659435 100644 --- a/packages/reflex-otel/README.md +++ b/packages/reflex-otel/README.md @@ -53,3 +53,6 @@ Plus the ASGI middleware's `http.server.*` metrics. `OTEL_PYTHON_REFLEX_EXCLUDED_URLS`, else `OTEL_PYTHON_EXCLUDED_URLS`, else `/ping`; pass `""` to exclude nothing) and the ASGI hooks `server_request_hook`, `client_request_hook`, `client_response_hook`. +Call `instrument()` before the app is served: `uninstrument()` turns the +framework trace points off again, but an ASGI middleware that was already +installed stays until the process restarts. From 2f9b04d0690b5d8c3c3771626b12cf60fb2013f3 Mon Sep 17 00:00:00 2001 From: Farhan Date: Tue, 18 Aug 2026 20:56:27 +0500 Subject: [PATCH 8/8] fix(otel): parse excluded_urls for the ASGI floor, record event duration inside the span opentelemetry-instrumentation-asgi < 0.56b0 stores excluded_urls verbatim and calls .url_disabled() on it, so the default "/ping" string made every request raise AttributeError. Parse with parse_excluded_urls first. Move the event-duration histogram record inside the event span so exporter latency is excluded from the sample and exemplars keep the trace/span IDs. --- packages/reflex-base/src/reflex_base/otel.py | 22 +++++---- .../reflex-otel/src/reflex_otel/__init__.py | 5 ++ tests/units/reflex_base/test_otel.py | 20 ++++++++ tests/units/reflex_otel/test_init.py | 47 +++++++++++++++++-- 4 files changed, 80 insertions(+), 14 deletions(-) diff --git a/packages/reflex-base/src/reflex_base/otel.py b/packages/reflex-base/src/reflex_base/otel.py index 2474c9518fd..79ca5dd93f5 100644 --- a/packages/reflex-base/src/reflex_base/otel.py +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -256,17 +256,19 @@ def event_span( if parent.is_valid and not parent.is_remote else SpanKind.SERVER ) - start = perf_counter() - try: - with _tracer.start_as_current_span( - event.name, context=ctx.otel_context, kind=kind, attributes=attributes - ) as span: + with _tracer.start_as_current_span( + event.name, context=ctx.otel_context, kind=kind, attributes=attributes + ) as span: + # Record while the span is still current: exporter latency stays out + # of the sample and exemplars keep the trace/span IDs. + start = perf_counter() + try: yield span - except Exception as ex: - metric_attributes[ATTR_ERROR_TYPE] = type(ex).__qualname__ - raise - finally: - _event_duration.record(perf_counter() - start, metric_attributes) + except Exception as ex: + metric_attributes[ATTR_ERROR_TYPE] = type(ex).__qualname__ + raise + finally: + _event_duration.record(perf_counter() - start, metric_attributes) def record_state_acquired(start: float, event: Event) -> None: diff --git a/packages/reflex-otel/src/reflex_otel/__init__.py b/packages/reflex-otel/src/reflex_otel/__init__.py index 8034b25b512..2d1323faacc 100644 --- a/packages/reflex-otel/src/reflex_otel/__init__.py +++ b/packages/reflex-otel/src/reflex_otel/__init__.py @@ -57,6 +57,7 @@ def _instrument(self, **kwargs: Any) -> None: # Imported here so `from reflex_otel import OtelPlugin` in rxconfig.py # stays cheap for CLI processes that never instrument. from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware + from opentelemetry.util.http import parse_excluded_urls tracer_provider = kwargs.get("tracer_provider") meter_provider = kwargs.get("meter_provider") @@ -67,6 +68,10 @@ def _instrument(self, **kwargs: Any) -> None: or os.environ.get("OTEL_PYTHON_EXCLUDED_URLS") or _DEFAULT_EXCLUDED_URLS ) + # opentelemetry-instrumentation-asgi < 0.56b0 stores a str verbatim and + # then calls .url_disabled() on it, failing every request; parse first. + if isinstance(excluded_urls, str): + excluded_urls = parse_excluded_urls(excluded_urls) def asgi_middleware(app: otel.ASGIApp) -> otel.ASGIApp: return OpenTelemetryMiddleware( diff --git a/tests/units/reflex_base/test_otel.py b/tests/units/reflex_base/test_otel.py index fce75140d9a..9f2243f556d 100644 --- a/tests/units/reflex_base/test_otel.py +++ b/tests/units/reflex_base/test_otel.py @@ -133,6 +133,26 @@ def test_event_span_records_duration_metric(otel_metrics: InMemoryMetricReader): assert ok.sum >= 0 +def test_event_duration_recorded_inside_span( + otel_exporter: InMemorySpanExporter, monkeypatch: pytest.MonkeyPatch +): + seen: list[trace.Span] = [] + + class Histogram: + def record(self, amount, attributes=None): + seen.append(trace.get_current_span()) + + monkeypatch.setattr(otel, "_event_duration", Histogram()) + registered = RegisteredEventHandler(handler=EventHandler(fn=_handler), states=()) + with otel.event_span(Event(name="ok"), _ctx(), registered) as span: + pass + # The sample must be taken while the event span is still current so + # exporter latency is excluded and exemplars keep the trace/span IDs. + assert len(seen) == 1 + assert seen[0] is span + assert seen[0].get_span_context().is_valid + + def test_metric_helpers_record(otel_metrics: InMemoryMetricReader): otel.record_state_acquired(perf_counter(), Event(name="e")) otel.record_message_size(42, "transmit") diff --git a/tests/units/reflex_otel/test_init.py b/tests/units/reflex_otel/test_init.py index df9f8ee4ae2..0ffe0bf1d64 100644 --- a/tests/units/reflex_otel/test_init.py +++ b/tests/units/reflex_otel/test_init.py @@ -5,6 +5,7 @@ import pytest from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.util.http import ExcludeList from reflex_base import otel from reflex_otel import ReflexInstrumentor @@ -66,7 +67,45 @@ async def app(scope, receive, send): ... wrapped = otel.asgi_middleware(app) assert isinstance(wrapped, OpenTelemetryMiddleware) - assert ( - bool(wrapped.excluded_urls and wrapped.excluded_urls.url_disabled("/ping")) - is ping_disabled - ) + # Older ASGI instrumentation stores strings verbatim and then calls + # .url_disabled() on them; always hand it a parsed ExcludeList. + assert isinstance(wrapped.excluded_urls, ExcludeList) + assert wrapped.excluded_urls.url_disabled("/ping") is ping_disabled + + +async def test_asgi_middleware_handles_requests(instrumentor: ReflexInstrumentor): + """The wrapped app must serve requests, including at the ASGI floor. + + Older ASGI instrumentation stores excluded_urls verbatim, so a raw string + made every request raise AttributeError. + """ + instrumentor.instrument(tracer_provider=TracerProvider()) + assert otel.asgi_middleware is not None + served: list[str] = [] + + async def app(scope, receive, send): + served.append(scope["path"]) + await send({"type": "http.response.start", "status": 200, "headers": []}) + await send({"type": "http.response.body", "body": b""}) + + async def receive(): # noqa: RUF029 + return {"type": "http.request", "body": b"", "more_body": False} + + async def send(message): + pass + + wrapped = otel.asgi_middleware(app) + for path in ("/ping", "/_event"): + scope = { + "type": "http", + "method": "GET", + "path": path, + "raw_path": path.encode(), + "query_string": b"", + "headers": [], + "scheme": "http", + "server": ("localhost", 80), + "http_version": "1.1", + } + await wrapped(scope, receive, send) + assert served == ["/ping", "/_event"]