214 lines
7.5 KiB
Python
214 lines
7.5 KiB
Python
"""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
|