diff --git a/src/analytics/worker.py b/src/analytics/worker.py index a0810be..18d8750 100644 --- a/src/analytics/worker.py +++ b/src/analytics/worker.py @@ -5,6 +5,7 @@ import os import requests from typing import Dict, List, Tuple, Optional import pandas as pd +from src.utils.validation import validate_table_name, validate_exchange logging.basicConfig( level=logging.INFO, @@ -499,8 +500,12 @@ class AnalyticsWorker: """ if exchange_filter: + try: + exchange_filter = validate_exchange(exchange_filter) + except ValueError: + continue query += f" and exchange = '{exchange_filter}'" - + query += f" group by date_trunc('day', timestamp), {group_by_field}" data = self.query_questdb(query) @@ -764,6 +769,7 @@ class AnalyticsWorker: def get_existing_dates(self, table_name: str) -> set: """Holt alle bereits berechneten Daten aus einer Analytics-Tabelle""" + table_name = validate_table_name(table_name) query = f"select distinct date_trunc('day', timestamp) as date from {table_name}" data = self.query_questdb(query) if not data: @@ -875,12 +881,13 @@ class AnalyticsWorker: for table in tables: try: + table = validate_table_name(table) # QuestDB DELETE syntax delete_query = f"DELETE FROM {table} WHERE timestamp >= '{date_str}' AND timestamp < '{next_day_str}'" response = requests.get( - f"{self.questdb_url}/exec", + f"{self.db_url}/exec", params={'query': delete_query}, - auth=self.auth, + auth=DB_AUTH, timeout=30 ) if response.status_code == 200: