import abc import time import uuid from collections import deque from dataclasses import dataclass, field from redis import Redis _RPS_WINDOW_S = 30 @dataclass class RequestStats: requests: int = field(default=0) retries: int = field(default=0) rate_limited: int = field(default=0) rps: float | None = field(default=None) @property def total_attempts(self) -> int: return self.requests + self.retries class StatsBackend(abc.ABC): @abc.abstractmethod def incr_requests(self) -> None: ... @abc.abstractmethod def incr_retries(self) -> None: ... @abc.abstractmethod def incr_rate_limited(self) -> None: ... @abc.abstractmethod def snapshot(self) -> RequestStats: ... class InMemoryStatsBackend(StatsBackend): def __init__(self) -> None: self._requests = 0 self._retries = 0 self._rate_limited = 0 self._window: deque[float] = deque() def incr_requests(self) -> None: self._requests += 1 self._window.append(time.monotonic()) def incr_retries(self) -> None: self._retries += 1 def incr_rate_limited(self) -> None: self._rate_limited += 1 def snapshot(self) -> RequestStats: now = time.monotonic() cutoff = now - _RPS_WINDOW_S while self._window and self._window[0] < cutoff: self._window.popleft() rps: float | None = None if len(self._window) >= 2: elapsed = now - self._window[0] if elapsed > 0: rps = round(len(self._window) / elapsed, 1) return RequestStats( requests=self._requests, retries=self._retries, rate_limited=self._rate_limited, rps=rps, ) class RedisStatsBackend(StatsBackend): def __init__(self, client: Redis[bytes], key_prefix: str) -> None: self._client = client self._keys = { "requests": f"{key_prefix}:requests", "retries": f"{key_prefix}:retries", "rate_limited": f"{key_prefix}:rate_limited", "rps_window": f"{key_prefix}:rps_window", } def incr_requests(self) -> None: now = time.time() pipe = self._client.pipeline() pipe.incr(self._keys["requests"]) pipe.zadd(self._keys["rps_window"], {uuid.uuid4().hex: now}) pipe.zremrangebyscore(self._keys["rps_window"], 0, now - _RPS_WINDOW_S) pipe.execute() def incr_retries(self) -> None: self._client.incr(self._keys["retries"]) def incr_rate_limited(self) -> None: self._client.incr(self._keys["rate_limited"]) def _get_int(self, key: str) -> int: value = self._client.get(key) return int(value) if value is not None else 0 def snapshot(self) -> RequestStats: now = time.time() pipe = self._client.pipeline() pipe.get(self._keys["requests"]) pipe.get(self._keys["retries"]) pipe.get(self._keys["rate_limited"]) pipe.zrangebyscore( self._keys["rps_window"], now - _RPS_WINDOW_S, now, withscores=True ) results = pipe.execute() timestamps = [score for _, score in (results[3] or [])] rps: float | None = None if len(timestamps) >= 2: elapsed = now - min(timestamps) if elapsed > 0: rps = round(len(timestamps) / elapsed, 1) return RequestStats( requests=self._get_int(self._keys["requests"]), retries=self._get_int(self._keys["retries"]), rate_limited=self._get_int(self._keys["rate_limited"]), rps=rps, )