diff --git a/daemon.py b/daemon.py index 842924e..68dca6b 100644 --- a/daemon.py +++ b/daemon.py @@ -17,6 +17,7 @@ from src.exchanges.boersenag import ( HAMAExchange, HAMBExchange, HANAExchange, HANBExchange ) from src.database.questdb_client import DatabaseClient +from src.utils.validation import validate_exchange logging.basicConfig( level=logging.INFO, @@ -67,6 +68,7 @@ STANDARD_EXCHANGES: List[Type[BaseExchange]] = [ # Cache für existierende Trades pro Tag (wird nach jedem Exchange geleert) _existing_trades_cache = {} +MAX_CACHE_SIZE = 50 def get_trade_hash(trade): """Erstellt einen eindeutigen Hash für einen Trade.""" @@ -75,8 +77,9 @@ def get_trade_hash(trade): def get_existing_trades_for_day(db_url, exchange_name, day): """Holt existierende Trades für einen Tag aus der DB (mit Caching).""" + exchange_name = validate_exchange(exchange_name) cache_key = f"{exchange_name}_{day.strftime('%Y-%m-%d')}" - + if cache_key in _existing_trades_cache: return _existing_trades_cache[cache_key] @@ -109,6 +112,11 @@ def get_existing_trades_for_day(db_url, exchange_name, day): except Exception as e: logger.warning(f"Error fetching existing trades for {day}: {e}") + # Cache-Groesse begrenzen (FIFO via dict ordering, Python 3.7+) + if len(_existing_trades_cache) >= MAX_CACHE_SIZE: + oldest_key = next(iter(_existing_trades_cache)) + del _existing_trades_cache[oldest_key] + _existing_trades_cache[cache_key] = existing_trades return existing_trades @@ -163,6 +171,7 @@ def filter_new_trades_batch(db_url, exchange_name, trades, batch_size=5000): def get_last_trade_timestamp(db_url: str, exchange_name: str) -> datetime.datetime: """Holt den Timestamp des letzten Trades für eine Exchange aus QuestDB.""" + exchange_name = validate_exchange(exchange_name) query = f"trades where exchange = '{exchange_name}' latest by timestamp" try: response = requests.get(f"{db_url}/exec", params={'query': query}, auth=DB_AUTH)