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/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 new file mode 100644 index 00000000000..0ed2e91a1e3 --- /dev/null +++ b/packages/reflex-base/news/6227.feature.md @@ -0,0 +1 @@ +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/pyproject.toml b/packages/reflex-base/pyproject.toml index 7e744790a18..310626adce4 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.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/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/base_state_processor.py b/packages/reflex-base/src/reflex_base/event/processor/base_state_processor.py index 9411d62927f..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 @@ -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 otel.enabled: + 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() @@ -411,6 +416,10 @@ 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: + # 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( events=events, 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..79ca5dd93f5 --- /dev/null +++ b/packages/reflex-base/src/reflex_base/otel.py @@ -0,0 +1,302 @@ +"""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, 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: 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. +""" + +from __future__ import annotations + +from collections.abc import Awaitable, Callable, Iterator, Mapping +from contextlib import contextmanager +from time import perf_counter +from typing import TYPE_CHECKING, Any + +from opentelemetry import context as otel_context +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 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 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" + +# 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) + +# 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[[ASGIApp], ASGIApp] | 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.", + 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, + unit="{connection}", + description="Number of open client socket connections.", + ) + + +def enable( + tracer_provider: trace.TracerProvider | None = None, + meter_provider: metrics.MeterProvider | None = None, + asgi_middleware_factory: Callable[[ASGIApp], ASGIApp] | 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, 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 and instruments.""" + global _tracer, enabled, asgi_middleware + enabled = False + asgi_middleware = None + _tracer = trace.NoOpTracer() + _create_instruments(_noop_meter) + + +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 + + +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.""" + + __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]) -> _AttachedContext: + """Make the trace context carried by a frontend event the current context. + + 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. + """ + 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 and record its duration. + + 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. + ctx: The event context for this execution. + registered_handler: The handler resolved for the event. + + Yields: + The active span. + """ + handler = registered_handler.handler + metric_attributes = { + 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 + ) + 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) + + +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/CHANGELOG.md b/packages/reflex-otel/CHANGELOG.md new file mode 100644 index 00000000000..8cf66aa33c9 --- /dev/null +++ b/packages/reflex-otel/CHANGELOG.md @@ -0,0 +1 @@ + diff --git a/packages/reflex-otel/README.md b/packages/reflex-otel/README.md new file mode 100644 index 00000000000..fbcea659435 --- /dev/null +++ b/packages/reflex-otel/README.md @@ -0,0 +1,58 @@ +# 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 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 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). + +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`); default `sio` only | +| `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 +`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. diff --git a/packages/reflex-otel/news/6227.feature.md b/packages/reflex-otel/news/6227.feature.md new file mode 100644 index 00000000000..fdc6ea8d17b --- /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 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 new file mode 100644 index 00000000000..df219b2bddc --- /dev/null +++ b/packages/reflex-otel/pyproject.toml @@ -0,0 +1,32 @@ +[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.30.0,<2.0", + "opentelemetry-instrumentation >=0.49b0,<1.0", + "opentelemetry-instrumentation-asgi >=0.49b0,<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..2d1323faacc --- /dev/null +++ b/packages/reflex-otel/src/reflex_otel/__init__.py @@ -0,0 +1,100 @@ +"""OpenTelemetry instrumentation for the Reflex framework.""" + +from __future__ import annotations + +import os +from collections.abc import Collection +from typing import Any, Literal + +from opentelemetry.instrumentation.instrumentor import BaseInstrumentor +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): + """Enable the trace points and metrics built into the Reflex runtime. + + Usage:: + + ReflexInstrumentor().instrument(tracer_provider=provider) + + 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]: + """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`` 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 + ``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. + """ + # 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") + 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 + ) + # 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( + app, + 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"), + 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. + + Args: + **kwargs: Ignored. + """ + otel.disable() diff --git a/pyproject.toml b/pyproject.toml index 4a3759b3775..be5e1f0681c 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", ] @@ -285,6 +288,7 @@ source = [ "reflex_components_sonner", "reflex_base", "reflex_docgen", + "reflex_otel", "reflex_release", "reflex_site_shared", "reflex_components_internal", @@ -427,6 +431,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/reflex/app.py b/reflex/app.py index 258410fe9f7..aa21898bde6 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,39 @@ 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.encode()), "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.encode()) if isinstance(data, str) else len(data), "receive" + ) + return json.loads(data, **kwargs) + + class EventNamespace(AsyncNamespace): """The event namespace.""" @@ -1997,6 +2032,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 +2044,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 +2180,11 @@ 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) + 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) 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..60012f2fed4 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,74 @@ def preserve_memo_registries(): finally: MEMOS.clear() MEMOS.update(memos) + + +@pytest.fixture +def otel_sdk() -> Generator[ + tuple[InMemorySpanExporter, InMemoryMetricReader], None, None +]: + """Enable the reflex_base.otel trace points and metrics against in-memory sinks. + + Yields: + The span exporter and the metric reader. + """ + exporter = InMemorySpanExporter() + 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, reader + finally: + otel.disable() + + +@pytest.fixture +def otel_exporter(otel_sdk) -> InMemorySpanExporter: + """The in-memory span exporter of the enabled otel_sdk. + + Args: + otel_sdk: The enabled sinks. + + Returns: + The span exporter. + """ + 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 6f89d9adf49..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 @@ -214,3 +216,34 @@ 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. + """ + + 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) + + 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 d5dda19dca3..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, @@ -792,3 +795,48 @@ 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, 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. + """ + 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): + """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. + """ + 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 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] + ) + 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..9f2243f556d --- /dev/null +++ b/tests/units/reflex_base/test_otel.py @@ -0,0 +1,217 @@ +"""Tests for the reflex_base.otel trace points.""" + +import asyncio +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 +from reflex_base.event.context import EventContext +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: + 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 + assert child.kind == SpanKind.INTERNAL + + +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 + 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, + 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_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") + 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_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.kind == SpanKind.SERVER + 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 + + +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/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..0ffe0bf1d64 --- /dev/null +++ b/tests/units/reflex_otel/test_init.py @@ -0,0 +1,111 @@ +"""Tests for the reflex_otel instrumentor.""" + +from collections.abc import Generator + +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 + + +@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") + + +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 + + +@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) + # 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"] diff --git a/tests/units/test_app.py b/tests/units/test_app.py index 60e3122d70a..773b463044f 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, @@ -4257,3 +4266,76 @@ 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.""" + 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.""" + 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_points(otel_metrics, otel.METRIC_WEBSOCKET_MESSAGE_SIZE) + } + 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.""" + 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.""" + 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 + (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() diff --git a/uv.lock b/uv.lock index 75a88ad310d..1f5421157fd 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", ] @@ -287,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" @@ -585,7 +598,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 +676,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 +1016,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 +1892,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 +1976,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 +2437,85 @@ 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-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" +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 = "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" @@ -2534,10 +2626,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 +2698,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 +3782,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 +3805,7 @@ dev = [ { name = "python-dotenv" }, { name = "pyyaml" }, { name = "reflex-docgen" }, + { name = "reflex-otel" }, { name = "reflex-release" }, { name = "reflex-site-shared" }, { name = "ruff" }, @@ -3772,6 +3866,7 @@ dev = [ { name = "hatchling" }, { name = "libsass" }, { name = "numpy" }, + { name = "opentelemetry-sdk" }, { name = "pandas" }, { name = "pillow" }, { name = "playwright" }, @@ -3793,6 +3888,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 +3907,7 @@ dev = [ name = "reflex-base" source = { editable = "packages/reflex-base" } dependencies = [ + { name = "opentelemetry-api" }, { name = "packaging" }, { name = "platformdirs" }, { name = "rich" }, @@ -3824,6 +3921,7 @@ pydantic = [ [package.metadata] requires-dist = [ + { 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" }, @@ -4149,6 +4247,31 @@ 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 = "opentelemetry-instrumentation-asgi" }, + { name = "reflex-base" }, +] + +[package.optional-dependencies] +instruments = [ + { name = "reflex-base" }, +] + +[package.metadata] +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 = "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 +4427,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 +4488,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 +4567,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 = [