feat: add volatile cache system for ephemeral data
Phase 2 of Memory Management System - volatile memory tier:
- VolatileCacheService: Redis-backed TTL storage
- Volatile router with full CRUD operations
- Predefined namespaces: weather, news, financial, transit, traffic,
air_quality, sports, social, system, context, custom
- Each namespace has appropriate default TTL (1min to 1hr)
- Refresh schedule support via cron expressions
- Scheduler integration endpoint: GET /volatile/scheduled
Endpoints:
- GET/POST/DELETE /volatile/{namespace}/{key}
- GET/DELETE /volatile/{namespace}
- GET /volatile/stats
- GET /volatile/scheduled
- GET /volatile/namespaces
🤖 Generated with [Claude Code](https://claude.com/claude-code)
Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,433 @@
|
||||
"""
|
||||
Volatile Cache service for Library Desk.
|
||||
|
||||
Provides ephemeral data storage with TTL for time-sensitive information:
|
||||
- Weather, news, financial data
|
||||
- Transit schedules, traffic conditions
|
||||
- System status, social notifications
|
||||
"""
|
||||
|
||||
import json
|
||||
import logging
|
||||
import hashlib
|
||||
from datetime import datetime
|
||||
from typing import List, Optional, Dict, Any
|
||||
|
||||
import redis.asyncio as aioredis
|
||||
|
||||
from src.config import Settings
|
||||
from src.models.volatile import (
|
||||
VolatileRecord,
|
||||
VolatileRecordResponse,
|
||||
VolatileNamespace,
|
||||
NAMESPACE_DEFAULT_TTL,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class VolatileCacheService:
|
||||
"""
|
||||
Service for volatile data with TTL.
|
||||
|
||||
Stores ephemeral data in Redis with automatic expiration.
|
||||
Supports multiple namespaces with configurable TTLs.
|
||||
"""
|
||||
|
||||
# Redis key prefix for volatile data
|
||||
KEY_PREFIX = "volatile"
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
redis_client: aioredis.Redis,
|
||||
settings: Settings
|
||||
):
|
||||
"""
|
||||
Initialize volatile cache service.
|
||||
|
||||
Args:
|
||||
redis_client: Async Redis client
|
||||
settings: Application settings
|
||||
"""
|
||||
self.redis = redis_client
|
||||
self.settings = settings
|
||||
|
||||
logger.info("Initialized VolatileCacheService")
|
||||
|
||||
def _build_key(self, user: str, namespace: str, key: str) -> str:
|
||||
"""
|
||||
Build Redis key for volatile record.
|
||||
|
||||
Pattern: {user}:volatile:{namespace}:{key_hash}
|
||||
Uses hash to ensure safe key characters and consistent length.
|
||||
"""
|
||||
key_hash = hashlib.md5(key.encode()).hexdigest()[:12]
|
||||
return f"{user}:{self.KEY_PREFIX}:{namespace}:{key_hash}"
|
||||
|
||||
def _build_pattern(self, user: str, namespace: Optional[str] = None) -> str:
|
||||
"""Build pattern for key scanning."""
|
||||
if namespace:
|
||||
return f"{user}:{self.KEY_PREFIX}:{namespace}:*"
|
||||
return f"{user}:{self.KEY_PREFIX}:*"
|
||||
|
||||
def _get_default_ttl(self, namespace: str) -> int:
|
||||
"""Get default TTL for a namespace."""
|
||||
try:
|
||||
ns = VolatileNamespace(namespace)
|
||||
return NAMESPACE_DEFAULT_TTL.get(ns, self.settings.volatile_default_ttl)
|
||||
except ValueError:
|
||||
return self.settings.volatile_default_ttl
|
||||
|
||||
def _serialize_record(self, record: VolatileRecord) -> str:
|
||||
"""Serialize record to JSON for storage."""
|
||||
return json.dumps({
|
||||
"key": record.key,
|
||||
"namespace": record.namespace,
|
||||
"data": record.data,
|
||||
"source": record.source,
|
||||
"created_at": record.created_at.isoformat(),
|
||||
"updated_at": record.updated_at.isoformat(),
|
||||
"ttl": record.ttl,
|
||||
"refresh_schedule": record.refresh_schedule,
|
||||
"user": record.user,
|
||||
})
|
||||
|
||||
def _deserialize_record(self, data: str) -> VolatileRecord:
|
||||
"""Deserialize record from JSON."""
|
||||
obj = json.loads(data)
|
||||
return VolatileRecord(
|
||||
key=obj["key"],
|
||||
namespace=obj["namespace"],
|
||||
data=obj["data"],
|
||||
source=obj.get("source"),
|
||||
created_at=datetime.fromisoformat(obj["created_at"]),
|
||||
updated_at=datetime.fromisoformat(obj["updated_at"]),
|
||||
ttl=obj["ttl"],
|
||||
refresh_schedule=obj.get("refresh_schedule"),
|
||||
user=obj["user"],
|
||||
)
|
||||
|
||||
async def get(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
key: str
|
||||
) -> Optional[VolatileRecordResponse]:
|
||||
"""
|
||||
Get a volatile record.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
key: Record key
|
||||
|
||||
Returns:
|
||||
Record if found and not expired, None otherwise
|
||||
"""
|
||||
redis_key = self._build_key(user, namespace, key)
|
||||
|
||||
try:
|
||||
data = await self.redis.get(redis_key)
|
||||
if not data:
|
||||
return None
|
||||
|
||||
record = self._deserialize_record(data)
|
||||
|
||||
# Get TTL remaining
|
||||
ttl_remaining = await self.redis.ttl(redis_key)
|
||||
if ttl_remaining < 0:
|
||||
return None
|
||||
|
||||
return VolatileRecordResponse(
|
||||
key=record.key,
|
||||
namespace=record.namespace,
|
||||
data=record.data,
|
||||
source=record.source,
|
||||
created_at=record.created_at,
|
||||
updated_at=record.updated_at,
|
||||
ttl=record.ttl,
|
||||
ttl_remaining=max(0, ttl_remaining),
|
||||
refresh_schedule=record.refresh_schedule,
|
||||
user=record.user,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to get volatile record {redis_key}: {e}")
|
||||
return None
|
||||
|
||||
async def set(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
key: str,
|
||||
data: Dict[str, Any],
|
||||
source: Optional[str] = None,
|
||||
ttl: Optional[int] = None,
|
||||
refresh_schedule: Optional[str] = None
|
||||
) -> VolatileRecordResponse:
|
||||
"""
|
||||
Store or update a volatile record.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
key: Record key
|
||||
data: Content to store
|
||||
source: Origin API/service
|
||||
ttl: TTL in seconds (uses namespace default if not set)
|
||||
refresh_schedule: Optional cron expression for refresh
|
||||
|
||||
Returns:
|
||||
The stored record
|
||||
"""
|
||||
redis_key = self._build_key(user, namespace, key)
|
||||
|
||||
# Use provided TTL or namespace default
|
||||
effective_ttl = ttl if ttl is not None else self._get_default_ttl(namespace)
|
||||
|
||||
# Check if record exists (for created_at)
|
||||
existing = await self.get(user, namespace, key)
|
||||
now = datetime.utcnow()
|
||||
|
||||
record = VolatileRecord(
|
||||
key=key,
|
||||
namespace=namespace,
|
||||
data=data,
|
||||
source=source,
|
||||
created_at=existing.created_at if existing else now,
|
||||
updated_at=now,
|
||||
ttl=effective_ttl,
|
||||
refresh_schedule=refresh_schedule,
|
||||
user=user,
|
||||
)
|
||||
|
||||
try:
|
||||
serialized = self._serialize_record(record)
|
||||
await self.redis.setex(redis_key, effective_ttl, serialized)
|
||||
|
||||
logger.debug(f"Stored volatile record {redis_key} with TTL {effective_ttl}s")
|
||||
|
||||
return VolatileRecordResponse(
|
||||
key=record.key,
|
||||
namespace=record.namespace,
|
||||
data=record.data,
|
||||
source=record.source,
|
||||
created_at=record.created_at,
|
||||
updated_at=record.updated_at,
|
||||
ttl=record.ttl,
|
||||
ttl_remaining=effective_ttl,
|
||||
refresh_schedule=record.refresh_schedule,
|
||||
user=record.user,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to store volatile record {redis_key}: {e}")
|
||||
raise
|
||||
|
||||
async def delete(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
key: str
|
||||
) -> bool:
|
||||
"""
|
||||
Delete a volatile record.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
key: Record key
|
||||
|
||||
Returns:
|
||||
True if record was deleted, False if not found
|
||||
"""
|
||||
redis_key = self._build_key(user, namespace, key)
|
||||
|
||||
try:
|
||||
deleted = await self.redis.delete(redis_key)
|
||||
if deleted:
|
||||
logger.debug(f"Deleted volatile record {redis_key}")
|
||||
return deleted > 0
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to delete volatile record {redis_key}: {e}")
|
||||
return False
|
||||
|
||||
async def list_namespace(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str
|
||||
) -> List[str]:
|
||||
"""
|
||||
List all keys in a namespace.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
|
||||
Returns:
|
||||
List of keys (original keys, not Redis keys)
|
||||
"""
|
||||
pattern = self._build_pattern(user, namespace)
|
||||
|
||||
try:
|
||||
keys = []
|
||||
async for redis_key in self.redis.scan_iter(match=pattern):
|
||||
# Get the record to retrieve original key
|
||||
data = await self.redis.get(redis_key)
|
||||
if data:
|
||||
record = self._deserialize_record(data)
|
||||
keys.append(record.key)
|
||||
|
||||
return keys
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to list namespace {namespace}: {e}")
|
||||
return []
|
||||
|
||||
async def get_scheduled(
|
||||
self,
|
||||
user: str
|
||||
) -> List[VolatileRecordResponse]:
|
||||
"""
|
||||
Get all records with refresh schedules.
|
||||
|
||||
Used by scheduler to determine what needs refreshing.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
List of records with refresh_schedule set
|
||||
"""
|
||||
pattern = self._build_pattern(user)
|
||||
|
||||
try:
|
||||
scheduled = []
|
||||
async for redis_key in self.redis.scan_iter(match=pattern):
|
||||
data = await self.redis.get(redis_key)
|
||||
if data:
|
||||
record = self._deserialize_record(data)
|
||||
if record.refresh_schedule:
|
||||
ttl_remaining = await self.redis.ttl(redis_key)
|
||||
scheduled.append(VolatileRecordResponse(
|
||||
key=record.key,
|
||||
namespace=record.namespace,
|
||||
data=record.data,
|
||||
source=record.source,
|
||||
created_at=record.created_at,
|
||||
updated_at=record.updated_at,
|
||||
ttl=record.ttl,
|
||||
ttl_remaining=max(0, ttl_remaining),
|
||||
refresh_schedule=record.refresh_schedule,
|
||||
user=record.user,
|
||||
))
|
||||
|
||||
return scheduled
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to get scheduled records: {e}")
|
||||
return []
|
||||
|
||||
async def get_stats(
|
||||
self,
|
||||
user: str
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Get cache statistics for user.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
Statistics dict
|
||||
"""
|
||||
pattern = self._build_pattern(user)
|
||||
|
||||
try:
|
||||
by_namespace: Dict[str, int] = {}
|
||||
total = 0
|
||||
scheduled = 0
|
||||
|
||||
async for redis_key in self.redis.scan_iter(match=pattern):
|
||||
data = await self.redis.get(redis_key)
|
||||
if data:
|
||||
record = self._deserialize_record(data)
|
||||
total += 1
|
||||
by_namespace[record.namespace] = by_namespace.get(record.namespace, 0) + 1
|
||||
if record.refresh_schedule:
|
||||
scheduled += 1
|
||||
|
||||
return {
|
||||
"total_records": total,
|
||||
"by_namespace": by_namespace,
|
||||
"scheduled_count": scheduled,
|
||||
"total_memory_bytes": None, # Could implement with DEBUG MEMORY
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to get stats: {e}")
|
||||
return {
|
||||
"total_records": 0,
|
||||
"by_namespace": {},
|
||||
"scheduled_count": 0,
|
||||
"total_memory_bytes": None,
|
||||
}
|
||||
|
||||
async def delete_namespace(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str
|
||||
) -> int:
|
||||
"""
|
||||
Delete all records in a namespace.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
|
||||
Returns:
|
||||
Number of records deleted
|
||||
"""
|
||||
pattern = self._build_pattern(user, namespace)
|
||||
|
||||
try:
|
||||
deleted = 0
|
||||
async for redis_key in self.redis.scan_iter(match=pattern):
|
||||
await self.redis.delete(redis_key)
|
||||
deleted += 1
|
||||
|
||||
logger.info(f"Deleted {deleted} records from namespace {namespace}")
|
||||
return deleted
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to delete namespace {namespace}: {e}")
|
||||
return 0
|
||||
|
||||
async def delete_all(
|
||||
self,
|
||||
user: str
|
||||
) -> int:
|
||||
"""
|
||||
Delete all volatile records for user.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
Number of records deleted
|
||||
"""
|
||||
pattern = self._build_pattern(user)
|
||||
|
||||
try:
|
||||
deleted = 0
|
||||
async for redis_key in self.redis.scan_iter(match=pattern):
|
||||
await self.redis.delete(redis_key)
|
||||
deleted += 1
|
||||
|
||||
logger.info(f"Deleted all {deleted} volatile records for user {user}")
|
||||
return deleted
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to delete all records: {e}")
|
||||
return 0
|
||||
Reference in New Issue
Block a user