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

265 statements  

1from abc import ABC, abstractmethod 

2from collections import OrderedDict 

3from dataclasses import dataclass 

4from enum import Enum 

5from typing import Any, List, Optional, Union 

6 

7from redis.commands.metadata import MetadataResolver, StaticMetadataResolver 

8from redis.observability.attributes import CSCReason 

9 

10 

11class CacheEntryStatus(Enum): 

12 VALID = "VALID" 

13 IN_PROGRESS = "IN_PROGRESS" 

14 

15 

16class EvictionPolicyType(Enum): 

17 time_based = "time_based" 

18 frequency_based = "frequency_based" 

19 

20 

21@dataclass(frozen=True) 

22class CacheKey: 

23 """ 

24 Represents a unique key for a cache entry. 

25 

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 """ 

34 

35 command: str 

36 redis_keys: tuple 

37 redis_args: tuple = () # Additional arguments for the Redis command; affects cache key uniqueness. 

38 

39 

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 

52 

53 def __hash__(self): 

54 return hash( 

55 (self.cache_key, self.cache_value, self.status, self.connection_ref) 

56 ) 

57 

58 def __eq__(self, other): 

59 return hash(self) == hash(other) 

60 

61 

62class EvictionPolicyInterface(ABC): 

63 @property 

64 @abstractmethod 

65 def cache(self): 

66 pass 

67 

68 @cache.setter 

69 @abstractmethod 

70 def cache(self, value): 

71 pass 

72 

73 @property 

74 @abstractmethod 

75 def type(self) -> EvictionPolicyType: 

76 pass 

77 

78 @abstractmethod 

79 def evict_next(self) -> CacheKey: 

80 pass 

81 

82 @abstractmethod 

83 def evict_many(self, count: int) -> List[CacheKey]: 

84 pass 

85 

86 @abstractmethod 

87 def touch(self, cache_key: CacheKey) -> None: 

88 pass 

89 

90 

91class CacheConfigurationInterface(ABC): 

92 @abstractmethod 

93 def get_cache_class(self): 

94 pass 

95 

96 @abstractmethod 

97 def get_max_size(self) -> int: 

98 pass 

99 

100 @abstractmethod 

101 def get_eviction_policy(self): 

102 pass 

103 

104 @abstractmethod 

105 def is_exceeds_max_size(self, count: int) -> bool: 

106 pass 

107 

108 @abstractmethod 

109 def is_allowed_to_cache(self, command: str) -> bool: 

110 pass 

111 

112 

113class CacheInterface(ABC): 

114 @property 

115 @abstractmethod 

116 def collection(self) -> OrderedDict: 

117 pass 

118 

119 @property 

120 @abstractmethod 

121 def config(self) -> CacheConfigurationInterface: 

122 pass 

123 

124 @property 

125 @abstractmethod 

126 def eviction_policy(self) -> EvictionPolicyInterface: 

127 pass 

128 

129 @property 

130 @abstractmethod 

131 def size(self) -> int: 

132 pass 

133 

134 @abstractmethod 

135 def get(self, key: CacheKey) -> Union[CacheEntry, None]: 

136 pass 

137 

138 @abstractmethod 

139 def set(self, entry: CacheEntry) -> bool: 

140 pass 

141 

142 @abstractmethod 

143 def delete_by_cache_keys(self, cache_keys: List[CacheKey]) -> List[bool]: 

144 pass 

145 

146 @abstractmethod 

147 def delete_by_redis_keys(self, redis_keys: List[bytes]) -> List[bool]: 

148 pass 

149 

150 @abstractmethod 

151 def flush(self) -> int: 

152 pass 

153 

154 @abstractmethod 

155 def is_cachable(self, key: CacheKey) -> bool: 

156 pass 

157 

158 

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 

168 

169 @property 

170 def collection(self) -> OrderedDict: 

171 return self._cache 

172 

173 @property 

174 def config(self) -> CacheConfigurationInterface: 

175 return self._cache_config 

176 

177 @property 

178 def eviction_policy(self) -> EvictionPolicyInterface: 

179 return self._eviction_policy 

180 

181 @property 

182 def size(self) -> int: 

183 return len(self._cache) 

184 

185 def set(self, entry: CacheEntry) -> bool: 

186 if not self.is_cachable(entry.cache_key): 

187 return False 

188 

189 self._cache[entry.cache_key] = entry 

190 self._eviction_policy.touch(entry.cache_key) 

191 

192 return True 

193 

194 def get(self, key: CacheKey) -> Union[CacheEntry, None]: 

195 entry = self._cache.get(key, None) 

196 

197 if entry is None: 

198 return None 

199 

200 self._eviction_policy.touch(key) 

201 return entry 

202 

203 def delete_by_cache_keys(self, cache_keys: List[CacheKey]) -> List[bool]: 

204 response = [] 

205 

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) 

212 

213 return response 

214 

215 def delete_by_redis_keys( 

216 self, redis_keys: Union[List[bytes], List[str]] 

217 ) -> List[bool]: 

218 response = [] 

219 keys_to_delete = [] 

220 

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 

231 

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) 

236 

237 for key in keys_to_delete: 

238 self._cache.pop(key) 

239 

240 return response 

241 

242 def flush(self) -> int: 

243 elem_count = len(self._cache) 

244 self._cache.clear() 

245 return elem_count 

246 

247 def is_cachable(self, key: CacheKey) -> bool: 

248 return self._cache_config.is_allowed_to_cache(key.command) 

249 

250 

251class CacheProxy(CacheInterface): 

252 """ 

253 Proxy object that wraps cache implementations to enable additional logic on top 

254 """ 

255 

256 def __init__(self, cache: CacheInterface): 

257 self._cache = cache 

258 

259 @property 

260 def collection(self) -> OrderedDict: 

261 return self._cache.collection 

262 

263 @property 

264 def config(self) -> CacheConfigurationInterface: 

265 return self._cache.config 

266 

267 @property 

268 def eviction_policy(self) -> EvictionPolicyInterface: 

269 return self._cache.eviction_policy 

270 

271 @property 

272 def size(self) -> int: 

273 return self._cache.size 

274 

275 def get(self, key: CacheKey) -> Union[CacheEntry, None]: 

276 return self._cache.get(key) 

277 

278 def set(self, entry: CacheEntry) -> bool: 

279 is_set = self._cache.set(entry) 

280 

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 

284 

285 record_csc_eviction( 

286 count=1, 

287 reason=CSCReason.FULL, 

288 ) 

289 self.eviction_policy.evict_next() 

290 

291 return is_set 

292 

293 def delete_by_cache_keys(self, cache_keys: List[CacheKey]) -> List[bool]: 

294 return self._cache.delete_by_cache_keys(cache_keys) 

295 

296 def delete_by_redis_keys(self, redis_keys: List[bytes]) -> List[bool]: 

297 return self._cache.delete_by_redis_keys(redis_keys) 

298 

299 def flush(self) -> int: 

300 return self._cache.flush() 

301 

302 def is_cachable(self, key: CacheKey) -> bool: 

303 return self._cache.is_cachable(key) 

304 

305 

306class LRUPolicy(EvictionPolicyInterface): 

307 def __init__(self): 

308 self.cache = None 

309 

310 @property 

311 def cache(self): 

312 return self._cache 

313 

314 @cache.setter 

315 def cache(self, cache: CacheInterface): 

316 self._cache = cache 

317 

318 @property 

319 def type(self) -> EvictionPolicyType: 

320 return EvictionPolicyType.time_based 

321 

322 def evict_next(self) -> CacheKey: 

323 self._assert_cache() 

324 popped_entry = self._cache.collection.popitem(last=False) 

325 return popped_entry[0] 

326 

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") 

331 

332 popped_keys = [] 

333 

334 for _ in range(count): 

335 popped_entry = self._cache.collection.popitem(last=False) 

336 popped_keys.append(popped_entry[0]) 

337 

338 return popped_keys 

339 

340 def touch(self, cache_key: CacheKey) -> None: 

341 self._assert_cache() 

342 

343 if self._cache.collection.get(cache_key) is None: 

344 raise ValueError("Given entry does not belong to the cache") 

345 

346 self._cache.collection.move_to_end(cache_key) 

347 

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.") 

351 

352 

353class EvictionPolicy(Enum): 

354 LRU = LRUPolicy 

355 

356 

357class CacheConfig(CacheConfigurationInterface): 

358 DEFAULT_CACHE_CLASS = DefaultCache 

359 DEFAULT_EVICTION_POLICY = EvictionPolicy.LRU 

360 DEFAULT_MAX_SIZE = 10000 

361 

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 ] 

450 

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() 

466 

467 def set_metadata_resolver(self, metadata_resolver: MetadataResolver) -> None: 

468 """ 

469 Set the metadata resolver that decides which commands may be cached. 

470 

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. 

475 

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: 

478 

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. 

487 

488 Args: 

489 metadata_resolver: The resolver to decide eligibility through. 

490 """ 

491 self._metadata_resolver = metadata_resolver 

492 

493 def get_cache_class(self): 

494 return self._cache_class 

495 

496 def get_max_size(self) -> int: 

497 return self._max_size 

498 

499 def get_eviction_policy(self) -> EvictionPolicy: 

500 return self._eviction_policy 

501 

502 def is_exceeds_max_size(self, count: int) -> bool: 

503 return count > self._max_size 

504 

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) 

511 

512 

513class CacheFactoryInterface(ABC): 

514 @abstractmethod 

515 def get_cache(self) -> CacheInterface: 

516 pass 

517 

518 

519class CacheFactory(CacheFactoryInterface): 

520 def __init__(self, cache_config: Optional[CacheConfig] = None): 

521 self._config = cache_config 

522 

523 if self._config is None: 

524 self._config = CacheConfig() 

525 

526 def get_cache(self) -> CacheInterface: 

527 cache_class = self._config.get_cache_class() 

528 return CacheProxy(cache_class(cache_config=self._config))