Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/redis_utils.py: 94%
29 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
1from typing import Generator
3import redis.asyncio as redis
4from opal_common.logger import logger
5from pydantic import BaseModel
8class RedisDB:
9 """Small utility class to persist objects in Redis."""
11 def __init__(self, redis_url):
12 self._url = redis_url
13 logger.debug("Connecting to Redis: {url}", url=self._url)
15 self._redis = redis.Redis.from_url(self._url)
17 @property
18 def redis_connection(self) -> redis.Redis:
19 return self._redis
21 async def set(self, key: str, value: BaseModel):
22 await self._redis.set(key, self._serialize(value))
24 async def set_if_not_exists(self, key: str, value: BaseModel) -> bool:
25 """:param key:
26 :param value:
27 :return: True if created, False if key already exists
28 """
30 return await self._redis.set(key, self._serialize(value), nx=True)
32 async def get(self, key: str) -> bytes:
33 return await self._redis.get(key)
35 async def scan(self, pattern: str) -> Generator[bytes, None, None]:
36 cur = b"0"
37 while cur:
38 cur, keys = await self._redis.scan(cur, match=pattern)
40 for key in keys:
41 value = await self._redis.get(key)
42 yield value
44 async def delete(self, key: str):
45 await self._redis.delete(key)
47 def _serialize(self, value: BaseModel) -> str:
48 return value.json()