mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-30 04:27:40 +08:00
1175 lines
40 KiB
Python
1175 lines
40 KiB
Python
import contextvars
|
|
import inspect
|
|
import logging
|
|
import threading
|
|
from abc import ABC, abstractmethod
|
|
from contextlib import contextmanager, asynccontextmanager
|
|
from functools import wraps
|
|
from pathlib import Path
|
|
from typing import Any, Callable, Dict, Optional, Generator, AsyncGenerator, Tuple, Literal, Union
|
|
|
|
from cachetools import LRUCache as MemoryLRUCache
|
|
from cachetools import TLRUCache as MemoryTLRUCache
|
|
from cachetools.keys import hashkey
|
|
|
|
# 默认缓存区
|
|
DEFAULT_CACHE_REGION = "DEFAULT"
|
|
# 默认缓存大小
|
|
DEFAULT_CACHE_SIZE = 1024
|
|
# 默认缓存有效期
|
|
DEFAULT_CACHE_TTL = 365 * 24 * 60 * 60
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_backend_type_provider: Callable[[], str] = lambda: "memory"
|
|
_redis_factory: Optional[Callable[[Optional[int]], "CacheBackend"]] = None
|
|
_async_redis_factory: Optional[
|
|
Callable[[Optional[int]], "AsyncCacheBackend"]
|
|
] = None
|
|
_file_factory: Optional[Callable[[Optional[Path]], "CacheBackend"]] = None
|
|
_async_file_factory: Optional[
|
|
Callable[[Optional[Path]], "AsyncCacheBackend"]
|
|
] = None
|
|
_file_ttl_provider: Callable[[], int] = lambda: DEFAULT_CACHE_TTL
|
|
|
|
|
|
def configure_cache_factories(
|
|
*,
|
|
backend_type_provider: Callable[[], str],
|
|
redis_factory: Callable[[Optional[int]], "CacheBackend"],
|
|
async_redis_factory: Callable[[Optional[int]], "AsyncCacheBackend"],
|
|
file_factory: Callable[[Optional[Path]], "CacheBackend"],
|
|
async_file_factory: Callable[[Optional[Path]], "AsyncCacheBackend"],
|
|
file_ttl_provider: Callable[[], int],
|
|
) -> None:
|
|
"""注入配置读取器和具体缓存适配器工厂。"""
|
|
global _backend_type_provider, _redis_factory, _async_redis_factory
|
|
global _file_factory, _async_file_factory, _file_ttl_provider
|
|
_backend_type_provider = backend_type_provider
|
|
_redis_factory = redis_factory
|
|
_async_redis_factory = async_redis_factory
|
|
_file_factory = file_factory
|
|
_async_file_factory = async_file_factory
|
|
_file_ttl_provider = file_ttl_provider
|
|
|
|
# 上下文变量来控制缓存行为
|
|
_fresh = contextvars.ContextVar('fresh', default=False)
|
|
|
|
|
|
class CacheBackend(ABC):
|
|
"""
|
|
缓存后端基类,定义通用的缓存接口
|
|
"""
|
|
|
|
def __getitem__(self, key: str) -> Any:
|
|
"""
|
|
获取缓存项,类似 dict[key]
|
|
"""
|
|
value = self.get(key)
|
|
if value is None:
|
|
raise KeyError(key)
|
|
return value
|
|
|
|
def __setitem__(self, key: str, value: Any) -> None:
|
|
"""
|
|
设置缓存项,类似 dict[key] = value
|
|
"""
|
|
self.set(key, value)
|
|
|
|
def __delitem__(self, key: str) -> None:
|
|
"""
|
|
删除缓存项,类似 del dict[key]
|
|
"""
|
|
if not self.exists(key):
|
|
raise KeyError(key)
|
|
self.delete(key)
|
|
|
|
def __contains__(self, key: str) -> bool:
|
|
"""
|
|
检查键是否存在,类似 key in dict
|
|
"""
|
|
return self.exists(key)
|
|
|
|
def __iter__(self):
|
|
"""
|
|
返回缓存的迭代器,类似 iter(dict)
|
|
"""
|
|
for key, _ in self.items():
|
|
yield key
|
|
|
|
def __len__(self) -> int:
|
|
"""
|
|
返回缓存项的数量,类似 len(dict)
|
|
"""
|
|
return sum(1 for _ in self.items())
|
|
|
|
@abstractmethod
|
|
def set(self, key: str, value: Any, ttl: Optional[int] = None,
|
|
region: Optional[str] = DEFAULT_CACHE_REGION, **kwargs) -> None:
|
|
"""
|
|
设置缓存
|
|
|
|
:param key: 缓存的键
|
|
:param value: 缓存的值
|
|
:param ttl: 缓存的存活时间,单位秒
|
|
:param region: 缓存的区
|
|
:param kwargs: 其他参数
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
def exists(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> bool:
|
|
"""
|
|
判断缓存键是否存在
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 存在返回 True,否则返回 False
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
def get(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> Any:
|
|
"""
|
|
获取缓存
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 返回缓存的值,如果缓存不存在返回 None
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
def delete(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> None:
|
|
"""
|
|
删除缓存
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
def clear(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> None:
|
|
"""
|
|
清除指定区域的缓存或全部缓存
|
|
|
|
:param region: 缓存的区,为None时清空所有区缓存
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
def items(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> Generator[Tuple[str, Any], None, None]:
|
|
"""
|
|
获取指定区域的所有缓存项
|
|
|
|
:param region: 缓存的区
|
|
:return: 返回一个字典,包含所有缓存键值对
|
|
"""
|
|
pass
|
|
|
|
def keys(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> Generator[str, None, None]:
|
|
"""
|
|
获取所有缓存键,类似 dict.keys()
|
|
"""
|
|
for key, _ in self.items(region=region):
|
|
yield key
|
|
|
|
def values(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> Generator[Any, None, None]:
|
|
"""
|
|
获取所有缓存值,类似 dict.values()
|
|
"""
|
|
for _, value in self.items(region=region):
|
|
yield value
|
|
|
|
def update(self, other: Dict[str, Any], region: Optional[str] = DEFAULT_CACHE_REGION,
|
|
ttl: Optional[int] = None, **kwargs) -> None:
|
|
"""
|
|
更新缓存,类似 dict.update()
|
|
"""
|
|
for key, value in other.items():
|
|
self.set(key, value, ttl=ttl, region=region, **kwargs)
|
|
|
|
def pop(self, key: str, default: Any = None, region: Optional[str] = DEFAULT_CACHE_REGION) -> Any:
|
|
"""
|
|
弹出缓存项,类似 dict.pop()
|
|
"""
|
|
value = self.get(key, region=region)
|
|
if value is not None:
|
|
self.delete(key, region=region)
|
|
return value
|
|
if default is not None:
|
|
return default
|
|
raise KeyError(key)
|
|
|
|
def popitem(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> Tuple[str, Any]:
|
|
"""
|
|
弹出最后一个缓存项,类似 dict.popitem()
|
|
"""
|
|
items = list(self.items(region=region))
|
|
if not items:
|
|
raise KeyError("popitem(): cache is empty")
|
|
key, value = items[-1]
|
|
self.delete(key, region=region)
|
|
return key, value
|
|
|
|
def setdefault(self, key: str, default: Any = None, region: Optional[str] = DEFAULT_CACHE_REGION,
|
|
ttl: Optional[int] = None, **kwargs) -> Any:
|
|
"""
|
|
设置默认值,类似 dict.setdefault()
|
|
"""
|
|
value = self.get(key, region=region)
|
|
if value is None:
|
|
self.set(key, default, ttl=ttl, region=region, **kwargs)
|
|
return default
|
|
return value
|
|
|
|
@abstractmethod
|
|
def close(self) -> None:
|
|
"""
|
|
关闭缓存连接
|
|
"""
|
|
pass
|
|
|
|
@staticmethod
|
|
def get_region(region: Optional[str] = None) -> str:
|
|
"""
|
|
获取缓存的区
|
|
"""
|
|
return f"region:{region}" if region else "region:DEFAULT"
|
|
|
|
@staticmethod
|
|
def is_redis() -> bool:
|
|
"""
|
|
判断当前缓存后端是否为 Redis
|
|
"""
|
|
return False
|
|
|
|
|
|
class AsyncCacheBackend(CacheBackend):
|
|
"""
|
|
缓存后端基类,定义通用的缓存接口(异步)
|
|
"""
|
|
|
|
@abstractmethod
|
|
async def set(self, key: str, value: Any, ttl: Optional[int] = None,
|
|
region: Optional[str] = DEFAULT_CACHE_REGION, **kwargs) -> None:
|
|
"""
|
|
设置缓存
|
|
|
|
:param key: 缓存的键
|
|
:param value: 缓存的值
|
|
:param ttl: 缓存的存活时间,单位秒
|
|
:param region: 缓存的区
|
|
:param kwargs: 其他参数
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def exists(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> bool:
|
|
"""
|
|
判断缓存键是否存在
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 存在返回 True,否则返回 False
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> Any:
|
|
"""
|
|
获取缓存
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 返回缓存的值,如果缓存不存在返回 None
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def delete(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> None:
|
|
"""
|
|
删除缓存
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def clear(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> None:
|
|
"""
|
|
清除指定区域的缓存或全部缓存
|
|
|
|
:param region: 缓存的区,为None时清空所有区缓存
|
|
"""
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def items(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> AsyncGenerator[Tuple[str, Any], None]:
|
|
"""
|
|
获取指定区域的所有缓存项
|
|
|
|
:param region: 缓存的区
|
|
:return: 返回一个字典,包含所有缓存键值对
|
|
"""
|
|
pass
|
|
|
|
async def keys(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> AsyncGenerator[str, None]:
|
|
"""
|
|
获取所有缓存键,类似 dict.keys()(异步)
|
|
"""
|
|
async for key, _ in self.items(region=region):
|
|
yield key
|
|
|
|
async def values(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> AsyncGenerator[Any, None]:
|
|
"""
|
|
获取所有缓存值,类似 dict.values()(异步)
|
|
"""
|
|
async for _, value in self.items(region=region):
|
|
yield value
|
|
|
|
async def update(self, other: Dict[str, Any], region: Optional[str] = DEFAULT_CACHE_REGION,
|
|
ttl: Optional[int] = None, **kwargs) -> None:
|
|
"""
|
|
更新缓存,类似 dict.update()(异步)
|
|
"""
|
|
for key, value in other.items():
|
|
await self.set(key, value, ttl=ttl, region=region, **kwargs)
|
|
|
|
async def pop(self, key: str, default: Any = None, region: Optional[str] = DEFAULT_CACHE_REGION) -> Any:
|
|
"""
|
|
弹出缓存项,类似 dict.pop()(异步)
|
|
"""
|
|
value = await self.get(key, region=region)
|
|
if value is not None:
|
|
await self.delete(key, region=region)
|
|
return value
|
|
if default is not None:
|
|
return default
|
|
raise KeyError(key)
|
|
|
|
async def popitem(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> Tuple[str, Any]:
|
|
"""
|
|
弹出最后一个缓存项,类似 dict.popitem()(异步)
|
|
"""
|
|
items = []
|
|
async for item in self.items(region=region):
|
|
items.append(item)
|
|
if not items:
|
|
raise KeyError("popitem(): cache is empty")
|
|
key, value = items[-1]
|
|
await self.delete(key, region=region)
|
|
return key, value
|
|
|
|
async def setdefault(self, key: str, default: Any = None, region: Optional[str] = DEFAULT_CACHE_REGION,
|
|
ttl: Optional[int] = None, **kwargs) -> Any:
|
|
"""
|
|
设置默认值,类似 dict.setdefault()(异步)
|
|
"""
|
|
value = await self.get(key, region=region)
|
|
if value is None:
|
|
await self.set(key, default, ttl=ttl, region=region, **kwargs)
|
|
return default
|
|
return value
|
|
|
|
@abstractmethod
|
|
async def close(self) -> None:
|
|
"""
|
|
关闭缓存连接
|
|
"""
|
|
pass
|
|
|
|
|
|
class _MemoryTLRUCache(MemoryTLRUCache):
|
|
"""
|
|
支持为每个 key 设置独立 TTL 的内存缓存
|
|
"""
|
|
|
|
def __init__(self, maxsize: int, ttl: int):
|
|
"""使用默认 TTL 初始化支持单键过期时间的内存缓存。"""
|
|
self.__ttl = ttl
|
|
self.__setting_ttls: Dict[str, int] = {}
|
|
super().__init__(maxsize=maxsize, ttu=self._get_expiration)
|
|
|
|
def _get_expiration(self, key: str, _value: Any, now: float) -> float:
|
|
"""计算指定键的绝对过期时间。"""
|
|
return now + self.__setting_ttls.get(key, self.__ttl)
|
|
|
|
@property
|
|
def ttl(self) -> int:
|
|
"""
|
|
默认缓存存活时间,单位秒
|
|
"""
|
|
return self.__ttl
|
|
|
|
def set(self, key: str, value: Any, ttl: int) -> None:
|
|
"""
|
|
使用指定 TTL 设置缓存值
|
|
"""
|
|
if ttl <= 0:
|
|
try:
|
|
del self[key]
|
|
except KeyError:
|
|
pass
|
|
return
|
|
self.__setting_ttls[key] = ttl
|
|
try:
|
|
super().__setitem__(key, value)
|
|
finally:
|
|
self.__setting_ttls.pop(key, None)
|
|
|
|
|
|
class MemoryBackend(CacheBackend):
|
|
"""
|
|
基于 `cachetools.TLRUCache` 实现的缓存后端
|
|
"""
|
|
|
|
# 类变量 _region_caches 的互斥锁
|
|
_lock = threading.Lock()
|
|
# 存储各个 region 的缓存实例,region -> TLRUCache/LRUCache
|
|
_region_caches: Dict[str, Union[_MemoryTLRUCache, MemoryLRUCache]] = {}
|
|
|
|
def __init__(self, cache_type: Literal['ttl', 'lru'] = 'ttl',
|
|
maxsize: Optional[int] = None, ttl: Optional[int] = None):
|
|
"""
|
|
初始化缓存实例
|
|
|
|
:param cache_type: 缓存类型,支持 'ttl'(默认)和 'lru'
|
|
:param maxsize: 缓存的最大条目数
|
|
:param ttl: 默认缓存存活时间,单位秒
|
|
"""
|
|
self.cache_type = cache_type
|
|
self.maxsize = maxsize or DEFAULT_CACHE_SIZE
|
|
self.ttl = DEFAULT_CACHE_TTL if ttl is None else ttl
|
|
|
|
def __get_region_cache(self, region: str) -> Optional[Union[_MemoryTLRUCache, MemoryLRUCache]]:
|
|
"""
|
|
获取指定区域的缓存实例,如果不存在则返回 None
|
|
"""
|
|
region = self.get_region(region)
|
|
return self._region_caches.get(region)
|
|
|
|
def set(self, key: str, value: Any, ttl: Optional[int] = None,
|
|
region: Optional[str] = DEFAULT_CACHE_REGION, **kwargs) -> None:
|
|
"""
|
|
设置缓存值支持每个 key 独立配置 TTL
|
|
|
|
:param key: 缓存的键
|
|
:param value: 缓存的值
|
|
:param ttl: 缓存的存活时间,未传入则使用 backend 默认值,单位秒
|
|
:param region: 缓存的区
|
|
"""
|
|
ttl = self.ttl if ttl is None else ttl
|
|
maxsize = kwargs.get("maxsize") or self.maxsize
|
|
region = self.get_region(region)
|
|
# 设置缓存值
|
|
with self._lock:
|
|
region_cache = self._region_caches.get(region)
|
|
if region_cache is None:
|
|
region_cache = (
|
|
_MemoryTLRUCache(maxsize=maxsize, ttl=ttl) if self.cache_type == 'ttl'
|
|
else MemoryLRUCache(maxsize=maxsize)
|
|
)
|
|
self._region_caches[region] = region_cache
|
|
elif isinstance(region_cache, _MemoryTLRUCache) != (self.cache_type == 'ttl'):
|
|
raise ValueError(
|
|
f"Cache region {region!r} already uses a different cache type"
|
|
)
|
|
if isinstance(region_cache, _MemoryTLRUCache):
|
|
region_cache.set(key, value, ttl=ttl)
|
|
else:
|
|
region_cache[key] = value
|
|
|
|
def exists(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> bool:
|
|
"""
|
|
判断缓存键是否存在
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 存在返回 True,否则返回 False
|
|
"""
|
|
region_cache = self.__get_region_cache(region)
|
|
if region_cache is None:
|
|
return False
|
|
with self._lock:
|
|
return key in region_cache
|
|
|
|
def get(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> Any:
|
|
"""
|
|
获取缓存的值
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 返回缓存的值,如果缓存不存在返回 None
|
|
"""
|
|
region_cache = self.__get_region_cache(region)
|
|
if region_cache is None:
|
|
return None
|
|
with self._lock:
|
|
return region_cache.get(key)
|
|
|
|
def delete(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION):
|
|
"""
|
|
删除缓存
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
"""
|
|
region_cache = self.__get_region_cache(region)
|
|
if region_cache is None:
|
|
return
|
|
with self._lock:
|
|
if key in region_cache:
|
|
del region_cache[key]
|
|
|
|
def clear(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> None:
|
|
"""
|
|
清除指定区域的缓存或全部缓存
|
|
|
|
:param region: 缓存的区,为None时清空所有区缓存
|
|
"""
|
|
with self._lock:
|
|
if region:
|
|
# 清理指定缓存区
|
|
region_cache = self.__get_region_cache(region)
|
|
if region_cache is not None:
|
|
region_cache.clear()
|
|
logger.debug(f"Cleared cache for region: {region}")
|
|
else:
|
|
# 清除所有区域的缓存
|
|
for region_cache in self._region_caches.values():
|
|
region_cache.clear()
|
|
logger.info("Cleared all cache")
|
|
|
|
def items(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> Generator[Tuple[str, Any], None, None]:
|
|
"""
|
|
获取指定区域的所有缓存项
|
|
|
|
:param region: 缓存的区
|
|
:return: 返回一个字典,包含所有缓存键值对
|
|
"""
|
|
region_cache = self.__get_region_cache(region)
|
|
if region_cache is None:
|
|
yield from ()
|
|
return
|
|
# 使用锁保护迭代过程,避免在迭代时缓存被修改
|
|
with self._lock:
|
|
# 创建快照避免并发修改问题
|
|
items_snapshot = list(region_cache.items())
|
|
for item in items_snapshot:
|
|
yield item
|
|
|
|
def close(self) -> None:
|
|
"""
|
|
内存缓存不需要关闭资源
|
|
"""
|
|
pass
|
|
|
|
|
|
class AsyncMemoryBackend(AsyncCacheBackend):
|
|
"""
|
|
基于 `cachetools.TTLCache` 实现的异步缓存后端
|
|
"""
|
|
|
|
def __init__(self, cache_type: Literal['ttl', 'lru'] = 'ttl',
|
|
maxsize: Optional[int] = None, ttl: Optional[int] = None):
|
|
"""
|
|
初始化缓存实例
|
|
|
|
:param cache_type: 缓存类型,支持 'ttl'(默认)和 'lru'
|
|
:param maxsize: 缓存的最大条目数
|
|
:param ttl: 默认缓存存活时间,单位秒
|
|
"""
|
|
self._backend = MemoryBackend(cache_type=cache_type, maxsize=maxsize, ttl=ttl)
|
|
|
|
async def set(self, key: str, value: Any, ttl: Optional[int] = None,
|
|
region: Optional[str] = DEFAULT_CACHE_REGION, **kwargs) -> None:
|
|
"""
|
|
设置缓存值支持每个 key 独立配置 TTL
|
|
|
|
:param key: 缓存的键
|
|
:param value: 缓存的值
|
|
:param ttl: 缓存的存活时间,未传入则使用 backend 默认值,单位秒
|
|
:param region: 缓存的区
|
|
"""
|
|
return self._backend.set(key=key, value=value, ttl=ttl, region=region, **kwargs)
|
|
|
|
async def exists(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> bool:
|
|
"""
|
|
判断缓存键是否存在
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 存在返回 True,否则返回 False
|
|
"""
|
|
return self._backend.exists(key=key, region=region)
|
|
|
|
async def get(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION) -> Any:
|
|
"""
|
|
获取缓存的值
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
:return: 返回缓存的值,如果缓存不存在返回 None
|
|
"""
|
|
return self._backend.get(key=key, region=region)
|
|
|
|
async def delete(self, key: str, region: Optional[str] = DEFAULT_CACHE_REGION):
|
|
"""
|
|
删除缓存
|
|
|
|
:param key: 缓存的键
|
|
:param region: 缓存的区
|
|
"""
|
|
return self._backend.delete(key=key, region=region)
|
|
|
|
async def clear(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> None:
|
|
"""
|
|
清除指定区域的缓存或全部缓存
|
|
|
|
:param region: 缓存的区,为None时清空所有区缓存
|
|
"""
|
|
return self._backend.clear(region=region)
|
|
|
|
async def items(self, region: Optional[str] = DEFAULT_CACHE_REGION) -> AsyncGenerator[Tuple[str, Any], None]:
|
|
"""
|
|
获取指定区域的所有缓存项
|
|
|
|
:param region: 缓存的区
|
|
:return: 返回一个字典,包含所有缓存键值对
|
|
"""
|
|
for item in self._backend.items(region):
|
|
yield item
|
|
|
|
async def close(self) -> None:
|
|
"""
|
|
内存缓存不需要关闭资源
|
|
"""
|
|
pass
|
|
|
|
|
|
@contextmanager
|
|
def fresh(fresh: bool = True):
|
|
"""
|
|
是否获取新数据(不使用缓存的值)
|
|
|
|
Usage:
|
|
with fresh():
|
|
result = some_cached_function()
|
|
"""
|
|
token = _fresh.set(fresh or is_fresh())
|
|
try:
|
|
yield
|
|
finally:
|
|
_fresh.reset(token)
|
|
|
|
@asynccontextmanager
|
|
async def async_fresh(fresh: bool = True):
|
|
"""
|
|
是否获取新数据(不使用缓存的值)
|
|
|
|
Usage:
|
|
async with async_fresh():
|
|
result = await some_async_cached_function()
|
|
"""
|
|
token = _fresh.set(fresh or is_fresh())
|
|
try:
|
|
yield
|
|
finally:
|
|
_fresh.reset(token)
|
|
|
|
def is_fresh() -> bool:
|
|
"""
|
|
是否获取新数据
|
|
"""
|
|
try:
|
|
return _fresh.get()
|
|
except LookupError:
|
|
return False
|
|
|
|
def FileCache(base: Optional[Path] = None, ttl: Optional[int] = None) -> CacheBackend:
|
|
"""
|
|
获取文件缓存后端实例(Redis或文件系统),ttl仅在Redis环境中有效
|
|
"""
|
|
if _backend_type_provider() == "redis":
|
|
if _redis_factory is None:
|
|
raise RuntimeError("Redis 缓存适配器尚未由启动层配置")
|
|
return _redis_factory(ttl if ttl is not None else _file_ttl_provider())
|
|
if _file_factory is None:
|
|
raise RuntimeError("文件缓存适配器尚未由启动层配置")
|
|
return _file_factory(base)
|
|
|
|
|
|
def AsyncFileCache(
|
|
base: Optional[Path] = None,
|
|
ttl: Optional[int] = None,
|
|
) -> AsyncCacheBackend:
|
|
"""
|
|
获取文件异步缓存后端实例(Redis或文件系统),ttl仅在Redis环境中有效
|
|
"""
|
|
if _backend_type_provider() == "redis":
|
|
if _async_redis_factory is None:
|
|
raise RuntimeError("异步 Redis 缓存适配器尚未由启动层配置")
|
|
return _async_redis_factory(
|
|
ttl if ttl is not None else _file_ttl_provider()
|
|
)
|
|
if _async_file_factory is None:
|
|
raise RuntimeError("异步文件缓存适配器尚未由启动层配置")
|
|
return _async_file_factory(base)
|
|
|
|
|
|
def Cache(cache_type: Literal['ttl', 'lru'] = 'ttl',
|
|
maxsize: Optional[int] = None,
|
|
ttl: Optional[int] = None) -> CacheBackend:
|
|
"""
|
|
根据配置获取缓存后端实例(内存或Redis),maxsize仅在未启用Redis时生效
|
|
|
|
:param cache_type: 缓存类型,仅使用内存缓存时生效,支持 'ttl'(默认)和 'lru'
|
|
:param maxsize: 缓存的最大条目数,仅使用cachetools时生效
|
|
:param ttl: 缓存的默认存活时间,单位秒
|
|
:return: 返回缓存后端实例
|
|
"""
|
|
if _backend_type_provider() == "redis":
|
|
if _redis_factory is None:
|
|
raise RuntimeError("Redis 缓存适配器尚未由启动层配置")
|
|
return _redis_factory(ttl)
|
|
return MemoryBackend(cache_type=cache_type, maxsize=maxsize, ttl=ttl)
|
|
|
|
|
|
def AsyncCache(cache_type: Literal['ttl', 'lru'] = 'ttl',
|
|
maxsize: Optional[int] = None,
|
|
ttl: Optional[int] = None) -> AsyncCacheBackend:
|
|
"""
|
|
根据配置获取异步缓存后端实例(内存或Redis),maxsize仅在未启用Redis时生效
|
|
|
|
:param cache_type: 缓存类型,仅使用内存缓存时生效,支持 'ttl'(默认)和 'lru'
|
|
:param maxsize: 缓存的最大条目数,仅使用cachetools时生效
|
|
:param ttl: 缓存的默认存活时间,单位秒
|
|
:return: 返回异步缓存后端实例
|
|
"""
|
|
if _backend_type_provider() == "redis":
|
|
if _async_redis_factory is None:
|
|
raise RuntimeError("异步 Redis 缓存适配器尚未由启动层配置")
|
|
return _async_redis_factory(ttl)
|
|
return AsyncMemoryBackend(cache_type=cache_type, maxsize=maxsize, ttl=ttl)
|
|
|
|
|
|
def cached(region: Optional[str] = None, maxsize: Optional[int] = 1024, ttl: Optional[int] = None,
|
|
skip_none: Optional[bool] = True, skip_empty: Optional[bool] = False, shared_key: Optional[str] = None,
|
|
skip_if: Optional[Callable[[Any], bool]] = None):
|
|
"""
|
|
自定义缓存装饰器,支持配置缓存区域的 maxsize 和每个 key 的 ttl
|
|
|
|
:param region: 缓存区域的标识符,默认根据模块名、函数名等自动生成标识
|
|
:param maxsize: 缓存区内的最大条目数
|
|
:param ttl: 缓存的存活时间,单位秒;未传入时使用 LRU 缓存
|
|
:param skip_none: 跳过 None 缓存,默认为 True
|
|
:param skip_empty: 跳过空值缓存(如 None, [], {}, "", set()),默认为 False
|
|
:param shared_key: 同步/异步函数共享缓存的键,默认使用函数名(异步函数名会标准化为同步格式,如移除 `async_` 前缀)
|
|
:param skip_if: 按返回值判断是否跳过缓存的谓词,返回真值时不缓存;用于
|
|
「结构合法但业务失败」的返回值(如 TMDB 的 success=false 响应),
|
|
这类值无法用 skip_none/skip_empty 表达
|
|
:return: 装饰器函数
|
|
"""
|
|
|
|
def decorator(func):
|
|
# 函数签名在装饰后不会变化,预计算可避免每次缓存访问都重复反射。
|
|
signature = inspect.signature(func)
|
|
parameter_names = list(signature.parameters.keys())
|
|
cache_parameter_names = (
|
|
parameter_names[1:]
|
|
if parameter_names and parameter_names[0] in ("self", "cls")
|
|
else parameter_names
|
|
)
|
|
|
|
def should_cache(value: Any) -> bool:
|
|
"""
|
|
判断是否应该缓存结果,如果返回值是 None 或空值则不缓存
|
|
|
|
:param value: 要判断的缓存值
|
|
:return: 是否缓存结果
|
|
"""
|
|
if skip_none and value is None:
|
|
return False
|
|
# if skip_empty and value in [None, [], {}, "", set()]:
|
|
if skip_empty and not value:
|
|
return False
|
|
if skip_if and skip_if(value):
|
|
return False
|
|
return True
|
|
|
|
def is_valid_cache_value(_cache_key: str, _cached_value: Any, _cache_region: str) -> bool:
|
|
"""
|
|
判断指定的值是否为一个有效的缓存值
|
|
|
|
:param _cache_key: 缓存的键
|
|
:param _cached_value: 缓存的值
|
|
:param _cache_region: 缓存的区
|
|
:return: 若值是有效的缓存值返回 True,否则返回 False
|
|
"""
|
|
# 如果 skip_none 为 False,且 value 为 None,需要判断缓存实际是否存在
|
|
if not skip_none and _cached_value is None:
|
|
if not cache_backend.exists(key=_cache_key, region=_cache_region):
|
|
return False
|
|
return True
|
|
|
|
async def async_is_valid_cache_value(_cache_key: str, _cached_value: Any, _cache_region: str) -> bool:
|
|
"""
|
|
判断指定的值是否为一个有效的缓存值(异步版本)
|
|
|
|
:param _cache_key: 缓存的键
|
|
:param _cached_value: 缓存的值
|
|
:param _cache_region: 缓存的区
|
|
:return: 若值是有效的缓存值返回 True,否则返回 False
|
|
"""
|
|
# 如果 skip_none 为 False,且 value 为 None,需要判断缓存实际是否存在
|
|
if not skip_none and _cached_value is None:
|
|
if not await cache_backend.exists(key=_cache_key, region=_cache_region):
|
|
return False
|
|
return True
|
|
|
|
def __standardize_func_name() -> str:
|
|
"""
|
|
将异步函数名标准化为同步函数的命名,以生成统一的缓存键
|
|
"""
|
|
# XXX 假设异步函数名与同步版本仅差`async_`前缀或`_async`后缀(当前MP代码大多符合),否则需通过`shared_key`参数显式指定
|
|
return (
|
|
func.__name__.removeprefix("async_").removesuffix("_async")
|
|
if is_async
|
|
else func.__name__
|
|
)
|
|
|
|
def __get_cache_key(args, kwargs) -> str:
|
|
"""
|
|
根据函数和参数生成缓存键
|
|
|
|
:param args: 位置参数
|
|
:param kwargs: 关键字参数
|
|
:return: 缓存键
|
|
"""
|
|
# 绑定传入的参数并应用默认值
|
|
bound = signature.bind(*args, **kwargs)
|
|
bound.apply_defaults()
|
|
# 按照函数签名顺序提取参数值列表
|
|
keys = [
|
|
bound.arguments[param] for param in cache_parameter_names if param in bound.arguments
|
|
]
|
|
# 使用有序参数生成缓存键
|
|
return f"{func_name}_{hashkey(*keys)}"
|
|
|
|
# 被装饰函数的上层名称(如类名或外层函数名)
|
|
enclosing_name = (
|
|
func.__qualname__[:last_dot]
|
|
if (last_dot := func.__qualname__.rfind(".")) != -1
|
|
else ""
|
|
)
|
|
# 检查是否为异步函数
|
|
is_async = inspect.iscoroutinefunction(func)
|
|
# 生成标准化后的函数名称,用于同步/异步函数共享缓存
|
|
func_name = shared_key if shared_key else __standardize_func_name()
|
|
# 获取缓存区
|
|
cache_region = (
|
|
region if region is not None else f"{func.__module__}:{enclosing_name}:{func_name}"
|
|
)
|
|
|
|
if is_async:
|
|
# 异步函数使用异步缓存后端
|
|
cache_backend = AsyncCache(cache_type="ttl" if ttl is not None else "lru", maxsize=maxsize, ttl=ttl)
|
|
# 异步函数的缓存装饰器
|
|
@wraps(func)
|
|
async def async_wrapper(*args, **kwargs):
|
|
# 获取缓存键
|
|
cache_key = __get_cache_key(args, kwargs)
|
|
|
|
if not is_fresh():
|
|
# 尝试获取缓存
|
|
cached_value = await cache_backend.get(cache_key, region=cache_region)
|
|
if should_cache(cached_value) and await async_is_valid_cache_value(cache_key, cached_value,
|
|
cache_region):
|
|
return cached_value
|
|
# 执行异步函数并缓存结果
|
|
result = await func(*args, **kwargs)
|
|
# 判断是否需要缓存
|
|
if not should_cache(result):
|
|
return result
|
|
# 设置缓存(如果有传入的 maxsize 和 ttl,则覆盖默认值)
|
|
await cache_backend.set(cache_key, result, ttl=ttl, maxsize=maxsize, region=cache_region)
|
|
return result
|
|
|
|
async def cache_clear():
|
|
"""
|
|
清理缓存区
|
|
"""
|
|
await cache_backend.clear(region=cache_region)
|
|
|
|
async def cache_exists(*args, **kwargs) -> bool:
|
|
"""
|
|
判断当前参数对应的有效缓存是否存在。
|
|
"""
|
|
cache_key = __get_cache_key(args, kwargs)
|
|
cached_value = await cache_backend.get(cache_key, region=cache_region)
|
|
return should_cache(cached_value) and await async_is_valid_cache_value(
|
|
cache_key, cached_value, cache_region
|
|
)
|
|
|
|
async def cache_delete(*args, **kwargs) -> None:
|
|
"""
|
|
删除当前参数对应的缓存。
|
|
"""
|
|
cache_key = __get_cache_key(args, kwargs)
|
|
await cache_backend.delete(cache_key, region=cache_region)
|
|
|
|
async_wrapper.cache_region = cache_region
|
|
async_wrapper.cache_clear = cache_clear
|
|
async_wrapper.cache_exists = cache_exists
|
|
async_wrapper.cache_delete = cache_delete
|
|
return async_wrapper
|
|
else:
|
|
# 同步函数使用同步缓存后端
|
|
cache_backend = Cache(cache_type="ttl" if ttl is not None else "lru", maxsize=maxsize, ttl=ttl)
|
|
# 同步函数的缓存装饰器
|
|
@wraps(func)
|
|
def wrapper(*args, **kwargs):
|
|
# 获取缓存键
|
|
cache_key = __get_cache_key(args, kwargs)
|
|
|
|
if not is_fresh():
|
|
# 尝试获取缓存
|
|
cached_value = cache_backend.get(cache_key, region=cache_region)
|
|
if should_cache(cached_value) and is_valid_cache_value(cache_key, cached_value, cache_region):
|
|
return cached_value
|
|
# 执行函数并缓存结果
|
|
result = func(*args, **kwargs)
|
|
# 判断是否需要缓存
|
|
if not should_cache(result):
|
|
return result
|
|
# 设置缓存(如果有传入的 maxsize 和 ttl,则覆盖默认值)
|
|
cache_backend.set(cache_key, result, ttl=ttl, maxsize=maxsize, region=cache_region)
|
|
return result
|
|
|
|
def cache_clear():
|
|
"""
|
|
清理缓存区
|
|
"""
|
|
cache_backend.clear(region=cache_region)
|
|
|
|
def cache_exists(*args, **kwargs) -> bool:
|
|
"""
|
|
判断当前参数对应的有效缓存是否存在。
|
|
"""
|
|
cache_key = __get_cache_key(args, kwargs)
|
|
cached_value = cache_backend.get(cache_key, region=cache_region)
|
|
return should_cache(cached_value) and is_valid_cache_value(
|
|
cache_key, cached_value, cache_region
|
|
)
|
|
|
|
def cache_delete(*args, **kwargs) -> None:
|
|
"""
|
|
删除当前参数对应的缓存。
|
|
"""
|
|
cache_key = __get_cache_key(args, kwargs)
|
|
cache_backend.delete(cache_key, region=cache_region)
|
|
|
|
wrapper.cache_region = cache_region
|
|
wrapper.cache_clear = cache_clear
|
|
wrapper.cache_exists = cache_exists
|
|
wrapper.cache_delete = cache_delete
|
|
return wrapper
|
|
|
|
return decorator
|
|
|
|
|
|
class CacheProxy:
|
|
"""
|
|
缓存代理类,将缓存后端的方法直接代理到实例上
|
|
"""
|
|
|
|
def __init__(self, cache_backend: CacheBackend, region: str):
|
|
"""
|
|
初始化缓存代理
|
|
|
|
:param cache_backend: 缓存后端实例
|
|
:param region: 缓存区域
|
|
"""
|
|
self._cache_backend = cache_backend
|
|
self._region = region
|
|
|
|
def __getitem__(self, key):
|
|
"""
|
|
获取缓存项
|
|
"""
|
|
value = self._cache_backend.get(key, region=self._region)
|
|
if value is None:
|
|
raise KeyError(key)
|
|
return value
|
|
|
|
def __setitem__(self, key, value):
|
|
"""
|
|
设置缓存项
|
|
"""
|
|
kwargs = {'region': self._region}
|
|
self._cache_backend.set(key, value, **kwargs)
|
|
|
|
def __delitem__(self, key):
|
|
"""
|
|
删除缓存项
|
|
"""
|
|
if not self._cache_backend.exists(key, region=self._region):
|
|
raise KeyError(key)
|
|
self._cache_backend.delete(key, region=self._region)
|
|
|
|
def __contains__(self, key):
|
|
"""
|
|
检查键是否存在
|
|
"""
|
|
return self._cache_backend.exists(key, region=self._region)
|
|
|
|
def __iter__(self):
|
|
"""
|
|
返回缓存的迭代器
|
|
"""
|
|
for key, _ in self._cache_backend.items(region=self._region):
|
|
yield key
|
|
|
|
def __len__(self):
|
|
"""
|
|
返回缓存项的数量
|
|
"""
|
|
return sum(1 for _ in self._cache_backend.items(region=self._region))
|
|
|
|
def is_redis(self) -> bool:
|
|
"""
|
|
检查当前缓存后端是否为 Redis
|
|
"""
|
|
return self._cache_backend.is_redis()
|
|
|
|
def get(self, key: str, **kwargs) -> Any:
|
|
"""
|
|
获取缓存值
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.get(key, **kwargs)
|
|
|
|
def set(self, key: str, value: Any, **kwargs) -> None:
|
|
"""
|
|
设置缓存值
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
self._cache_backend.set(key, value, **kwargs)
|
|
|
|
def delete(self, key: str, **kwargs) -> None:
|
|
"""
|
|
删除缓存值
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
self._cache_backend.delete(key, **kwargs)
|
|
|
|
def exists(self, key: str, **kwargs) -> bool:
|
|
"""
|
|
检查缓存键是否存在
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.exists(key, **kwargs)
|
|
|
|
def clear(self, **kwargs) -> None:
|
|
"""
|
|
清除缓存
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
self._cache_backend.clear(**kwargs)
|
|
|
|
def items(self, **kwargs):
|
|
"""
|
|
获取所有缓存项
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.items(**kwargs)
|
|
|
|
def keys(self, **kwargs):
|
|
"""
|
|
获取所有缓存键
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.keys(**kwargs)
|
|
|
|
def values(self, **kwargs):
|
|
"""
|
|
获取所有缓存值
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.values(**kwargs)
|
|
|
|
def update(self, other: Dict[str, Any], **kwargs) -> None:
|
|
"""
|
|
更新缓存
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
self._cache_backend.update(other, **kwargs)
|
|
|
|
def pop(self, key: str, default: Any = None, **kwargs) -> Any:
|
|
"""
|
|
弹出缓存项
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.pop(key, default, **kwargs)
|
|
|
|
def popitem(self, **kwargs) -> Tuple[str, Any]:
|
|
"""
|
|
弹出最后一个缓存项
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.popitem(**kwargs)
|
|
|
|
def setdefault(self, key: str, default: Any = None, **kwargs) -> Any:
|
|
"""
|
|
设置默认值
|
|
"""
|
|
kwargs.setdefault('region', self._region)
|
|
return self._cache_backend.setdefault(key, default, **kwargs)
|
|
|
|
def close(self) -> None:
|
|
"""
|
|
关闭缓存连接
|
|
"""
|
|
self._cache_backend.close()
|
|
|
|
|
|
class TTLCache(CacheProxy):
|
|
"""
|
|
基于 TTL 的缓存类,兼容 cachetools.TTLCache 接口
|
|
使用项目的缓存后端实现,支持 Redis 和内存缓存
|
|
"""
|
|
|
|
def __init__(self,
|
|
region: Optional[str] = DEFAULT_CACHE_REGION,
|
|
maxsize: Optional[int] = DEFAULT_CACHE_SIZE,
|
|
ttl: Optional[int] = DEFAULT_CACHE_TTL):
|
|
"""
|
|
初始化 TTL 缓存
|
|
|
|
:param maxsize: 缓存的最大条目数
|
|
:param ttl: 缓存的存活时间,单位秒
|
|
:param region: 缓存的区,为 None 时使用默认区
|
|
"""
|
|
super().__init__(Cache(cache_type='ttl', maxsize=maxsize, ttl=ttl), region)
|
|
|
|
|
|
class LRUCache(CacheProxy):
|
|
"""
|
|
基于 LRU 的缓存类,兼容 cachetools.LRUCache 接口
|
|
使用项目的缓存后端实现,支持 Redis 和内存缓存
|
|
"""
|
|
|
|
def __init__(self,
|
|
region: Optional[str] = DEFAULT_CACHE_REGION,
|
|
maxsize: Optional[int] = DEFAULT_CACHE_SIZE
|
|
):
|
|
"""
|
|
初始化 LRU 缓存
|
|
|
|
:param maxsize: 缓存的最大条目数
|
|
:param region: 缓存的区,为 None 时使用默认区
|
|
"""
|
|
super().__init__(Cache(cache_type='lru', maxsize=maxsize), region)
|