"""System and observability API service functions.""" from __future__ import annotations import time import os from collections import Counter from typing import Any, Dict, Optional from fastapi import BackgroundTasks, Request from fastapi.concurrency import run_in_threadpool from fastapi.responses import PlainTextResponse from loguru import logger from src.database.db_manager import DBManager from src.utils.metrics import export_prometheus_metrics, gauge_set from web.core import build_health_payload, build_system_status_payload import web.routes as legacy_routes def get_health_payload() -> Dict[str, Any]: return build_health_payload() async def get_system_status_payload(request: Request) -> Dict[str, Any]: legacy_routes._require_ops_admin(request) payload = await run_in_threadpool(build_system_status_payload) payload["realtime"] = await run_in_threadpool(_realtime_status_payload) return payload def _realtime_status_payload() -> Dict[str, Any]: try: from web.routers import sse_router store = sse_router.event_store status_fn = getattr(store, "status", None) if callable(status_fn): status = dict(status_fn()) else: store_name = "degraded_sqlite" if getattr(store, "degraded_from", None) == "redis" else "sqlite" status = { "store": store_name, "latest_revision": int(store.latest_revision()), } connection_count = getattr(sse_router.sse_manager, "connection_count", None) status["sse_connections"] = int(connection_count()) if callable(connection_count) else 0 return status except Exception as exc: return { "store": "unknown", "latest_revision": 0, "sse_connections": 0, "error": str(exc), } def get_system_cache_status(request: Request, cities: Optional[str] = None) -> Dict[str, Any]: legacy_routes._require_ops_admin(request) selected = legacy_routes._normalize_city_list(cities) if not selected: selected = legacy_routes._normalize_city_list(None) kinds = { "summary": legacy_routes.CITY_SUMMARY_CACHE_TTL_SEC, "panel": legacy_routes.CITY_PANEL_CACHE_TTL_SEC, "nearby": legacy_routes.CITY_NEARBY_CACHE_TTL_SEC, "market": legacy_routes.CITY_MARKET_CACHE_TTL_SEC, "full": legacy_routes.CITY_FULL_CACHE_TTL_SEC, } items = [] for city in selected: row = {"city": city} for kind, ttl_sec in kinds.items(): entry = legacy_routes._CACHE_DB.get_city_cache(kind, city) row[kind] = { "exists": bool(entry), "fresh": legacy_routes._city_cache_is_fresh(entry, ttl_sec), "updated_at": entry.get("updated_at") if entry else None, "age_sec": round(max(0.0, time.time() - float(entry.get("updated_at_ts") or 0.0)), 1) if entry else None, "ttl_sec": ttl_sec, } items.append(row) return {"cities": items} def run_system_priority_warm( request: Request, background_tasks: BackgroundTasks, timezone: Optional[str] = None, ) -> Dict[str, Any]: legacy_routes._assert_entitlement(request) batches = legacy_routes._select_priority_city_batches(timezone) primary = list(batches.get("primary") or []) secondary = list(batches.get("secondary") or []) def _queue_city_refresh(city: str, *, priority: str) -> None: enqueue = getattr(legacy_routes._CACHE_DB, "enqueue_observation_refresh_request", None) if not callable(enqueue): logger.warning("priority warm queue unavailable city={} timezone={}", city, timezone) return enqueue( city=city, kind="panel", priority=priority, reason="system_priority_warm", ) def _runner() -> None: for city in primary: try: _queue_city_refresh(city, priority="high") except Exception as exc: logger.warning("priority warm primary failed city={} timezone={}: {}", city, timezone, exc) for city in secondary: try: _queue_city_refresh(city, priority="normal") except Exception as exc: logger.warning("priority warm secondary failed city={} timezone={}: {}", city, timezone, exc) background_tasks.add_task(_runner) return { "ok": True, "region": batches.get("region"), "timezone": batches.get("timezone"), "primary": primary, "secondary": secondary, } def get_prometheus_metrics_response(request: Request) -> PlainTextResponse: legacy_routes._require_ops_admin(request) _refresh_operational_metrics() return PlainTextResponse( export_prometheus_metrics(), media_type="text/plain; version=0.0.4; charset=utf-8", ) def _payment_event_reason(row: Dict[str, Any]) -> str: payload = row.get("payload") if isinstance(row, dict) else {} payload = payload if isinstance(payload, dict) else {} confirm_failure = ( payload.get("confirm_failure") if isinstance(payload.get("confirm_failure"), dict) else {} ) return str( payload.get("reason") or confirm_failure.get("reason") or payload.get("error") or "unknown" ).strip().lower() or "unknown" def _payment_event_is_resolved(row: Dict[str, Any]) -> bool: payload = row.get("payload") if isinstance(row, dict) else {} payload = payload if isinstance(payload, dict) else {} return bool(str(payload.get("resolved_at") or "").strip()) def _refresh_operational_metrics() -> None: try: db = DBManager() except Exception: return payment_events = [] for event_type in ("payment_intent_failed", "payment_refund_required"): try: payment_events.extend( db.list_payment_audit_events(limit=500, event_type=event_type) ) except Exception: continue open_events = [row for row in payment_events if not _payment_event_is_resolved(row)] gauge_set("polyweather_payment_incidents_open", len(open_events)) for reason, count in Counter(_payment_event_reason(row) for row in open_events).items(): gauge_set("polyweather_payment_incidents_by_reason", count, reason=reason) try: refund_cases = db.list_refund_cases(limit=500) except Exception: refund_cases = [] terminal_statuses = {"refunded", "rejected", "closed"} open_refunds = [ case for case in refund_cases if str(case.get("status") or "").strip().lower() not in terminal_statuses ] gauge_set("polyweather_refund_cases_open", len(open_refunds)) realtime = _realtime_status_payload() gauge_set("polyweather_sse_connections", int(realtime.get("sse_connections") or 0)) gauge_set( "polyweather_realtime_latest_revision", int(realtime.get("latest_revision") or 0), ) degraded_from = str(realtime.get("degraded_from") or "").strip().lower() store = str(realtime.get("store") or "").strip().lower() gauge_set( "polyweather_realtime_redis_fallback", 1 if degraded_from == "redis" or store == "degraded_sqlite" else 0, ) db_path = str(getattr(db, "db_path", "") or "").strip() if db_path: try: gauge_set("polyweather_sqlite_db_size_bytes", os.path.getsize(db_path)) except OSError: pass