Exchange name validation in get_existing_trades_for_day() and get_last_trade_timestamp(), cache size limit (50 entries with FIFO eviction)
This commit is contained in:
11
daemon.py
11
daemon.py
@@ -17,6 +17,7 @@ from src.exchanges.boersenag import (
|
|||||||
HAMAExchange, HAMBExchange, HANAExchange, HANBExchange
|
HAMAExchange, HAMBExchange, HANAExchange, HANBExchange
|
||||||
)
|
)
|
||||||
from src.database.questdb_client import DatabaseClient
|
from src.database.questdb_client import DatabaseClient
|
||||||
|
from src.utils.validation import validate_exchange
|
||||||
|
|
||||||
logging.basicConfig(
|
logging.basicConfig(
|
||||||
level=logging.INFO,
|
level=logging.INFO,
|
||||||
@@ -67,6 +68,7 @@ STANDARD_EXCHANGES: List[Type[BaseExchange]] = [
|
|||||||
|
|
||||||
# Cache für existierende Trades pro Tag (wird nach jedem Exchange geleert)
|
# Cache für existierende Trades pro Tag (wird nach jedem Exchange geleert)
|
||||||
_existing_trades_cache = {}
|
_existing_trades_cache = {}
|
||||||
|
MAX_CACHE_SIZE = 50
|
||||||
|
|
||||||
def get_trade_hash(trade):
|
def get_trade_hash(trade):
|
||||||
"""Erstellt einen eindeutigen Hash für einen 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):
|
def get_existing_trades_for_day(db_url, exchange_name, day):
|
||||||
"""Holt existierende Trades für einen Tag aus der DB (mit Caching)."""
|
"""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')}"
|
cache_key = f"{exchange_name}_{day.strftime('%Y-%m-%d')}"
|
||||||
|
|
||||||
if cache_key in _existing_trades_cache:
|
if cache_key in _existing_trades_cache:
|
||||||
return _existing_trades_cache[cache_key]
|
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:
|
except Exception as e:
|
||||||
logger.warning(f"Error fetching existing trades for {day}: {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
|
_existing_trades_cache[cache_key] = existing_trades
|
||||||
return 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:
|
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."""
|
"""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"
|
query = f"trades where exchange = '{exchange_name}' latest by timestamp"
|
||||||
try:
|
try:
|
||||||
response = requests.get(f"{db_url}/exec", params={'query': query}, auth=DB_AUTH)
|
response = requests.get(f"{db_url}/exec", params={'query': query}, auth=DB_AUTH)
|
||||||
|
|||||||
Reference in New Issue
Block a user