Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/redis/cache.py: 52%
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
1from abc import ABC, abstractmethod
2from collections import OrderedDict
3from dataclasses import dataclass
4from enum import Enum
5from typing import Any, List, Optional, Union
7from redis.commands.metadata import MetadataResolver, StaticMetadataResolver
8from redis.observability.attributes import CSCReason
11class CacheEntryStatus(Enum):
12 VALID = "VALID"
13 IN_PROGRESS = "IN_PROGRESS"
16class EvictionPolicyType(Enum):
17 time_based = "time_based"
18 frequency_based = "frequency_based"
21@dataclass(frozen=True)
22class CacheKey:
23 """
24 Represents a unique key for a cache entry.
26 Attributes:
27 command (str): The Redis command being cached.
28 redis_keys (tuple): The Redis keys involved in the command.
29 redis_args (tuple): Additional arguments for the Redis command.
30 This field is included in the cache key to ensure uniqueness
31 when commands have the same keys but different arguments.
32 Changing this field will affect cache key uniqueness.
33 """
35 command: str
36 redis_keys: tuple
37 redis_args: tuple = () # Additional arguments for the Redis command; affects cache key uniqueness.
40class CacheEntry:
41 def __init__(
42 self,
43 cache_key: CacheKey,
44 cache_value: bytes,
45 status: CacheEntryStatus,
46 connection_ref,
47 ):
48 self.cache_key = cache_key
49 self.cache_value = cache_value
50 self.status = status
51 self.connection_ref = connection_ref
53 def __hash__(self):
54 return hash(
55 (self.cache_key, self.cache_value, self.status, self.connection_ref)
56 )
58 def __eq__(self, other):
59 return hash(self) == hash(other)
62class EvictionPolicyInterface(ABC):
63 @property
64 @abstractmethod
65 def cache(self):
66 pass
68 @cache.setter
69 @abstractmethod
70 def cache(self, value):
71 pass
73 @property
74 @abstractmethod
75 def type(self) -> EvictionPolicyType:
76 pass
78 @abstractmethod
79 def evict_next(self) -> CacheKey:
80 pass
82 @abstractmethod
83 def evict_many(self, count: int) -> List[CacheKey]:
84 pass
86 @abstractmethod
87 def touch(self, cache_key: CacheKey) -> None:
88 pass
91class CacheConfigurationInterface(ABC):
92 @abstractmethod
93 def get_cache_class(self):
94 pass
96 @abstractmethod
97 def get_max_size(self) -> int:
98 pass
100 @abstractmethod
101 def get_eviction_policy(self):
102 pass
104 @abstractmethod
105 def is_exceeds_max_size(self, count: int) -> bool:
106 pass
108 @abstractmethod
109 def is_allowed_to_cache(self, command: str) -> bool:
110 pass
113class CacheInterface(ABC):
114 @property
115 @abstractmethod
116 def collection(self) -> OrderedDict:
117 pass
119 @property
120 @abstractmethod
121 def config(self) -> CacheConfigurationInterface:
122 pass
124 @property
125 @abstractmethod
126 def eviction_policy(self) -> EvictionPolicyInterface:
127 pass
129 @property
130 @abstractmethod
131 def size(self) -> int:
132 pass
134 @abstractmethod
135 def get(self, key: CacheKey) -> Union[CacheEntry, None]:
136 pass
138 @abstractmethod
139 def set(self, entry: CacheEntry) -> bool:
140 pass
142 @abstractmethod
143 def delete_by_cache_keys(self, cache_keys: List[CacheKey]) -> List[bool]:
144 pass
146 @abstractmethod
147 def delete_by_redis_keys(self, redis_keys: List[bytes]) -> List[bool]:
148 pass
150 @abstractmethod
151 def flush(self) -> int:
152 pass
154 @abstractmethod
155 def is_cachable(self, key: CacheKey) -> bool:
156 pass
159class DefaultCache(CacheInterface):
160 def __init__(
161 self,
162 cache_config: CacheConfigurationInterface,
163 ) -> None:
164 self._cache = OrderedDict()
165 self._cache_config = cache_config
166 self._eviction_policy = self._cache_config.get_eviction_policy().value()
167 self._eviction_policy.cache = self
169 @property
170 def collection(self) -> OrderedDict:
171 return self._cache
173 @property
174 def config(self) -> CacheConfigurationInterface:
175 return self._cache_config
177 @property
178 def eviction_policy(self) -> EvictionPolicyInterface:
179 return self._eviction_policy
181 @property
182 def size(self) -> int:
183 return len(self._cache)
185 def set(self, entry: CacheEntry) -> bool:
186 if not self.is_cachable(entry.cache_key):
187 return False
189 self._cache[entry.cache_key] = entry
190 self._eviction_policy.touch(entry.cache_key)
192 return True
194 def get(self, key: CacheKey) -> Union[CacheEntry, None]:
195 entry = self._cache.get(key, None)
197 if entry is None:
198 return None
200 self._eviction_policy.touch(key)
201 return entry
203 def delete_by_cache_keys(self, cache_keys: List[CacheKey]) -> List[bool]:
204 response = []
206 for key in cache_keys:
207 if self.get(key) is not None:
208 self._cache.pop(key)
209 response.append(True)
210 else:
211 response.append(False)
213 return response
215 def delete_by_redis_keys(
216 self, redis_keys: Union[List[bytes], List[str]]
217 ) -> List[bool]:
218 response = []
219 keys_to_delete = []
221 for redis_key in redis_keys:
222 # Prepare both versions for lookup
223 candidates = [redis_key]
224 if isinstance(redis_key, str):
225 candidates.append(redis_key.encode("utf-8"))
226 elif isinstance(redis_key, bytes):
227 try:
228 candidates.append(redis_key.decode("utf-8"))
229 except UnicodeDecodeError:
230 pass # Non-UTF-8 bytes, skip str version
232 for cache_key in self._cache:
233 if any(candidate in cache_key.redis_keys for candidate in candidates):
234 keys_to_delete.append(cache_key)
235 response.append(True)
237 for key in keys_to_delete:
238 self._cache.pop(key)
240 return response
242 def flush(self) -> int:
243 elem_count = len(self._cache)
244 self._cache.clear()
245 return elem_count
247 def is_cachable(self, key: CacheKey) -> bool:
248 return self._cache_config.is_allowed_to_cache(key.command)
251class CacheProxy(CacheInterface):
252 """
253 Proxy object that wraps cache implementations to enable additional logic on top
254 """
256 def __init__(self, cache: CacheInterface):
257 self._cache = cache
259 @property
260 def collection(self) -> OrderedDict:
261 return self._cache.collection
263 @property
264 def config(self) -> CacheConfigurationInterface:
265 return self._cache.config
267 @property
268 def eviction_policy(self) -> EvictionPolicyInterface:
269 return self._cache.eviction_policy
271 @property
272 def size(self) -> int:
273 return self._cache.size
275 def get(self, key: CacheKey) -> Union[CacheEntry, None]:
276 return self._cache.get(key)
278 def set(self, entry: CacheEntry) -> bool:
279 is_set = self._cache.set(entry)
281 if self.config.is_exceeds_max_size(self.size):
282 # Lazy import to avoid circular dependency
283 from redis.observability.recorder import record_csc_eviction
285 record_csc_eviction(
286 count=1,
287 reason=CSCReason.FULL,
288 )
289 self.eviction_policy.evict_next()
291 return is_set
293 def delete_by_cache_keys(self, cache_keys: List[CacheKey]) -> List[bool]:
294 return self._cache.delete_by_cache_keys(cache_keys)
296 def delete_by_redis_keys(self, redis_keys: List[bytes]) -> List[bool]:
297 return self._cache.delete_by_redis_keys(redis_keys)
299 def flush(self) -> int:
300 return self._cache.flush()
302 def is_cachable(self, key: CacheKey) -> bool:
303 return self._cache.is_cachable(key)
306class LRUPolicy(EvictionPolicyInterface):
307 def __init__(self):
308 self.cache = None
310 @property
311 def cache(self):
312 return self._cache
314 @cache.setter
315 def cache(self, cache: CacheInterface):
316 self._cache = cache
318 @property
319 def type(self) -> EvictionPolicyType:
320 return EvictionPolicyType.time_based
322 def evict_next(self) -> CacheKey:
323 self._assert_cache()
324 popped_entry = self._cache.collection.popitem(last=False)
325 return popped_entry[0]
327 def evict_many(self, count: int) -> List[CacheKey]:
328 self._assert_cache()
329 if count > len(self._cache.collection):
330 raise ValueError("Evictions count is above cache size")
332 popped_keys = []
334 for _ in range(count):
335 popped_entry = self._cache.collection.popitem(last=False)
336 popped_keys.append(popped_entry[0])
338 return popped_keys
340 def touch(self, cache_key: CacheKey) -> None:
341 self._assert_cache()
343 if self._cache.collection.get(cache_key) is None:
344 raise ValueError("Given entry does not belong to the cache")
346 self._cache.collection.move_to_end(cache_key)
348 def _assert_cache(self):
349 if self.cache is None or not isinstance(self.cache, CacheInterface):
350 raise ValueError("Eviction policy should be associated with valid cache.")
353class EvictionPolicy(Enum):
354 LRU = LRUPolicy
357class CacheConfig(CacheConfigurationInterface):
358 DEFAULT_CACHE_CLASS = DefaultCache
359 DEFAULT_EVICTION_POLICY = EvictionPolicy.LRU
360 DEFAULT_MAX_SIZE = 10000
362 # DEPRECATED - no longer consulted, and it will be removed in a future release.
363 #
364 # Command eligibility is now decided from command metadata by the metadata resolver this
365 # config holds, so this list no longer describes what gets cached. It is kept as a public
366 # attribute only so an external caller reading it keeps working; editing it changes
367 # nothing. The effective set it is replaced by differs from it by ``+FT.SUGGET``,
368 # ``+FT.SUGLEN``, ``+DIGEST``, ``+EXPIRETIME``, ``+PEXPIRETIME``, ``+HEXPIRETIME``,
369 # ``+HPEXPIRETIME``, ``+SDIFFCARD``, ``+SUNIONCARD`` and ``-XPENDING``, ``-TS.INFO``,
370 # ``-XREAD`` - the three removals being commands the server itself reports as not cacheable.
371 #
372 # To change eligibility, edit ``redis.commands.metadata._STATIC_COMMAND_METADATA`` or pass
373 # a ``metadata_resolver`` to the client. Nothing here.
374 DEFAULT_ALLOW_LIST = [
375 "BITCOUNT",
376 "BITFIELD_RO",
377 "BITPOS",
378 "EXISTS",
379 "GEODIST",
380 "GEOHASH",
381 "GEOPOS",
382 "GEORADIUSBYMEMBER_RO",
383 "GEORADIUS_RO",
384 "GEOSEARCH",
385 "GET",
386 "GETBIT",
387 "GETRANGE",
388 "HEXISTS",
389 "HGET",
390 "HGETALL",
391 "HKEYS",
392 "HLEN",
393 "HMGET",
394 "HSTRLEN",
395 "HVALS",
396 "JSON.ARRINDEX",
397 "JSON.ARRLEN",
398 "JSON.GET",
399 "JSON.MGET",
400 "JSON.OBJKEYS",
401 "JSON.OBJLEN",
402 "JSON.RESP",
403 "JSON.STRLEN",
404 "JSON.TYPE",
405 "LCS",
406 "LINDEX",
407 "LLEN",
408 "LPOS",
409 "LRANGE",
410 "MGET",
411 "SCARD",
412 "SDIFF",
413 "SINTER",
414 "SINTERCARD",
415 "SISMEMBER",
416 "SMEMBERS",
417 "SMISMEMBER",
418 "SORT_RO",
419 "STRLEN",
420 "SUBSTR",
421 "SUNION",
422 "TS.GET",
423 "TS.INFO",
424 "TS.RANGE",
425 "TS.REVRANGE",
426 "TYPE",
427 "XLEN",
428 "XPENDING",
429 "XRANGE",
430 "XREAD",
431 "XREVRANGE",
432 "ZCARD",
433 "ZCOUNT",
434 "ZDIFF",
435 "ZINTER",
436 "ZINTERCARD",
437 "ZLEXCOUNT",
438 "ZMSCORE",
439 "ZRANGE",
440 "ZRANGEBYLEX",
441 "ZRANGEBYSCORE",
442 "ZRANK",
443 "ZREVRANGE",
444 "ZREVRANGEBYLEX",
445 "ZREVRANGEBYSCORE",
446 "ZREVRANK",
447 "ZSCORE",
448 "ZUNION",
449 ]
451 def __init__(
452 self,
453 max_size: int = DEFAULT_MAX_SIZE,
454 cache_class: Any = DEFAULT_CACHE_CLASS,
455 eviction_policy: EvictionPolicy = DEFAULT_EVICTION_POLICY,
456 ):
457 self._cache_class = cache_class
458 self._max_size = max_size
459 self._eviction_policy = eviction_policy
460 # Defaulted here rather than taken as a constructor argument: eligibility is
461 # configured at client level, through the ``metadata_resolver`` of the client or the
462 # pool, which injects it below. Defaulting it means a config built standalone - in a
463 # test, or by a user who configures nothing else - is fully functional, and decides
464 # eligibility from the command metadata this library ships.
465 self._metadata_resolver: MetadataResolver = StaticMetadataResolver()
467 def set_metadata_resolver(self, metadata_resolver: MetadataResolver) -> None:
468 """
469 Set the metadata resolver that decides which commands may be cached.
471 Called by the connection pool with the client-level resolver, so that one object
472 serves both cluster routing and cache eligibility. Deliberately not part of
473 :class:`CacheConfigurationInterface`: that ABC is public and implemented by third
474 parties, so a custom configuration keeps whatever eligibility logic it has.
476 Which object the pool calls this on depends on how the cache was supplied, and the
477 difference is observable to a caller who reuses one configuration:
479 - Given ``cache_config=``, the pool copies the configuration first, so the caller's
480 object keeps the resolver it had and two clients sharing it stay independent. A
481 later call to this method on the caller's object does not reach a pool already
482 built from it.
483 - Given ``cache=`` or ``cache_factory=``, the caller supplied a whole cache that
484 reads its configuration on every lookup and cannot be handed a different one, so
485 the pool calls this on the configuration inside it. One configuration reused that
486 way therefore ends up with whichever resolver was injected last.
488 Args:
489 metadata_resolver: The resolver to decide eligibility through.
490 """
491 self._metadata_resolver = metadata_resolver
493 def get_cache_class(self):
494 return self._cache_class
496 def get_max_size(self) -> int:
497 return self._max_size
499 def get_eviction_policy(self) -> EvictionPolicy:
500 return self._eviction_policy
502 def is_exceeds_max_size(self, count: int) -> bool:
503 return count > self._max_size
505 def is_allowed_to_cache(self, command: str) -> bool:
506 # Fails closed on everything the resolver cannot decide: an unknown command, a name
507 # the record tables cannot be keyed by, and a record built from incomplete metadata
508 # all resolve to False. The verdict is memoized per command name, so this is a dict
509 # hit on the command execution path.
510 return self._metadata_resolver.is_cacheable(command)
513class CacheFactoryInterface(ABC):
514 @abstractmethod
515 def get_cache(self) -> CacheInterface:
516 pass
519class CacheFactory(CacheFactoryInterface):
520 def __init__(self, cache_config: Optional[CacheConfig] = None):
521 self._config = cache_config
523 if self._config is None:
524 self._config = CacheConfig()
526 def get_cache(self) -> CacheInterface:
527 cache_class = self._config.get_cache_class()
528 return CacheProxy(cache_class(cache_config=self._config))