crowdata / app /cache /redis_client.py
YOSOYYONOSOYOTRO's picture
Upload folder using huggingface_hub
89e2f2e verified
Raw
History Blame Contribute Delete
5.88 kB
"""
Redis Client Pool — CrowData
Centralized Redis connection management with in-memory LRU fallback.
"""
import asyncio
import json
import logging
import time
from collections import OrderedDict
from typing import Optional
try:
import redis.asyncio as aioredis
_REDIS_MODULE_AVAILABLE = True
except ImportError:
aioredis = None
_REDIS_MODULE_AVAILABLE = False
from app.config import get_settings
logger = logging.getLogger(__name__)
settings = get_settings()
# Global connection pools
_redis_client = None
_redis_is_available = True
_reconnect_task = None
# In-memory LRU cache fallback
class LRUCache:
def __init__(self, maxsize: int = 5000):
self.maxsize = maxsize
self._cache = OrderedDict()
self._ttls = {}
self._lock = asyncio.Lock()
async def get(self, key: str):
async with self._lock:
if key not in self._cache:
return None
expire_at = self._ttls.get(key, 0)
if expire_at and expire_at < time.time():
# Expired
self._cache.pop(key, None)
self._ttls.pop(key, None)
return None
# Move to end (most recently used)
value = self._cache.pop(key)
self._cache[key] = value
return value
async def set(self, key: str, value: dict, ttl: int):
async with self._lock:
# Evict if at maxsize
if len(self._cache) >= self.maxsize and key not in self._cache:
self._cache.popitem(last=False) # Remove LRU
self._cache[key] = value
self._ttls[key] = time.time() + ttl if ttl > 0 else 0
async def delete(self, key: str):
async with self._lock:
self._cache.pop(key, None)
self._ttls.pop(key, None)
async def cleanup_expired(self):
"""Remove expired entries. Called periodically."""
async with self._lock:
now = time.time()
expired = [k for k, exp in self._ttls.items() if exp and exp < now]
for k in expired:
self._cache.pop(k, None)
self._ttls.pop(k, None)
def __len__(self):
return len(self._cache)
_in_memory_cache = LRUCache(maxsize=5000)
_redis_client = None
_redis_is_available = _REDIS_MODULE_AVAILABLE # False si el módulo no está instalado
_reconnect_task = None
async def get_redis():
global _redis_client
if not _REDIS_MODULE_AVAILABLE:
raise RuntimeError("redis module not installed; using in-memory fallback")
if _redis_client is None:
from app.config import get_settings
settings = get_settings()
_redis_client = await aioredis.from_url(
settings.redis_url,
encoding="utf-8",
decode_responses=True,
max_connections=getattr(settings, 'redis_max_connections', 20) or 20,
socket_keepalive=True,
socket_connect_timeout=5,
socket_timeout=5,
retry_on_timeout=True,
)
return _redis_client
async def _try_reconnect():
"""Background task to attempt Redis reconnection."""
global _redis_is_available, _reconnect_task
while True:
await asyncio.sleep(30) # Try every 30 seconds
if not _redis_is_available:
try:
r = await get_redis()
await r.ping()
_redis_is_available = True
logger.info("Redis reconnected successfully")
except Exception:
pass
async def _start_reconnect_task():
global _reconnect_task
if _reconnect_task is None or _reconnect_task.done():
_reconnect_task = asyncio.create_task(_try_reconnect())
async def cache_get(key: str):
global _redis_is_available
if _redis_is_available:
try:
r = await get_redis()
data = await r.get(key)
if data:
return json.loads(data)
except Exception as e:
logger.warning(f"Redis GET error for key {key}: {e}. Falling back to in-memory cache.")
_redis_is_available = False
# Start background reconnection
asyncio.create_task(_try_reconnect())
# In-memory fallback
return await _in_memory_cache.get(key)
async def cache_set(key: str, value: dict, ttl: int = None):
global _redis_is_available
ttl = ttl or 86400
if _redis_is_available:
try:
r = await get_redis()
await r.setex(key, ttl, json.dumps(value, ensure_ascii=False))
return
except Exception as e:
logger.warning(f"Redis SET error for key {key}: {e}. Saving in memory.")
_redis_is_available = False
# Start background reconnection
asyncio.create_task(_try_reconnect())
# In-memory LRU fallback
await _in_memory_cache.set(key, value, ttl=60)
async def cache_delete(key: str):
global _redis_is_available
if _redis_is_available:
try:
r = await get_redis()
await r.delete(key)
except Exception as e:
logger.warning(f"Redis DELETE error for key {key}: {e}.")
_redis_is_available = False
asyncio.create_task(_try_reconnect())
# In-memory fallback
await _in_memory_cache.delete(key)
async def cleanup_expired():
"""Periodic cleanup of expired in-memory cache entries."""
await _in_memory_cache.cleanup_expired()
async def _start_reconnect_task():
"""Start the background reconnection task."""
pass # Task is already started at module level
# Module-level initialization - start background tasks
try:
loop = asyncio.get_running_loop()
loop.create_task(_try_reconnect())
except RuntimeError:
# No event loop running, will be created later
pass