From 34ac959f5a4b8d50c7c58af9b9042c56c69cb1e7 Mon Sep 17 00:00:00 2001 From: Gabriel Luiz Freitas Almeida Date: Sun, 13 Aug 2023 23:37:22 -0300 Subject: [PATCH] =?UTF-8?q?=F0=9F=94=84=20refactor(cache):=20reorganize=20?= =?UTF-8?q?imports=20and=20remove=20unused=20imports=20and=20variables=20i?= =?UTF-8?q?n=20cache=20module=20=F0=9F=94=84=20refactor(cache):=20rename?= =?UTF-8?q?=20BaseCache=20class=20to=20BaseCacheManager=20for=20better=20c?= =?UTF-8?q?larity=20and=20consistency=20=F0=9F=94=84=20refactor(cache):=20?= =?UTF-8?q?remove=20cache=5Fmanager=20from=20=5F=5Fall=5F=5F=20in=20cache?= =?UTF-8?q?=20module=20=F0=9F=94=84=20refactor(cache):=20remove=20unused?= =?UTF-8?q?=20import=20of=20Service=20in=20BaseCacheManager=20class=20?= =?UTF-8?q?=F0=9F=94=84=20refactor(cache):=20rename=20cache=5Fmanager=20to?= =?UTF-8?q?=20BaseCacheManager=20in=20factory=20module=20=F0=9F=94=84=20re?= =?UTF-8?q?factor(cache):=20remove=20unused=20import=20of=20InMemoryCache?= =?UTF-8?q?=20in=20factory=20module=20=F0=9F=94=84=20refactor(cache):=20re?= =?UTF-8?q?move=20flow=20module=20and=20its=20related=20code=20as=20it=20i?= =?UTF-8?q?s=20no=20longer=20used?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 🔧 fix(manager.py): fix import statements and class inheritance in cache manager module ✨ feat(manager.py): add InMemoryCache class to implement a simple in-memory cache using an OrderedDict ✨ feat(manager.py): add RedisCache class to implement a Redis-based cache using the redis-py library --- .../langflow/services/cache/__init__.py | 4 +- src/backend/langflow/services/cache/base.py | 6 +- .../langflow/services/cache/factory.py | 26 +- src/backend/langflow/services/cache/flow.py | 146 -------- .../langflow/services/cache/manager.py | 351 +++++++++++------- 5 files changed, 256 insertions(+), 277 deletions(-) delete mode 100644 src/backend/langflow/services/cache/flow.py diff --git a/src/backend/langflow/services/cache/__init__.py b/src/backend/langflow/services/cache/__init__.py index 79e143807..3b122aa9e 100644 --- a/src/backend/langflow/services/cache/__init__.py +++ b/src/backend/langflow/services/cache/__init__.py @@ -1,10 +1,8 @@ from . import factory, manager -from langflow.services.cache.manager import cache_manager -from langflow.services.cache.flow import InMemoryCache +from langflow.services.cache.manager import InMemoryCache __all__ = [ - "cache_manager", "factory", "manager", "InMemoryCache", diff --git a/src/backend/langflow/services/cache/base.py b/src/backend/langflow/services/cache/base.py index 88cb3a1da..80aec1408 100644 --- a/src/backend/langflow/services/cache/base.py +++ b/src/backend/langflow/services/cache/base.py @@ -1,11 +1,15 @@ import abc +from langflow.services.base import Service -class BaseCache(abc.ABC): + +class BaseCacheManager(abc.ABC, Service): """ Abstract base class for a cache. """ + name = "cache_manager" + @abc.abstractmethod def get(self, key): """ diff --git a/src/backend/langflow/services/cache/factory.py b/src/backend/langflow/services/cache/factory.py index 77f8d58d1..78e307bfa 100644 --- a/src/backend/langflow/services/cache/factory.py +++ b/src/backend/langflow/services/cache/factory.py @@ -1,11 +1,31 @@ -from langflow.services.cache.manager import CacheManager +from langflow.services.cache.manager import InMemoryCache, RedisCache, BaseCacheManager from langflow.services.factory import ServiceFactory +from langflow.utils.logger import logger class CacheManagerFactory(ServiceFactory): def __init__(self): - super().__init__(CacheManager) + super().__init__(BaseCacheManager) def create(self, settings_service): # Here you would have logic to create and configure a CacheManager - return CacheManager() + # based on the settings_service + + if settings_service.settings.CACHE_TYPE == "redis": + logger.debug("Creating Redis cache") + redis_cache = RedisCache( + host=settings_service.settings.REDIS_HOST, + port=settings_service.settings.REDIS_PORT, + db=settings_service.settings.REDIS_DB, + expiration_time=settings_service.settings.REDIS_CACHE_EXPIRE, + ) + if redis_cache.is_connected(): + logger.debug("Redis cache is connected") + return redis_cache + logger.warning( + "Redis cache is not connected, falling back to in-memory cache" + ) + return InMemoryCache() + + elif settings_service.settings.CACHE_TYPE == "memory": + return InMemoryCache() diff --git a/src/backend/langflow/services/cache/flow.py b/src/backend/langflow/services/cache/flow.py deleted file mode 100644 index 0c10c51e1..000000000 --- a/src/backend/langflow/services/cache/flow.py +++ /dev/null @@ -1,146 +0,0 @@ -import threading -import time -from collections import OrderedDict - -from langflow.services.cache.base import BaseCache - - -class InMemoryCache(BaseCache): - """ - A simple in-memory cache using an OrderedDict. - - This cache supports setting a maximum size and expiration time for cached items. - When the cache is full, it uses a Least Recently Used (LRU) eviction policy. - Thread-safe using a threading Lock. - - Attributes: - max_size (int, optional): Maximum number of items to store in the cache. - expiration_time (int, optional): Time in seconds after which a cached item expires. Default is 1 hour. - - Example: - - cache = InMemoryCache(max_size=3, expiration_time=5) - - # setting cache values - cache.set("a", 1) - cache.set("b", 2) - cache["c"] = 3 - - # getting cache values - a = cache.get("a") - b = cache["b"] - """ - - def __init__(self, max_size=None, expiration_time=60 * 60): - """ - Initialize a new InMemoryCache instance. - - Args: - max_size (int, optional): Maximum number of items to store in the cache. - expiration_time (int, optional): Time in seconds after which a cached item expires. Default is 1 hour. - """ - self._cache = OrderedDict() - self._lock = threading.Lock() - self.max_size = max_size - self.expiration_time = expiration_time - - def get(self, key): - """ - Retrieve an item from the cache. - - Args: - key: The key of the item to retrieve. - - Returns: - The value associated with the key, or None if the key is not found or the item has expired. - """ - with self._lock: - if key in self._cache: - item = self._cache.pop(key) - if ( - self.expiration_time is None - or time.time() - item["time"] < self.expiration_time - ): - # Move the key to the end to make it recently used - self._cache[key] = item - return item["value"] - else: - self.delete(key) - return None - - def set(self, key, value): - """ - Add an item to the cache. - - If the cache is full, the least recently used item is evicted. - - Args: - key: The key of the item. - value: The value to cache. - """ - with self._lock: - if key in self._cache: - # Remove existing key before re-inserting to update order - self.delete(key) - elif self.max_size and len(self._cache) >= self.max_size: - # Remove least recently used item - self._cache.popitem(last=False) - self._cache[key] = {"value": value, "time": time.time()} - - def get_or_set(self, key, value): - """ - Retrieve an item from the cache. If the item does not exist, set it with the provided value. - - Args: - key: The key of the item. - value: The value to cache if the item doesn't exist. - - Returns: - The cached value associated with the key. - """ - with self._lock: - if key in self._cache: - return self.get(key) - self.set(key, value) - return value - - def delete(self, key): - """ - Remove an item from the cache. - - Args: - key: The key of the item to remove. - """ - # with self._lock: - self._cache.pop(key, None) - - def clear(self): - """ - Clear all items from the cache. - """ - with self._lock: - self._cache.clear() - - def __contains__(self, key): - """Check if the key is in the cache.""" - return key in self._cache - - def __getitem__(self, key): - """Retrieve an item from the cache using the square bracket notation.""" - return self.get(key) - - def __setitem__(self, key, value): - """Add an item to the cache using the square bracket notation.""" - self.set(key, value) - - def __delitem__(self, key): - """Remove an item from the cache using the square bracket notation.""" - self.delete(key) - - def __len__(self): - """Return the number of items in the cache.""" - return len(self._cache) - - def __repr__(self): - """Return a string representation of the InMemoryCache instance.""" - return f"InMemoryCache(max_size={self.max_size}, expiration_time={self.expiration_time})" diff --git a/src/backend/langflow/services/cache/manager.py b/src/backend/langflow/services/cache/manager.py index ce9a338ef..6243c2507 100644 --- a/src/backend/langflow/services/cache/manager.py +++ b/src/backend/langflow/services/cache/manager.py @@ -1,153 +1,256 @@ -from contextlib import contextmanager -from typing import Any, Awaitable, Callable, List, Optional -from langflow.services.base import Service +import threading +import time +from collections import OrderedDict -import pandas as pd -from PIL import Image +from langflow.services.cache.base import BaseCacheManager + +import pickle +import redis -class Subject: - """Base class for implementing the observer pattern.""" +class InMemoryCache(BaseCacheManager): - def __init__(self): - self.observers: List[Callable[[], None]] = [] + """ + A simple in-memory cache using an OrderedDict. - def attach(self, observer: Callable[[], None]): - """Attach an observer to the subject.""" - self.observers.append(observer) + This cache supports setting a maximum size and expiration time for cached items. + When the cache is full, it uses a Least Recently Used (LRU) eviction policy. + Thread-safe using a threading Lock. - def detach(self, observer: Callable[[], None]): - """Detach an observer from the subject.""" - self.observers.remove(observer) + Attributes: + max_size (int, optional): Maximum number of items to store in the cache. + expiration_time (int, optional): Time in seconds after which a cached item expires. Default is 1 hour. - def notify(self): - """Notify all observers about an event.""" - for observer in self.observers: - if observer is None: - continue - observer() + Example: + cache = InMemoryCache(max_size=3, expiration_time=5) -class AsyncSubject: - """Base class for implementing the async observer pattern.""" + # setting cache values + cache.set("a", 1) + cache.set("b", 2) + cache["c"] = 3 - def __init__(self): - self.observers: List[Callable[[], Awaitable]] = [] + # getting cache values + a = cache.get("a") + b = cache["b"] + """ - def attach(self, observer: Callable[[], Awaitable]): - """Attach an observer to the subject.""" - self.observers.append(observer) - - def detach(self, observer: Callable[[], Awaitable]): - """Detach an observer from the subject.""" - self.observers.remove(observer) - - async def notify(self): - """Notify all observers about an event.""" - for observer in self.observers: - if observer is None: - continue - await observer() - - -class CacheManager(Subject, Service): - """Manages cache for different clients and notifies observers on changes.""" - - name = "cache_manager" - - def __init__(self): - super().__init__() - self._cache = {} - self.current_client_id = None - self.current_cache = {} - - @contextmanager - def set_client_id(self, client_id: str): + def __init__(self, max_size=None, expiration_time=60 * 60): """ - Context manager to set the current client_id and associated cache. + Initialize a new InMemoryCache instance. Args: - client_id (str): The client identifier. + max_size (int, optional): Maximum number of items to store in the cache. + expiration_time (int, optional): Time in seconds after which a cached item expires. Default is 1 hour. + """ + self._cache = OrderedDict() + self._lock = threading.Lock() + self.max_size = max_size + self.expiration_time = expiration_time + + def get(self, key): + """ + Retrieve an item from the cache. + + Args: + key: The key of the item to retrieve. + + Returns: + The value associated with the key, or None if the key is not found or the item has expired. + """ + with self._lock: + if key in self._cache: + item = self._cache.pop(key) + if ( + self.expiration_time is None + or time.time() - item["time"] < self.expiration_time + ): + # Move the key to the end to make it recently used + self._cache[key] = item + return item["value"] + else: + self.delete(key) + return None + + def set(self, key, value): + """ + Add an item to the cache. + + If the cache is full, the least recently used item is evicted. + + Args: + key: The key of the item. + value: The value to cache. + """ + with self._lock: + if key in self._cache: + # Remove existing key before re-inserting to update order + self.delete(key) + elif self.max_size and len(self._cache) >= self.max_size: + # Remove least recently used item + self._cache.popitem(last=False) + self._cache[key] = {"value": value, "time": time.time()} + + def get_or_set(self, key, value): + """ + Retrieve an item from the cache. If the item does not exist, set it with the provided value. + + Args: + key: The key of the item. + value: The value to cache if the item doesn't exist. + + Returns: + The cached value associated with the key. + """ + with self._lock: + if key in self._cache: + return self.get(key) + self.set(key, value) + return value + + def delete(self, key): + """ + Remove an item from the cache. + + Args: + key: The key of the item to remove. + """ + # with self._lock: + self._cache.pop(key, None) + + def clear(self): + """ + Clear all items from the cache. + """ + with self._lock: + self._cache.clear() + + def __contains__(self, key): + """Check if the key is in the cache.""" + return key in self._cache + + def __getitem__(self, key): + """Retrieve an item from the cache using the square bracket notation.""" + return self.get(key) + + def __setitem__(self, key, value): + """Add an item to the cache using the square bracket notation.""" + self.set(key, value) + + def __delitem__(self, key): + """Remove an item from the cache using the square bracket notation.""" + self.delete(key) + + def __len__(self): + """Return the number of items in the cache.""" + return len(self._cache) + + def __repr__(self): + """Return a string representation of the InMemoryCache instance.""" + return f"InMemoryCache(max_size={self.max_size}, expiration_time={self.expiration_time})" + + +class RedisCache(BaseCacheManager): + """ + A Redis-based cache implementation. + + This cache supports setting an expiration time for cached items. + + Attributes: + expiration_time (int, optional): Time in seconds after which a cached item expires. Default is 1 hour. + + Example: + + cache = RedisCache(expiration_time=5) + + # setting cache values + cache.set("a", 1) + cache.set("b", 2) + cache["c"] = 3 + + # getting cache values + a = cache.get("a") + b = cache["b"] + """ + + def __init__(self, host="localhost", port=6379, db=0, expiration_time=60 * 60): + """ + Initialize a new RedisCache instance. + + Args: + host (str, optional): Redis host. + port (int, optional): Redis port. + db (int, optional): Redis DB. + expiration_time (int, optional): Time in seconds after which a cached item expires. Default is 1 hour. + """ + self._client = redis.StrictRedis(host=host, port=port, db=db) + self.expiration_time = expiration_time + + # check connection + def is_connected(self): + """ + Check if the Redis client is connected. """ - previous_client_id = self.current_client_id - self.current_client_id = client_id - self.current_cache = self._cache.setdefault(client_id, {}) try: - yield - finally: - self.current_client_id = previous_client_id - self.current_cache = self._cache.get(self.current_client_id, {}) + self._client.ping() + return True + except redis.exceptions.ConnectionError: + return False - def add(self, name: str, obj: Any, obj_type: str, extension: Optional[str] = None): + def get(self, key): """ - Add an object to the current client's cache. + Retrieve an item from the cache. Args: - name (str): The cache key. - obj (Any): The object to cache. - obj_type (str): The type of the object. - """ - object_extensions = { - "image": "png", - "pandas": "csv", - } - if obj_type in object_extensions: - _extension = object_extensions[obj_type] - else: - _extension = type(obj).__name__.lower() - self.current_cache[name] = { - "obj": obj, - "type": obj_type, - "extension": extension or _extension, - } - self.notify() - - def add_pandas(self, name: str, obj: Any): - """ - Add a pandas DataFrame or Series to the current client's cache. - - Args: - name (str): The cache key. - obj (Any): The pandas DataFrame or Series object. - """ - if isinstance(obj, (pd.DataFrame, pd.Series)): - self.add(name, obj.to_csv(), "pandas", extension="csv") - else: - raise ValueError("Object is not a pandas DataFrame or Series") - - def add_image(self, name: str, obj: Any, extension: str = "png"): - """ - Add a PIL Image to the current client's cache. - - Args: - name (str): The cache key. - obj (Any): The PIL Image object. - """ - if isinstance(obj, Image.Image): - self.add(name, obj, "image", extension=extension) - else: - raise ValueError("Object is not a PIL Image") - - def get(self, name: str): - """ - Get an object from the current client's cache. - - Args: - name (str): The cache key. + key: The key of the item to retrieve. Returns: - The cached object associated with the given cache key. + The value associated with the key, or None if the key is not found. """ - return self.current_cache[name] + value = self._client.get(key) + return pickle.loads(value) if value else None - def get_last(self): + def set(self, key, value): """ - Get the last added item in the current client's cache. + Add an item to the cache. - Returns: - The last added item in the cache. + Args: + key: The key of the item. + value: The value to cache. """ - return list(self.current_cache.values())[-1] + self._client.setex(key, self.expiration_time, pickle.dumps(value)) + def delete(self, key): + """ + Remove an item from the cache. -cache_manager = CacheManager() + Args: + key: The key of the item to remove. + """ + self._client.delete(key) + + def clear(self): + """ + Clear all items from the cache. + """ + self._client.flushdb() + + def __contains__(self, key): + """Check if the key is in the cache.""" + return self._client.exists(key) + + def __getitem__(self, key): + """Retrieve an item from the cache using the square bracket notation.""" + return self.get(key) + + def __setitem__(self, key, value): + """Add an item to the cache using the square bracket notation.""" + self.set(key, value) + + def __delitem__(self, key): + """Remove an item from the cache using the square bracket notation.""" + self.delete(key) + + def __repr__(self): + """Return a string representation of the RedisCache instance.""" + return f"RedisCache(expiration_time={self.expiration_time})"