340 lines
11 KiB
Python
Raw Normal View History

import asyncio
import hashlib
import time
from typing import Dict, List, Any, Optional, Tuple, Union
from dataclasses import dataclass, field
from collections import OrderedDict
import logging
logger = logging.getLogger(__name__)
@dataclass
class CacheEntry:
value: Any
created_at: float = field(default_factory=time.time)
ttl: float = 300.0
access_count: int = 0
@property
def is_expired(self) -> bool:
return time.time() - self.created_at > self.ttl
class DALCache:
def __init__(self, maxsize: int = 50000):
self.maxsize = maxsize
self._data: OrderedDict[str, CacheEntry] = OrderedDict()
self._lock = asyncio.Lock()
self._stats = {"hits": 0, "misses": 0, "evictions": 0}
async def get(self, key: str) -> Optional[Any]:
async with self._lock:
if key in self._data:
entry = self._data[key]
if entry.is_expired:
del self._data[key]
self._stats["misses"] += 1
return None
entry.access_count += 1
self._data.move_to_end(key)
self._stats["hits"] += 1
return entry.value
self._stats["misses"] += 1
return None
async def set(self, key: str, value: Any, ttl: float = 300.0):
async with self._lock:
if len(self._data) >= self.maxsize:
oldest_key = next(iter(self._data))
del self._data[oldest_key]
self._stats["evictions"] += 1
self._data[key] = CacheEntry(value=value, ttl=ttl)
async def delete(self, key: str):
async with self._lock:
self._data.pop(key, None)
async def invalidate_prefix(self, prefix: str) -> int:
async with self._lock:
keys_to_delete = [k for k in self._data.keys() if k.startswith(prefix)]
for k in keys_to_delete:
del self._data[k]
return len(keys_to_delete)
async def clear(self):
async with self._lock:
self._data.clear()
def get_stats(self) -> Dict[str, Any]:
total = self._stats["hits"] + self._stats["misses"]
hit_rate = (self._stats["hits"] / total * 100) if total > 0 else 0
return {
**self._stats,
"size": len(self._data),
"hit_rate": round(hit_rate, 2),
}
class DataAccessLayer:
TTL_USER = 600.0
TTL_FOLDER = 120.0
TTL_FILE = 120.0
TTL_PATH = 60.0
TTL_CONTENTS = 30.0
def __init__(self, cache_size: int = 50000):
self._cache = DALCache(maxsize=cache_size)
self._path_cache: Dict[str, Tuple[Any, Any, bool]] = {}
self._path_lock = asyncio.Lock()
async def start(self):
logger.info("DataAccessLayer started")
async def stop(self):
await self._cache.clear()
async with self._path_lock:
self._path_cache.clear()
logger.info("DataAccessLayer stopped")
def _user_key(self, user_id: int) -> str:
return f"user:{user_id}"
def _user_by_name_key(self, username: str) -> str:
return f"user:name:{username}"
def _folder_key(self, user_id: int, folder_id: Optional[int]) -> str:
return f"u:{user_id}:f:{folder_id or 'root'}"
def _file_key(self, user_id: int, file_id: int) -> str:
return f"u:{user_id}:file:{file_id}"
def _folder_contents_key(self, user_id: int, folder_id: Optional[int]) -> str:
return f"u:{user_id}:contents:{folder_id or 'root'}"
def _path_key(self, user_id: int, path: str) -> str:
path_hash = hashlib.md5(path.encode()).hexdigest()[:12]
return f"u:{user_id}:path:{path_hash}"
async def get_user_by_id(self, user_id: int) -> Optional[Any]:
key = self._user_key(user_id)
cached = await self._cache.get(key)
if cached is not None:
return cached
from ..models import User
user = await User.get_or_none(id=user_id)
if user:
await self._cache.set(key, user, ttl=self.TTL_USER)
await self._cache.set(self._user_by_name_key(user.username), user, ttl=self.TTL_USER)
return user
async def get_user_by_username(self, username: str) -> Optional[Any]:
key = self._user_by_name_key(username)
cached = await self._cache.get(key)
if cached is not None:
return cached
from ..models import User
user = await User.get_or_none(username=username)
if user:
await self._cache.set(key, user, ttl=self.TTL_USER)
await self._cache.set(self._user_key(user.id), user, ttl=self.TTL_USER)
return user
async def invalidate_user(self, user_id: int, username: Optional[str] = None):
await self._cache.delete(self._user_key(user_id))
if username:
await self._cache.delete(self._user_by_name_key(username))
await self._cache.invalidate_prefix(f"u:{user_id}:")
async def refresh_user_cache(self, user):
await self._cache.set(self._user_key(user.id), user, ttl=self.TTL_USER)
await self._cache.set(self._user_by_name_key(user.username), user, ttl=self.TTL_USER)
async def get_folder(
self,
user_id: int,
folder_id: Optional[int] = None,
name: Optional[str] = None,
parent_id: Optional[int] = None,
) -> Optional[Any]:
if folder_id:
key = self._folder_key(user_id, folder_id)
cached = await self._cache.get(key)
if cached is not None:
return cached
from ..models import Folder
if folder_id:
folder = await Folder.get_or_none(id=folder_id, owner_id=user_id, is_deleted=False)
elif name is not None:
folder = await Folder.get_or_none(
name=name, parent_id=parent_id, owner_id=user_id, is_deleted=False
)
else:
return None
if folder:
await self._cache.set(self._folder_key(user_id, folder.id), folder, ttl=self.TTL_FOLDER)
return folder
async def get_file(
self,
user_id: int,
file_id: Optional[int] = None,
name: Optional[str] = None,
parent_id: Optional[int] = None,
) -> Optional[Any]:
if file_id:
key = self._file_key(user_id, file_id)
cached = await self._cache.get(key)
if cached is not None:
return cached
from ..models import File
if file_id:
file = await File.get_or_none(id=file_id, owner_id=user_id, is_deleted=False)
elif name is not None:
file = await File.get_or_none(
name=name, parent_id=parent_id, owner_id=user_id, is_deleted=False
)
else:
return None
if file:
await self._cache.set(self._file_key(user_id, file.id), file, ttl=self.TTL_FILE)
return file
async def get_folder_contents(
self,
user_id: int,
folder_id: Optional[int] = None,
) -> Tuple[List[Any], List[Any]]:
key = self._folder_contents_key(user_id, folder_id)
cached = await self._cache.get(key)
if cached is not None:
return cached
from ..models import Folder, File
folders = await Folder.filter(
owner_id=user_id, parent_id=folder_id, is_deleted=False
).all()
files = await File.filter(
owner_id=user_id, parent_id=folder_id, is_deleted=False
).all()
result = (folders, files)
await self._cache.set(key, result, ttl=self.TTL_CONTENTS)
return result
async def resolve_path(
self,
user_id: int,
path_str: str,
) -> Tuple[Optional[Any], Optional[Any], bool]:
path_str = path_str.strip("/")
if not path_str:
return None, None, True
key = self._path_key(user_id, path_str)
cached = await self._cache.get(key)
if cached is not None:
return cached
from ..models import Folder, File
parts = [p for p in path_str.split("/") if p]
current_folder = None
current_folder_id = None
for i, part in enumerate(parts[:-1]):
folder = await Folder.get_or_none(
name=part, parent_id=current_folder_id, owner_id=user_id, is_deleted=False
)
if not folder:
result = (None, None, False)
await self._cache.set(key, result, ttl=self.TTL_PATH)
return result
current_folder = folder
current_folder_id = folder.id
await self._cache.set(
self._folder_key(user_id, folder.id), folder, ttl=self.TTL_FOLDER
)
last_part = parts[-1]
folder = await Folder.get_or_none(
name=last_part, parent_id=current_folder_id, owner_id=user_id, is_deleted=False
)
if folder:
await self._cache.set(
self._folder_key(user_id, folder.id), folder, ttl=self.TTL_FOLDER
)
result = (folder, current_folder, True)
await self._cache.set(key, result, ttl=self.TTL_PATH)
return result
file = await File.get_or_none(
name=last_part, parent_id=current_folder_id, owner_id=user_id, is_deleted=False
)
if file:
await self._cache.set(
self._file_key(user_id, file.id), file, ttl=self.TTL_FILE
)
result = (file, current_folder, True)
await self._cache.set(key, result, ttl=self.TTL_PATH)
return result
result = (None, current_folder, False)
await self._cache.set(key, result, ttl=self.TTL_PATH)
return result
async def invalidate_folder(self, user_id: int, folder_id: Optional[int] = None):
await self._cache.delete(self._folder_key(user_id, folder_id))
await self._cache.delete(self._folder_contents_key(user_id, folder_id))
await self._cache.invalidate_prefix(f"u:{user_id}:path:")
async def invalidate_file(self, user_id: int, file_id: int, parent_id: Optional[int] = None):
await self._cache.delete(self._file_key(user_id, file_id))
if parent_id is not None:
await self._cache.delete(self._folder_contents_key(user_id, parent_id))
await self._cache.invalidate_prefix(f"u:{user_id}:path:")
async def invalidate_path(self, user_id: int, path: str):
key = self._path_key(user_id, path)
await self._cache.delete(key)
async def invalidate_user_paths(self, user_id: int):
await self._cache.invalidate_prefix(f"u:{user_id}:path:")
await self._cache.invalidate_prefix(f"u:{user_id}:contents:")
def get_stats(self) -> Dict[str, Any]:
return self._cache.get_stats()
_dal: Optional[DataAccessLayer] = None
async def init_dal(cache_size: int = 50000) -> DataAccessLayer:
global _dal
_dal = DataAccessLayer(cache_size=cache_size)
await _dal.start()
return _dal
async def shutdown_dal():
global _dal
if _dal:
await _dal.stop()
_dal = None
def get_dal() -> DataAccessLayer:
if not _dal:
raise RuntimeError("DataAccessLayer not initialized")
return _dal