mirror of
https://github.com/checktheroads/hyperglass
synced 2024-05-11 05:55:08 +00:00
Redis caching improvements
This commit is contained in:
112
hyperglass/cache/sync.py
vendored
Normal file
112
hyperglass/cache/sync.py
vendored
Normal file
@@ -0,0 +1,112 @@
|
||||
"""Non-asyncio Redis cache handler."""
|
||||
|
||||
# Standard Library
|
||||
import json
|
||||
import time
|
||||
import pickle
|
||||
from typing import Any, Dict
|
||||
|
||||
# Third Party
|
||||
from redis import Redis as SyncRedis
|
||||
from redis.client import PubSub as SyncPubsSub
|
||||
|
||||
# Project
|
||||
from hyperglass.cache.base import BaseCache
|
||||
|
||||
|
||||
class SyncCache(BaseCache):
|
||||
"""Synchronous Redis cache handler."""
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
"""Initialize Redis connection."""
|
||||
super().__init__(*args, **kwargs)
|
||||
self.instance: SyncRedis = SyncRedis(
|
||||
db=self.db,
|
||||
host=self.host,
|
||||
port=self.port,
|
||||
decode_responses=self.decode_responses,
|
||||
**self.redis_args,
|
||||
)
|
||||
|
||||
def get(self, *args: str) -> Any:
|
||||
"""Get item(s) from cache."""
|
||||
if len(args) == 1:
|
||||
raw = self.instance.get(args[0])
|
||||
else:
|
||||
raw = self.instance.mget(args)
|
||||
return self.parse_types(raw)
|
||||
|
||||
def get_dict(self, key: str, field: str = "") -> Any:
|
||||
"""Get hash map (dict) item(s)."""
|
||||
if not field:
|
||||
raw = self.instance.hgetall(key)
|
||||
else:
|
||||
raw = self.instance.hget(key, str(field))
|
||||
|
||||
return self.parse_types(raw)
|
||||
|
||||
def set(self, key: str, value: str) -> bool:
|
||||
"""Set cache values."""
|
||||
return self.instance.set(key, str(value))
|
||||
|
||||
def set_dict(self, key: str, field: str, value: str) -> bool:
|
||||
"""Set hash map (dict) values."""
|
||||
success = False
|
||||
|
||||
if isinstance(value, Dict):
|
||||
value = json.dumps(value)
|
||||
else:
|
||||
value = str(value)
|
||||
|
||||
response = self.instance.hset(key, str(field), value)
|
||||
|
||||
if response in (0, 1):
|
||||
success = True
|
||||
|
||||
return success
|
||||
|
||||
def wait(self, pubsub: SyncPubsSub, timeout: int = 30, **kwargs) -> Any:
|
||||
"""Wait for pub/sub messages & return posted message."""
|
||||
now = time.time()
|
||||
timeout = now + timeout
|
||||
|
||||
while now < timeout:
|
||||
|
||||
message = pubsub.get_message(ignore_subscribe_messages=True, **kwargs)
|
||||
|
||||
if message is not None and message["type"] == "message":
|
||||
data = message["data"]
|
||||
return self.parse_types(data)
|
||||
|
||||
time.sleep(0.01)
|
||||
now = time.time()
|
||||
|
||||
return None
|
||||
|
||||
def pubsub(self) -> SyncPubsSub:
|
||||
"""Provide a redis.client.Pubsub instance."""
|
||||
return self.instance.pubsub()
|
||||
|
||||
def pub(self, key: str, value: str) -> None:
|
||||
"""Publish a value."""
|
||||
time.sleep(1)
|
||||
self.instance.publish(key, value)
|
||||
|
||||
def clear(self) -> None:
|
||||
"""Clear the cache."""
|
||||
self.instance.flushdb()
|
||||
|
||||
def delete(self, *keys: str) -> None:
|
||||
"""Delete a cache key."""
|
||||
self.instance.delete(*keys)
|
||||
|
||||
def expire(self, *keys: str, seconds: int) -> None:
|
||||
"""Set timeout of key in seconds."""
|
||||
for key in keys:
|
||||
self.instance.expire(key, seconds)
|
||||
|
||||
def get_config(self) -> Dict:
|
||||
"""Get picked config object from cache."""
|
||||
|
||||
pickled = self.instance.get("HYPERGLASS_CONFIG")
|
||||
return pickle.loads(pickled)
|
Reference in New Issue
Block a user