diff --git a/backend/open_webui/utils/telemetry/metrics.py b/backend/open_webui/utils/telemetry/metrics.py index 26216b6ca..a1d1dcb7c 100644 --- a/backend/open_webui/utils/telemetry/metrics.py +++ b/backend/open_webui/utils/telemetry/metrics.py @@ -17,8 +17,10 @@ high-cardinality label sets. from __future__ import annotations +import datetime +import logging import time -from typing import Dict, List, Sequence, Any +from typing import Dict, Iterable, List, Optional from base64 import b64encode from fastapi import FastAPI, Request @@ -36,6 +38,8 @@ from opentelemetry.sdk.metrics.export import ( PeriodicExportingMetricReader, ) from opentelemetry.sdk.resources import Resource +from sqlalchemy import Engine, func, select +from sqlalchemy.orm import Session from open_webui.env import ( OTEL_SERVICE_NAME, @@ -46,7 +50,47 @@ from open_webui.env import ( OTEL_METRICS_EXPORTER_OTLP_INSECURE, OTEL_METRICS_EXPORT_INTERVAL_MILLIS, ) -from open_webui.models.users import Users +from open_webui.models.users import User + +logger = logging.getLogger(__name__) + + +# --------------------------------------------------------------------------- +# Sync DB helpers for OTel gauge callbacks +# +# The OTel Python SDK calls observable-instrument callbacks *synchronously* +# from a background collection thread — async callbacks are NOT supported +# (the SDK does not ``await`` the return value). +# +# Rather than bridging into the async event loop, we run plain synchronous +# SQL queries using the sync engine that is already available at setup time. +# This avoids any cross-thread / cross-loop concerns entirely. +# --------------------------------------------------------------------------- + + +def _count_total_users(db_engine: Engine) -> Optional[int]: + """Return the total number of registered users (sync).""" + with Session(db_engine) as session: + return session.execute(select(func.count()).select_from(User)).scalar() + + +def _count_active_users(db_engine: Engine) -> Optional[int]: + """Return the number of users active within the last 3 minutes (sync).""" + three_minutes_ago = int(time.time()) - 180 + with Session(db_engine) as session: + return session.execute( + select(func.count()).select_from(User).filter(User.last_active_at >= three_minutes_ago) + ).scalar() + + +def _count_users_active_today(db_engine: Engine) -> Optional[int]: + """Return the number of users active since midnight today (sync).""" + now = int(datetime.datetime.now().timestamp()) + today_midnight = now - (now % 86400) + with Session(db_engine) as session: + return session.execute( + select(func.count()).select_from(User).filter(User.last_active_at > today_midnight) + ).scalar() def _build_meter_provider(resource: Resource) -> MeterProvider: @@ -106,7 +150,7 @@ def _build_meter_provider(resource: Resource) -> MeterProvider: return provider -def setup_metrics(app: FastAPI, resource: Resource) -> None: +def setup_metrics(app: FastAPI, resource: Resource, db_engine: Engine) -> None: """Attach OTel metrics middleware to *app* and initialise provider.""" metrics.set_meter_provider(_build_meter_provider(resource)) @@ -124,32 +168,46 @@ def setup_metrics(app: FastAPI, resource: Resource) -> None: unit='ms', ) - async def observe_active_users( - options: metrics.CallbackOptions, - ) -> Sequence[metrics.Observation]: - return [ - metrics.Observation( - value=await Users.get_active_user_count(), - ) - ] + # -- Observable gauge callbacks ---------------------------------------- + # These are called synchronously by the OTel SDK from a background + # collection thread. They use the sync DB engine directly — no async + # bridging required. - async def observe_total_registered_users( + def observe_total_users( options: metrics.CallbackOptions, - ) -> Sequence[metrics.Observation]: - # IMPORTANT: Use get_num_users() for efficient COUNT(*) query. - # Do NOT use len(get_users()["users"]) - it loads ALL user records into memory, - # causing connection pool exhaustion on high-latency databases (e.g., Aurora). - return [ - metrics.Observation( - value=await Users.get_num_users() or 0, - ) - ] + ) -> Iterable[metrics.Observation]: + try: + value = _count_total_users(db_engine) + if value is not None: + yield metrics.Observation(value=value) + except Exception: + logger.debug('Failed to observe total users', exc_info=True) + + def observe_active_users( + options: metrics.CallbackOptions, + ) -> Iterable[metrics.Observation]: + try: + value = _count_active_users(db_engine) + if value is not None: + yield metrics.Observation(value=value) + except Exception: + logger.debug('Failed to observe active users', exc_info=True) + + def observe_users_active_today( + options: metrics.CallbackOptions, + ) -> Iterable[metrics.Observation]: + try: + value = _count_users_active_today(db_engine) + if value is not None: + yield metrics.Observation(value=value) + except Exception: + logger.debug('Failed to observe users active today', exc_info=True) meter.create_observable_gauge( name='webui.users.total', description='Total number of registered users', unit='users', - callbacks=[observe_total_registered_users], + callbacks=[observe_total_users], ) meter.create_observable_gauge( @@ -159,11 +217,6 @@ def setup_metrics(app: FastAPI, resource: Resource) -> None: callbacks=[observe_active_users], ) - async def observe_users_active_today( - options: metrics.CallbackOptions, - ) -> Sequence[metrics.Observation]: - return [metrics.Observation(value=await Users.get_num_users_active_today())] - meter.create_observable_gauge( name='webui.users.active.today', description='Number of users active since midnight today', diff --git a/backend/open_webui/utils/telemetry/setup.py b/backend/open_webui/utils/telemetry/setup.py index 744dced2d..14f10ef97 100644 --- a/backend/open_webui/utils/telemetry/setup.py +++ b/backend/open_webui/utils/telemetry/setup.py @@ -55,4 +55,4 @@ def setup(app: FastAPI, db_engine: Engine): # set up metrics only if enabled if ENABLE_OTEL_METRICS: - setup_metrics(app, resource) + setup_metrics(app, resource, db_engine)