Files
Northern-Thailand-Ping-Rive…/src/metrics.py
T
grabowski 9cac9c4d2a style: apply black/isort across the repo; make CI mypy advisory
The push-CI gates (black/isort/mypy) had never actually run before the
branch-trigger fix, and the codebase predates them. Formatting is now
black/isort clean repo-wide. mypy keeps running but non-blocking: 86
pre-existing errors are a separate cleanup, not a gate to hold hostage.
2026-08-10 15:57:00 +07:00

199 lines
6.3 KiB
Python

#!/usr/bin/env python3
"""
Metrics collection and monitoring for water monitoring system
"""
import logging
import threading
import time
from collections import defaultdict, deque
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from typing import Any, Dict, List, Optional
logger = logging.getLogger(__name__)
@dataclass
class MetricPoint:
"""Single metric data point"""
timestamp: datetime
value: float
labels: Dict[str, str] = field(default_factory=dict)
class MetricsCollector:
"""Collects and manages application metrics"""
def __init__(self, retention_hours: int = 24):
self.retention_hours = retention_hours
self.metrics: Dict[str, deque] = defaultdict(lambda: deque(maxlen=1000))
self.counters: Dict[str, float] = defaultdict(float)
self.gauges: Dict[str, float] = defaultdict(float)
self.histograms: Dict[str, List[float]] = defaultdict(list)
self._lock = threading.Lock()
# Start cleanup thread
self._cleanup_thread = threading.Thread(
target=self._cleanup_old_metrics, daemon=True
)
self._cleanup_thread.start()
def increment_counter(
self, name: str, value: float = 1.0, labels: Optional[Dict[str, str]] = None
):
"""Increment a counter metric"""
with self._lock:
key = self._make_key(name, labels)
self.counters[key] += value
self.metrics[key].append(
MetricPoint(datetime.now(), self.counters[key], labels or {})
)
def set_gauge(
self, name: str, value: float, labels: Optional[Dict[str, str]] = None
):
"""Set a gauge metric"""
with self._lock:
key = self._make_key(name, labels)
self.gauges[key] = value
self.metrics[key].append(MetricPoint(datetime.now(), value, labels or {}))
def record_histogram(
self, name: str, value: float, labels: Optional[Dict[str, str]] = None
):
"""Record a histogram value"""
with self._lock:
key = self._make_key(name, labels)
self.histograms[key].append(value)
# Keep only recent values
if len(self.histograms[key]) > 1000:
self.histograms[key] = self.histograms[key][-1000:]
self.metrics[key].append(MetricPoint(datetime.now(), value, labels or {}))
def get_counter(self, name: str, labels: Optional[Dict[str, str]] = None) -> float:
"""Get current counter value"""
key = self._make_key(name, labels)
return self.counters.get(key, 0.0)
def get_gauge(self, name: str, labels: Optional[Dict[str, str]] = None) -> float:
"""Get current gauge value"""
key = self._make_key(name, labels)
return self.gauges.get(key, 0.0)
def get_histogram_stats(
self, name: str, labels: Optional[Dict[str, str]] = None
) -> Dict[str, float]:
"""Get histogram statistics"""
key = self._make_key(name, labels)
values = self.histograms.get(key, [])
if not values:
return {"count": 0, "sum": 0, "avg": 0, "min": 0, "max": 0}
return {
"count": len(values),
"sum": sum(values),
"avg": sum(values) / len(values),
"min": min(values),
"max": max(values),
}
def get_all_metrics(self) -> Dict[str, Any]:
"""Get all current metrics"""
with self._lock:
return {
"counters": dict(self.counters),
"gauges": dict(self.gauges),
"histograms": {k: self.get_histogram_stats(k) for k in self.histograms},
}
def _make_key(self, name: str, labels: Optional[Dict[str, str]]) -> str:
"""Create a unique key for metric with labels"""
if not labels:
return name
label_str = ",".join(f"{k}={v}" for k, v in sorted(labels.items()))
return f"{name}{{{label_str}}}"
def _cleanup_old_metrics(self):
"""Clean up old metric data points"""
while True:
try:
cutoff_time = datetime.now() - timedelta(hours=self.retention_hours)
with self._lock:
for metric_name, points in self.metrics.items():
# Remove old points
while points and points[0].timestamp < cutoff_time:
points.popleft()
time.sleep(3600) # Run cleanup every hour
except Exception as e:
logger.error(f"Error in metrics cleanup: {e}")
time.sleep(60) # Wait a minute before retrying
# Global metrics collector instance
_metrics_collector = None
def get_metrics_collector() -> MetricsCollector:
"""Get the global metrics collector instance"""
global _metrics_collector
if _metrics_collector is None:
_metrics_collector = MetricsCollector()
return _metrics_collector
# Convenience functions
def increment_counter(
name: str, value: float = 1.0, labels: Optional[Dict[str, str]] = None
):
"""Increment a counter metric"""
get_metrics_collector().increment_counter(name, value, labels)
def set_gauge(name: str, value: float, labels: Optional[Dict[str, str]] = None):
"""Set a gauge metric"""
get_metrics_collector().set_gauge(name, value, labels)
def record_histogram(name: str, value: float, labels: Optional[Dict[str, str]] = None):
"""Record a histogram value"""
get_metrics_collector().record_histogram(name, value, labels)
class Timer:
"""Context manager for timing operations"""
def __init__(self, metric_name: str, labels: Optional[Dict[str, str]] = None):
self.metric_name = metric_name
self.labels = labels
self.start_time = None
def __enter__(self):
self.start_time = time.time()
return self
def __exit__(self, exc_type, exc_val, exc_tb):
if self.start_time:
duration = time.time() - self.start_time
record_histogram(self.metric_name, duration, self.labels)
def timer(metric_name: str, labels: Optional[Dict[str, str]] = None):
"""Decorator for timing function execution"""
def decorator(func):
def wrapper(*args, **kwargs):
with Timer(metric_name, labels):
return func(*args, **kwargs)
return wrapper
return decorator