Compare commits
8 Commits
f0acbf39bc
...
819b6c2232
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
819b6c2232 | ||
|
|
24da2a5861 | ||
|
|
18801ec27e | ||
|
|
8634c01ec0 | ||
|
|
71f8614cb5 | ||
|
|
80d8801728 | ||
|
|
846f5e76fe | ||
|
|
7c802a3907 |
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)
|
||||||
|
|||||||
@@ -6,6 +6,11 @@ import requests
|
|||||||
import os
|
import os
|
||||||
import logging
|
import logging
|
||||||
from typing import Optional, Dict, Any
|
from typing import Optional, Dict, Any
|
||||||
|
from src.utils.validation import (
|
||||||
|
validate_isin, validate_exchange, validate_date,
|
||||||
|
validate_int_range, sanitize_sql_string,
|
||||||
|
validate_isin_list, validate_exchange_list,
|
||||||
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -59,10 +64,19 @@ async def get_trades(isin: str = None, days: int = 7):
|
|||||||
Gibt aggregierte Analyse aller Trades zurück (nicht einzelne Trades).
|
Gibt aggregierte Analyse aller Trades zurück (nicht einzelne Trades).
|
||||||
Nutzt vorberechnete Daten aus analytics_exchange_daily.
|
Nutzt vorberechnete Daten aus analytics_exchange_daily.
|
||||||
"""
|
"""
|
||||||
|
try:
|
||||||
|
days = validate_int_range(days, 1, 365)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger days-Parameter (1-365)")
|
||||||
|
|
||||||
if isin:
|
if isin:
|
||||||
|
try:
|
||||||
|
isin = validate_isin(isin)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger ISIN-Wert")
|
||||||
# Für spezifische ISIN: hole aus trades Tabelle
|
# Für spezifische ISIN: hole aus trades Tabelle
|
||||||
query = f"""
|
query = f"""
|
||||||
select
|
select
|
||||||
date_trunc('day', timestamp) as date,
|
date_trunc('day', timestamp) as date,
|
||||||
count(*) as trade_count,
|
count(*) as trade_count,
|
||||||
sum(price * quantity) as volume,
|
sum(price * quantity) as volume,
|
||||||
@@ -109,6 +123,11 @@ async def get_summary(days: int = None):
|
|||||||
Gibt Zusammenfassung zurück. Nutzt analytics_daily_summary für total_trades.
|
Gibt Zusammenfassung zurück. Nutzt analytics_daily_summary für total_trades.
|
||||||
Optional: days Parameter für Zeitraum-basierte Zusammenfassung.
|
Optional: days Parameter für Zeitraum-basierte Zusammenfassung.
|
||||||
"""
|
"""
|
||||||
|
if days:
|
||||||
|
try:
|
||||||
|
days = validate_int_range(days, 1, 3650)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger days-Parameter (1-3650)")
|
||||||
if days:
|
if days:
|
||||||
# Zeitraum-basierte Zusammenfassung
|
# Zeitraum-basierte Zusammenfassung
|
||||||
query = f"""
|
query = f"""
|
||||||
@@ -169,6 +188,11 @@ async def get_summary(days: int = None):
|
|||||||
@app.get("/api/statistics/total-trades")
|
@app.get("/api/statistics/total-trades")
|
||||||
async def get_total_trades(days: int = None):
|
async def get_total_trades(days: int = None):
|
||||||
"""Gibt Gesamtzahl aller Trades zurück (aus analytics_daily_summary). Optional: days Parameter für Zeitraum."""
|
"""Gibt Gesamtzahl aller Trades zurück (aus analytics_daily_summary). Optional: days Parameter für Zeitraum."""
|
||||||
|
if days:
|
||||||
|
try:
|
||||||
|
days = validate_int_range(days, 1, 3650)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger days-Parameter (1-3650)")
|
||||||
if days:
|
if days:
|
||||||
query = f"select sum(total_trades) as total from analytics_daily_summary where timestamp >= dateadd('d', -{days}, now())"
|
query = f"select sum(total_trades) as total from analytics_daily_summary where timestamp >= dateadd('d', -{days}, now())"
|
||||||
else:
|
else:
|
||||||
@@ -201,17 +225,30 @@ async def get_custom_analytics(
|
|||||||
- exchanges: Komma-separierte Liste von Exchanges (optional)
|
- exchanges: Komma-separierte Liste von Exchanges (optional)
|
||||||
"""
|
"""
|
||||||
# Validiere Parameter
|
# Validiere Parameter
|
||||||
|
try:
|
||||||
|
date_from = validate_date(date_from)
|
||||||
|
date_to = validate_date(date_to)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiges Datumsformat (erwartet: YYYY-MM-DD)")
|
||||||
|
|
||||||
valid_x_axis = ["date", "exchange", "isin"]
|
valid_x_axis = ["date", "exchange", "isin"]
|
||||||
valid_y_axis = ["volume", "trade_count", "avg_price"]
|
valid_y_axis = ["volume", "trade_count", "avg_price"]
|
||||||
valid_group_by = ["exchange", "sector", "date"]
|
valid_group_by = ["exchange", "sector", "date"]
|
||||||
|
|
||||||
if x_axis not in valid_x_axis:
|
if x_axis not in valid_x_axis:
|
||||||
raise HTTPException(status_code=400, detail=f"Invalid x_axis. Must be one of: {valid_x_axis}")
|
raise HTTPException(status_code=400, detail=f"Invalid x_axis. Must be one of: {valid_x_axis}")
|
||||||
if y_axis not in valid_y_axis:
|
if y_axis not in valid_y_axis:
|
||||||
raise HTTPException(status_code=400, detail=f"Invalid y_axis. Must be one of: {valid_y_axis}")
|
raise HTTPException(status_code=400, detail=f"Invalid y_axis. Must be one of: {valid_y_axis}")
|
||||||
if group_by not in valid_group_by:
|
if group_by not in valid_group_by:
|
||||||
raise HTTPException(status_code=400, detail=f"Invalid group_by. Must be one of: {valid_group_by}")
|
raise HTTPException(status_code=400, detail=f"Invalid group_by. Must be one of: {valid_group_by}")
|
||||||
|
|
||||||
|
validated_exchanges = None
|
||||||
|
if exchanges:
|
||||||
|
try:
|
||||||
|
validated_exchanges = validate_exchange_list(exchanges)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger Exchange-Name in der Liste")
|
||||||
|
|
||||||
# Für Sektor-Gruppierung: direkter JOIN mit metadata (nicht vorberechnet)
|
# Für Sektor-Gruppierung: direkter JOIN mit metadata (nicht vorberechnet)
|
||||||
if group_by == "sector":
|
if group_by == "sector":
|
||||||
y_axis_map = {
|
y_axis_map = {
|
||||||
@@ -233,8 +270,8 @@ async def get_custom_analytics(
|
|||||||
and t.timestamp <= '{date_to}'
|
and t.timestamp <= '{date_to}'
|
||||||
"""
|
"""
|
||||||
|
|
||||||
if exchanges:
|
if validated_exchanges:
|
||||||
exchange_list = ",".join([f"'{e.strip()}'" for e in exchanges.split(",")])
|
exchange_list = ",".join([f"'{e}'" for e in validated_exchanges])
|
||||||
query += f" and t.exchange in ({exchange_list})"
|
query += f" and t.exchange in ({exchange_list})"
|
||||||
|
|
||||||
query += f" group by date_trunc('day', t.timestamp), coalesce(m.sector, 'Unbekannt') order by x_value asc, group_value asc"
|
query += f" group by date_trunc('day', t.timestamp), coalesce(m.sector, 'Unbekannt') order by x_value asc, group_value asc"
|
||||||
@@ -252,12 +289,9 @@ async def get_custom_analytics(
|
|||||||
|
|
||||||
# Nutze vorberechnete Daten aus analytics_custom
|
# Nutze vorberechnete Daten aus analytics_custom
|
||||||
exchange_filter = "all"
|
exchange_filter = "all"
|
||||||
if exchanges:
|
if validated_exchanges:
|
||||||
# Wenn mehrere Exchanges angegeben, müssen wir kombinieren
|
if len(validated_exchanges) == 1:
|
||||||
# Für jetzt: nutze nur wenn ein Exchange angegeben ist
|
exchange_filter = validated_exchanges[0]
|
||||||
exchange_list = [e.strip() for e in exchanges.split(",")]
|
|
||||||
if len(exchange_list) == 1:
|
|
||||||
exchange_filter = exchange_list[0]
|
|
||||||
else:
|
else:
|
||||||
# Bei mehreren Exchanges: gib Fehler zurück, da dies nicht vorberechnet wird
|
# Bei mehreren Exchanges: gib Fehler zurück, da dies nicht vorberechnet wird
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
@@ -317,8 +351,12 @@ async def get_moving_average(days: int = 7, exchange: str = None):
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
if exchange:
|
if exchange:
|
||||||
|
try:
|
||||||
|
exchange = validate_exchange(exchange)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger Exchange-Name")
|
||||||
query += f" and exchange = '{exchange}'"
|
query += f" and exchange = '{exchange}'"
|
||||||
|
|
||||||
query += " order by date asc, exchange asc"
|
query += " order by date asc, exchange asc"
|
||||||
|
|
||||||
data = query_questdb(query, timeout=5)
|
data = query_questdb(query, timeout=5)
|
||||||
@@ -393,6 +431,10 @@ async def get_stock_trends(days: int = 7, limit: int = 20):
|
|||||||
"""
|
"""
|
||||||
if days not in [7, 30, 42, 69, 180, 365]:
|
if days not in [7, 30, 42, 69, 180, 365]:
|
||||||
raise HTTPException(status_code=400, detail="Invalid days parameter. Must be one of: 7, 30, 42, 69, 180, 365")
|
raise HTTPException(status_code=400, detail="Invalid days parameter. Must be one of: 7, 30, 42, 69, 180, 365")
|
||||||
|
try:
|
||||||
|
limit = validate_int_range(limit, 1, 1000)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger limit-Parameter (1-1000)")
|
||||||
|
|
||||||
query = f"""
|
query = f"""
|
||||||
select
|
select
|
||||||
@@ -423,22 +465,55 @@ async def get_analytics(
|
|||||||
continents: str = None
|
continents: str = None
|
||||||
):
|
):
|
||||||
"""Analytics Endpunkt für Report Builder"""
|
"""Analytics Endpunkt für Report Builder"""
|
||||||
|
# Validiere optionale Parameter
|
||||||
|
if date_from:
|
||||||
|
try:
|
||||||
|
date_from = validate_date(date_from)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiges date_from Format (erwartet: YYYY-MM-DD)")
|
||||||
|
if date_to:
|
||||||
|
try:
|
||||||
|
date_to = validate_date(date_to)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiges date_to Format (erwartet: YYYY-MM-DD)")
|
||||||
|
|
||||||
|
validated_isins = None
|
||||||
|
if isins:
|
||||||
|
try:
|
||||||
|
validated_isins = validate_isin_list(isins)
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(status_code=400, detail="Ungueltiger ISIN-Wert in der Liste")
|
||||||
|
|
||||||
|
sanitized_continents = None
|
||||||
|
if continents:
|
||||||
|
sanitized_continents = [sanitize_sql_string(c.strip(), max_length=50) for c in continents.split(",") if c.strip()]
|
||||||
|
|
||||||
composite_keys = ["exchange_continent", "exchange_sector"]
|
composite_keys = ["exchange_continent", "exchange_sector"]
|
||||||
|
valid_metrics = ["volume", "count", "avg_price", "all"]
|
||||||
|
valid_groups = ["day", "month", "exchange", "isin", "name", "continent", "sector", "exchange_continent", "exchange_sector"]
|
||||||
|
|
||||||
|
if metric not in valid_metrics:
|
||||||
|
raise HTTPException(status_code=400, detail=f"Ungueltiger metric-Parameter. Erlaubt: {valid_metrics}")
|
||||||
|
if group_by not in valid_groups:
|
||||||
|
raise HTTPException(status_code=400, detail=f"Ungueltiger group_by-Parameter. Erlaubt: {valid_groups}")
|
||||||
|
if sub_group_by and sub_group_by not in valid_groups:
|
||||||
|
raise HTTPException(status_code=400, detail=f"Ungueltiger sub_group_by-Parameter. Erlaubt: {valid_groups}")
|
||||||
|
|
||||||
needs_metadata = any([
|
needs_metadata = any([
|
||||||
group_by in ["name", "continent", "sector"] + composite_keys,
|
group_by in ["name", "continent", "sector"] + composite_keys,
|
||||||
sub_group_by in ["name", "continent", "sector"] + composite_keys,
|
sub_group_by in ["name", "continent", "sector"] + composite_keys,
|
||||||
continents is not None
|
continents is not None
|
||||||
])
|
])
|
||||||
|
|
||||||
t_prefix = "t." if needs_metadata else ""
|
t_prefix = "t." if needs_metadata else ""
|
||||||
m_prefix = "m." if needs_metadata else ""
|
m_prefix = "m." if needs_metadata else ""
|
||||||
|
|
||||||
metrics_map = {
|
metrics_map = {
|
||||||
"volume": f"sum({t_prefix}price * {t_prefix}quantity)",
|
"volume": f"sum({t_prefix}price * {t_prefix}quantity)",
|
||||||
"count": f"count(*)",
|
"count": f"count(*)",
|
||||||
"avg_price": f"avg({t_prefix}price)"
|
"avg_price": f"avg({t_prefix}price)"
|
||||||
}
|
}
|
||||||
|
|
||||||
groups_map = {
|
groups_map = {
|
||||||
"day": f"date_trunc('day', {t_prefix}timestamp)",
|
"day": f"date_trunc('day', {t_prefix}timestamp)",
|
||||||
"month": f"date_trunc('month', {t_prefix}timestamp)",
|
"month": f"date_trunc('month', {t_prefix}timestamp)",
|
||||||
@@ -450,35 +525,35 @@ async def get_analytics(
|
|||||||
"exchange_continent": f"concat({t_prefix}exchange, ' - ', coalesce({m_prefix}continent, 'Unknown'))" if needs_metadata else "'Unknown'",
|
"exchange_continent": f"concat({t_prefix}exchange, ' - ', coalesce({m_prefix}continent, 'Unknown'))" if needs_metadata else "'Unknown'",
|
||||||
"exchange_sector": f"concat({t_prefix}exchange, ' - ', coalesce({m_prefix}sector, 'Unknown'))" if needs_metadata else "'Unknown'"
|
"exchange_sector": f"concat({t_prefix}exchange, ' - ', coalesce({m_prefix}sector, 'Unknown'))" if needs_metadata else "'Unknown'"
|
||||||
}
|
}
|
||||||
|
|
||||||
selected_metric = metrics_map.get(metric, metrics_map["volume"])
|
selected_metric = metrics_map.get(metric, metrics_map["volume"])
|
||||||
selected_group = groups_map.get(group_by, groups_map["day"])
|
selected_group = groups_map.get(group_by, groups_map["day"])
|
||||||
|
|
||||||
query = f"select {selected_group} as label"
|
query = f"select {selected_group} as label"
|
||||||
|
|
||||||
if sub_group_by and sub_group_by in groups_map:
|
if sub_group_by and sub_group_by in groups_map:
|
||||||
query += f", {groups_map[sub_group_by]} as sub_label"
|
query += f", {groups_map[sub_group_by]} as sub_label"
|
||||||
|
|
||||||
if metric == 'all':
|
if metric == 'all':
|
||||||
query += f", count(*) as value_count, sum({t_prefix}price * {t_prefix}quantity) as value_volume from trades"
|
query += f", count(*) as value_count, sum({t_prefix}price * {t_prefix}quantity) as value_volume from trades"
|
||||||
else:
|
else:
|
||||||
query += f", {selected_metric} as value from trades"
|
query += f", {selected_metric} as value from trades"
|
||||||
if needs_metadata:
|
if needs_metadata:
|
||||||
query += " t left join metadata m on t.isin = m.isin"
|
query += " t left join metadata m on t.isin = m.isin"
|
||||||
|
|
||||||
query += " where 1=1"
|
query += " where 1=1"
|
||||||
|
|
||||||
if date_from:
|
if date_from:
|
||||||
query += f" and {t_prefix}timestamp >= '{date_from}'"
|
query += f" and {t_prefix}timestamp >= '{date_from}'"
|
||||||
if date_to:
|
if date_to:
|
||||||
query += f" and {t_prefix}timestamp <= '{date_to}'"
|
query += f" and {t_prefix}timestamp <= '{date_to}'"
|
||||||
|
|
||||||
if isins:
|
if validated_isins:
|
||||||
isins_list = ",".join([f"'{i.strip()}'" for i in isins.split(",")])
|
isins_list = ",".join([f"'{i}'" for i in validated_isins])
|
||||||
query += f" and {t_prefix}isin in ({isins_list})"
|
query += f" and {t_prefix}isin in ({isins_list})"
|
||||||
|
|
||||||
if continents and needs_metadata:
|
if sanitized_continents and needs_metadata:
|
||||||
cont_list = ",".join([f"'{c.strip()}'" for c in continents.split(",")])
|
cont_list = ",".join([f"'{c}'" for c in sanitized_continents])
|
||||||
query += f" and {m_prefix}continent in ({cont_list})"
|
query += f" and {m_prefix}continent in ({cont_list})"
|
||||||
|
|
||||||
query += f" group by {selected_group}"
|
query += f" group by {selected_group}"
|
||||||
@@ -493,7 +568,8 @@ async def get_analytics(
|
|||||||
@app.get("/api/metadata/search")
|
@app.get("/api/metadata/search")
|
||||||
async def search_metadata(q: str):
|
async def search_metadata(q: str):
|
||||||
"""Case-insensitive search for ISIN or Name"""
|
"""Case-insensitive search for ISIN or Name"""
|
||||||
query = f"select isin, name from metadata where isin ilike '%{q}%' or name ilike '%{q}%' limit 10"
|
q_safe = sanitize_sql_string(q, max_length=100)
|
||||||
|
query = f"select isin, name from metadata where isin ilike '%{q_safe}%' or name ilike '%{q_safe}%' limit 10"
|
||||||
data = query_questdb(query)
|
data = query_questdb(query)
|
||||||
return format_questdb_response(data)
|
return format_questdb_response(data)
|
||||||
|
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import os
|
|||||||
import requests
|
import requests
|
||||||
from typing import Dict, List, Tuple, Optional
|
from typing import Dict, List, Tuple, Optional
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
|
from src.utils.validation import validate_table_name, validate_exchange
|
||||||
|
|
||||||
logging.basicConfig(
|
logging.basicConfig(
|
||||||
level=logging.INFO,
|
level=logging.INFO,
|
||||||
@@ -499,8 +500,12 @@ class AnalyticsWorker:
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
if exchange_filter:
|
if exchange_filter:
|
||||||
|
try:
|
||||||
|
exchange_filter = validate_exchange(exchange_filter)
|
||||||
|
except ValueError:
|
||||||
|
continue
|
||||||
query += f" and exchange = '{exchange_filter}'"
|
query += f" and exchange = '{exchange_filter}'"
|
||||||
|
|
||||||
query += f" group by date_trunc('day', timestamp), {group_by_field}"
|
query += f" group by date_trunc('day', timestamp), {group_by_field}"
|
||||||
|
|
||||||
data = self.query_questdb(query)
|
data = self.query_questdb(query)
|
||||||
@@ -764,6 +769,7 @@ class AnalyticsWorker:
|
|||||||
|
|
||||||
def get_existing_dates(self, table_name: str) -> set:
|
def get_existing_dates(self, table_name: str) -> set:
|
||||||
"""Holt alle bereits berechneten Daten aus einer Analytics-Tabelle"""
|
"""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}"
|
query = f"select distinct date_trunc('day', timestamp) as date from {table_name}"
|
||||||
data = self.query_questdb(query)
|
data = self.query_questdb(query)
|
||||||
if not data:
|
if not data:
|
||||||
@@ -875,12 +881,13 @@ class AnalyticsWorker:
|
|||||||
|
|
||||||
for table in tables:
|
for table in tables:
|
||||||
try:
|
try:
|
||||||
|
table = validate_table_name(table)
|
||||||
# QuestDB DELETE syntax
|
# QuestDB DELETE syntax
|
||||||
delete_query = f"DELETE FROM {table} WHERE timestamp >= '{date_str}' AND timestamp < '{next_day_str}'"
|
delete_query = f"DELETE FROM {table} WHERE timestamp >= '{date_str}' AND timestamp < '{next_day_str}'"
|
||||||
response = requests.get(
|
response = requests.get(
|
||||||
f"{self.questdb_url}/exec",
|
f"{self.db_url}/exec",
|
||||||
params={'query': delete_query},
|
params={'query': delete_query},
|
||||||
auth=self.auth,
|
auth=DB_AUTH,
|
||||||
timeout=30
|
timeout=30
|
||||||
)
|
)
|
||||||
if response.status_code == 200:
|
if response.status_code == 200:
|
||||||
|
|||||||
@@ -1,8 +1,11 @@
|
|||||||
import requests
|
import requests
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from typing import List
|
from typing import List
|
||||||
from ..exchanges.base import Trade
|
from ..exchanges.base import Trade
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
class DatabaseClient:
|
class DatabaseClient:
|
||||||
def __init__(self, host: str = "localhost", port: int = 9000, user: str = None, password: str = None):
|
def __init__(self, host: str = "localhost", port: int = 9000, user: str = None, password: str = None):
|
||||||
self.host = host
|
self.host = host
|
||||||
@@ -15,8 +18,8 @@ class DatabaseClient:
|
|||||||
return
|
return
|
||||||
|
|
||||||
total_trades = len(trades)
|
total_trades = len(trades)
|
||||||
print(f"Saving {total_trades} trades to QuestDB in batches of {batch_size}...")
|
logger.info(f"Speichere {total_trades} Trades in QuestDB (Batches von {batch_size})...")
|
||||||
|
|
||||||
for i in range(0, total_trades, batch_size):
|
for i in range(0, total_trades, batch_size):
|
||||||
batch = trades[i:i + batch_size]
|
batch = trades[i:i + batch_size]
|
||||||
lines = []
|
lines = []
|
||||||
@@ -25,34 +28,34 @@ class DatabaseClient:
|
|||||||
try:
|
try:
|
||||||
symbol = trade.symbol.replace(" ", "\\ ").replace(",", "\\,")
|
symbol = trade.symbol.replace(" ", "\\ ").replace(",", "\\,")
|
||||||
exchange = trade.exchange
|
exchange = trade.exchange
|
||||||
|
|
||||||
line = f"trades,exchange={exchange},symbol={symbol},isin={trade.isin} " \
|
line = f"trades,exchange={exchange},symbol={symbol},isin={trade.isin} " \
|
||||||
f"price={trade.price},quantity={trade.quantity} " \
|
f"price={trade.price},quantity={trade.quantity} " \
|
||||||
f"{int(trade.timestamp.timestamp() * 1e9)}"
|
f"{int(trade.timestamp.timestamp() * 1e9)}"
|
||||||
lines.append(line)
|
lines.append(line)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"Error formating trade {trade}: {e}")
|
logger.error(f"Fehler beim Formatieren von Trade {trade}: {e}")
|
||||||
continue
|
continue
|
||||||
|
|
||||||
if not lines:
|
if not lines:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
payload = "\n".join(lines) + "\n"
|
payload = "\n".join(lines) + "\n"
|
||||||
|
|
||||||
try:
|
try:
|
||||||
response = requests.post(
|
response = requests.post(
|
||||||
self.url,
|
self.url,
|
||||||
data=payload,
|
data=payload,
|
||||||
params={'precision': 'ns'},
|
params={'precision': 'ns'},
|
||||||
auth=self.auth
|
auth=self.auth
|
||||||
)
|
)
|
||||||
if response.status_code not in [204, 200]:
|
if response.status_code not in [204, 200]:
|
||||||
print(f"Error saving batch {i//batch_size + 1} to QuestDB: {response.text}")
|
logger.error(f"Fehler beim Speichern von Batch {i//batch_size + 1}: {response.text}")
|
||||||
else:
|
else:
|
||||||
print(f"Saved batch {i//batch_size + 1} ({len(batch)} trades)")
|
logger.info(f"Batch {i//batch_size + 1} gespeichert ({len(batch)} Trades)")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"Could not connect to QuestDB at {self.url}: {e}")
|
logger.error(f"Verbindung zu QuestDB fehlgeschlagen ({self.url}): {e}")
|
||||||
# Fallback: print to console or save to file
|
# Fallback: in Datei speichern
|
||||||
self._fallback_save(batch)
|
self._fallback_save(batch)
|
||||||
|
|
||||||
def _fallback_save(self, trades: List[Trade]):
|
def _fallback_save(self, trades: List[Trade]):
|
||||||
|
|||||||
@@ -8,11 +8,14 @@ URL-Format: https://cld42.boersenag.de/m13data/data/Mifir13DelayedData_{MIC}_{SE
|
|||||||
|
|
||||||
import requests
|
import requests
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
from typing import List, Optional
|
from typing import List, Optional
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
import re
|
import re
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Rate-Limiting Konfiguration
|
# Rate-Limiting Konfiguration
|
||||||
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
||||||
|
|
||||||
@@ -176,9 +179,9 @@ class BoersenagBase(BaseExchange):
|
|||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[{self.name}] HTTP error: {e}")
|
logger.error(f"[{self.name}] HTTP error: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[{self.name}] Error downloading {url}: {e}")
|
logger.error(f"[{self.name}] Error downloading {url}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|
||||||
@@ -265,9 +268,9 @@ class BoersenagBase(BaseExchange):
|
|||||||
target_date = self._get_last_trading_day(target_date)
|
target_date = self._get_last_trading_day(target_date)
|
||||||
|
|
||||||
if target_date != original_date:
|
if target_date != original_date:
|
||||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
logger.info(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||||
|
|
||||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||||
|
|
||||||
# Generiere mögliche URLs
|
# Generiere mögliche URLs
|
||||||
urls = self._generate_file_urls(target_date)
|
urls = self._generate_file_urls(target_date)
|
||||||
@@ -281,7 +284,7 @@ class BoersenagBase(BaseExchange):
|
|||||||
if trades:
|
if trades:
|
||||||
all_trades.extend(trades)
|
all_trades.extend(trades)
|
||||||
successful += 1
|
successful += 1
|
||||||
print(f"[{self.name}] Found {len(trades)} trades from: {url.split('/')[-1]}")
|
logger.info(f"[{self.name}] Found {len(trades)} trades from: {url.split('/')[-1]}")
|
||||||
# Bei Erfolg müssen wir nicht alle anderen URLs probieren
|
# Bei Erfolg müssen wir nicht alle anderen URLs probieren
|
||||||
break
|
break
|
||||||
|
|
||||||
@@ -293,7 +296,7 @@ class BoersenagBase(BaseExchange):
|
|||||||
if i > 20 and successful == 0:
|
if i > 20 and successful == 0:
|
||||||
break
|
break
|
||||||
|
|
||||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
logger.info(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -3,11 +3,14 @@ import gzip
|
|||||||
import csv
|
import csv
|
||||||
import io
|
import io
|
||||||
import time
|
import time
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
from typing import List, Optional
|
from typing import List, Optional
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
from bs4 import BeautifulSoup
|
from bs4 import BeautifulSoup
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Rate-Limiting
|
# Rate-Limiting
|
||||||
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
RATE_LIMIT_DELAY = 0.3 # Sekunden zwischen Requests
|
||||||
|
|
||||||
@@ -79,10 +82,10 @@ class GettexExchange(BaseExchange):
|
|||||||
url = f"https://www.gettex.de/fileadmin/posttrade-data/{filename}"
|
url = f"https://www.gettex.de/fileadmin/posttrade-data/{filename}"
|
||||||
files.append({'filename': filename, 'url': url})
|
files.append({'filename': filename, 'url': url})
|
||||||
|
|
||||||
print(f"[GETTEX] Found {len(files)} files on page")
|
logger.info(f"[GETTEX] Found {len(files)} files on page")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[GETTEX] Error fetching page: {e}")
|
logger.error(f"[GETTEX] Error fetching page: {e}")
|
||||||
|
|
||||||
return files
|
return files
|
||||||
|
|
||||||
@@ -148,11 +151,11 @@ class GettexExchange(BaseExchange):
|
|||||||
date_str = parts[1] # YYYYMMDD
|
date_str = parts[1] # YYYYMMDD
|
||||||
|
|
||||||
if not date_str:
|
if not date_str:
|
||||||
print(f"[GETTEX] WARNING: Could not extract date from filename: {filename}")
|
logger.warning(f"[GETTEX] Could not extract date from filename: {filename}")
|
||||||
|
|
||||||
# Debug: Zeige erste Zeile
|
# Debug: Zeige erste Zeile
|
||||||
if lines and len(lines) > 0:
|
if lines and len(lines) > 0:
|
||||||
print(f"[GETTEX] First line sample: {lines[0][:100]}")
|
logger.debug(f"[GETTEX] First line sample: {lines[0][:100]}")
|
||||||
|
|
||||||
# Gettex CSV hat KEINEN Header!
|
# Gettex CSV hat KEINEN Header!
|
||||||
# Format: ISIN,Zeit,Währung,Preis,Menge
|
# Format: ISIN,Zeit,Währung,Preis,Menge
|
||||||
@@ -166,24 +169,24 @@ class GettexExchange(BaseExchange):
|
|||||||
if trade:
|
if trade:
|
||||||
trades.append(trade)
|
trades.append(trade)
|
||||||
else:
|
else:
|
||||||
if i < 3: # Zeige nur erste paar Fehler
|
if i < 3:
|
||||||
print(f"[GETTEX] Failed to parse line {i+1}: {line[:80]}")
|
logger.debug(f"[GETTEX] Failed to parse line {i+1}: {line[:80]}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
parse_errors += 1
|
parse_errors += 1
|
||||||
if i < 3:
|
if i < 3:
|
||||||
print(f"[GETTEX] Exception parsing line {i+1}: {e}, line: {line[:80]}")
|
logger.debug(f"[GETTEX] Exception parsing line {i+1}: {e}, line: {line[:80]}")
|
||||||
continue
|
continue
|
||||||
|
|
||||||
if trades:
|
if trades:
|
||||||
print(f"[GETTEX] Parsed {len(trades)} trades from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
logger.info(f"[GETTEX] Parsed {len(trades)} trades from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
||||||
elif len(lines) > 0:
|
elif len(lines) > 0:
|
||||||
print(f"[GETTEX] No trades parsed from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
logger.warning(f"[GETTEX] No trades parsed from {filename} ({len(lines)} lines, {parse_errors} errors)")
|
||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[GETTEX] HTTP error downloading {filename}: {e}")
|
logger.error(f"[GETTEX] HTTP error downloading {filename}: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[GETTEX] Error downloading {filename}: {e}")
|
logger.error(f"[GETTEX] Error downloading {filename}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|
||||||
@@ -406,9 +409,9 @@ class GettexExchange(BaseExchange):
|
|||||||
target_date = self._get_last_trading_day(target_date)
|
target_date = self._get_last_trading_day(target_date)
|
||||||
|
|
||||||
if target_date != original_date:
|
if target_date != original_date:
|
||||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
logger.info(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||||
|
|
||||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||||
|
|
||||||
# Versuche zuerst, Dateien von der Webseite zu laden
|
# Versuche zuerst, Dateien von der Webseite zu laden
|
||||||
page_files = self._get_file_list_from_page()
|
page_files = self._get_file_list_from_page()
|
||||||
@@ -434,10 +437,10 @@ class GettexExchange(BaseExchange):
|
|||||||
hour = int(parts[2])
|
hour = int(parts[2])
|
||||||
if hour < 3:
|
if hour < 3:
|
||||||
target_files.append(f)
|
target_files.append(f)
|
||||||
except:
|
except (ValueError, IndexError):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
print(f"[{self.name}] Found {len(target_files)} files for target date from page")
|
logger.info(f"[{self.name}] Found {len(target_files)} files for target date from page")
|
||||||
|
|
||||||
# Lade Dateien von der Webseite (mit Rate-Limiting)
|
# Lade Dateien von der Webseite (mit Rate-Limiting)
|
||||||
for i, f in enumerate(target_files):
|
for i, f in enumerate(target_files):
|
||||||
@@ -450,9 +453,9 @@ class GettexExchange(BaseExchange):
|
|||||||
|
|
||||||
# Fallback: Versuche erwartete Dateinamen
|
# Fallback: Versuche erwartete Dateinamen
|
||||||
if not all_trades:
|
if not all_trades:
|
||||||
print(f"[{self.name}] No files from page, trying generated filenames...")
|
logger.info(f"[{self.name}] No files from page, trying generated filenames...")
|
||||||
expected_files = self._generate_expected_files(target_date)
|
expected_files = self._generate_expected_files(target_date)
|
||||||
print(f"[{self.name}] Trying {len(expected_files)} potential files")
|
logger.info(f"[{self.name}] Trying {len(expected_files)} potential files")
|
||||||
|
|
||||||
successful_files = 0
|
successful_files = 0
|
||||||
for filename in expected_files:
|
for filename in expected_files:
|
||||||
@@ -461,9 +464,9 @@ class GettexExchange(BaseExchange):
|
|||||||
all_trades.extend(trades)
|
all_trades.extend(trades)
|
||||||
successful_files += 1
|
successful_files += 1
|
||||||
|
|
||||||
print(f"[{self.name}] Successfully downloaded {successful_files} files")
|
logger.info(f"[{self.name}] Successfully downloaded {successful_files} files")
|
||||||
|
|
||||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
logger.info(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||||
|
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|
||||||
@@ -506,12 +509,12 @@ class GettexExchange(BaseExchange):
|
|||||||
continue
|
continue
|
||||||
|
|
||||||
if trades:
|
if trades:
|
||||||
print(f"[{self.name}] Parsed {len(trades)} trades from {filename}")
|
logger.info(f"[{self.name}] Parsed {len(trades)} trades from {filename}")
|
||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[{self.name}] HTTP error downloading {url}: {e}")
|
logger.error(f"[{self.name}] HTTP error downloading {url}: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[{self.name}] Error downloading {url}: {e}")
|
logger.error(f"[{self.name}] Error downloading {url}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|||||||
@@ -1,10 +1,13 @@
|
|||||||
import requests
|
import requests
|
||||||
import csv
|
import csv
|
||||||
import io
|
import io
|
||||||
|
import logging
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from typing import List
|
from typing import List
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
class LSExchange(BaseExchange):
|
class LSExchange(BaseExchange):
|
||||||
@property
|
@property
|
||||||
def name(self) -> str:
|
def name(self) -> str:
|
||||||
@@ -15,23 +18,23 @@ class LSExchange(BaseExchange):
|
|||||||
if include_yesterday:
|
if include_yesterday:
|
||||||
endpoints.append("https://www.ls-x.de/_rpc/json/.lstc/instrument/list/lstctradesyesterday")
|
endpoints.append("https://www.ls-x.de/_rpc/json/.lstc/instrument/list/lstctradesyesterday")
|
||||||
endpoints.append("https://www.ls-x.de/_rpc/json/.lstc/instrument/list/lsxtradesyesterday")
|
endpoints.append("https://www.ls-x.de/_rpc/json/.lstc/instrument/list/lsxtradesyesterday")
|
||||||
|
|
||||||
headers = {
|
headers = {
|
||||||
'User-Agent': 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
'User-Agent': 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||||
'Accept': 'application/json',
|
'Accept': 'application/json',
|
||||||
'Referer': 'https://www.ls-tc.de/'
|
'Referer': 'https://www.ls-tc.de/'
|
||||||
}
|
}
|
||||||
|
|
||||||
all_trades = []
|
all_trades = []
|
||||||
for url in endpoints:
|
for url in endpoints:
|
||||||
try:
|
try:
|
||||||
response = requests.get(url, headers=headers)
|
response = requests.get(url, headers=headers)
|
||||||
response.raise_for_status()
|
response.raise_for_status()
|
||||||
|
|
||||||
f = io.StringIO(response.text)
|
f = io.StringIO(response.text)
|
||||||
# Header: isin;displayName;tradeTime;price;currency;size;orderId
|
# Header: isin;displayName;tradeTime;price;currency;size;orderId
|
||||||
reader = csv.DictReader(f, delimiter=';')
|
reader = csv.DictReader(f, delimiter=';')
|
||||||
|
|
||||||
for item in reader:
|
for item in reader:
|
||||||
try:
|
try:
|
||||||
price = float(item['price'].replace(',', '.'))
|
price = float(item['price'].replace(',', '.'))
|
||||||
@@ -39,11 +42,11 @@ class LSExchange(BaseExchange):
|
|||||||
isin = item['isin']
|
isin = item['isin']
|
||||||
symbol = item['displayName']
|
symbol = item['displayName']
|
||||||
time_str = item['tradeTime']
|
time_str = item['tradeTime']
|
||||||
|
|
||||||
# Format: 2026-01-23T07:30:00.992000Z
|
# Format: 2026-01-23T07:30:00.992000Z
|
||||||
ts_str = time_str.replace('Z', '+00:00')
|
ts_str = time_str.replace('Z', '+00:00')
|
||||||
timestamp = datetime.fromisoformat(ts_str)
|
timestamp = datetime.fromisoformat(ts_str)
|
||||||
|
|
||||||
all_trades.append(Trade(
|
all_trades.append(Trade(
|
||||||
exchange=self.name,
|
exchange=self.name,
|
||||||
symbol=symbol,
|
symbol=symbol,
|
||||||
@@ -52,8 +55,9 @@ class LSExchange(BaseExchange):
|
|||||||
quantity=quantity,
|
quantity=quantity,
|
||||||
timestamp=timestamp
|
timestamp=timestamp
|
||||||
))
|
))
|
||||||
except Exception:
|
except (ValueError, KeyError) as e:
|
||||||
|
logger.debug(f"Fehler beim Parsen einer LS-Zeile: {e}")
|
||||||
continue
|
continue
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"Error fetching LS data from {url}: {e}")
|
logger.error(f"Fehler beim Abrufen von LS-Daten von {url}: {e}")
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|||||||
@@ -3,11 +3,14 @@ import gzip
|
|||||||
import json
|
import json
|
||||||
import csv
|
import csv
|
||||||
import io
|
import io
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
from typing import List, Optional
|
from typing import List, Optional
|
||||||
from .base import BaseExchange, Trade
|
from .base import BaseExchange, Trade
|
||||||
from bs4 import BeautifulSoup
|
from bs4 import BeautifulSoup
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Browser User-Agent (Vollständiger Browser-Fingerprint für Stuttgart)
|
# Browser User-Agent (Vollständiger Browser-Fingerprint für Stuttgart)
|
||||||
HEADERS = {
|
HEADERS = {
|
||||||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||||
@@ -84,7 +87,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
files = self._generate_expected_urls()
|
files = self._generate_expected_urls()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error fetching page: {e}")
|
logger.error(f"[STU] Error fetching page: {e}")
|
||||||
files = self._generate_expected_urls()
|
files = self._generate_expected_urls()
|
||||||
|
|
||||||
return files
|
return files
|
||||||
@@ -194,13 +197,13 @@ class StuttgartExchange(BaseExchange):
|
|||||||
if trade:
|
if trade:
|
||||||
trades.append(trade)
|
trades.append(trade)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Could not parse {url}: {e}")
|
logger.error(f"[STU] Could not parse {url}: {e}")
|
||||||
|
|
||||||
except requests.exceptions.HTTPError as e:
|
except requests.exceptions.HTTPError as e:
|
||||||
if e.response.status_code != 404:
|
if e.response.status_code != 404:
|
||||||
print(f"[STU] HTTP error downloading {url}: {e}")
|
logger.error(f"[STU] HTTP error downloading {url}: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error downloading {url}: {e}")
|
logger.error(f"[STU] Error downloading {url}: {e}")
|
||||||
|
|
||||||
return trades
|
return trades
|
||||||
|
|
||||||
@@ -275,7 +278,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error parsing JSON record: {e}")
|
logger.debug(f"[STU] Error parsing JSON record: {e}")
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def _parse_csv_row(self, row: dict) -> Optional[Trade]:
|
def _parse_csv_row(self, row: dict) -> Optional[Trade]:
|
||||||
@@ -331,7 +334,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"[STU] Error parsing CSV row: {e}")
|
logger.debug(f"[STU] Error parsing CSV row: {e}")
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def _get_last_trading_day(self, from_date) -> datetime.date:
|
def _get_last_trading_day(self, from_date) -> datetime.date:
|
||||||
@@ -365,13 +368,13 @@ class StuttgartExchange(BaseExchange):
|
|||||||
target_date = self._get_last_trading_day(target_date)
|
target_date = self._get_last_trading_day(target_date)
|
||||||
|
|
||||||
if target_date != original_date:
|
if target_date != original_date:
|
||||||
print(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
logger.info(f"[{self.name}] Skipping weekend: {original_date} -> {target_date}")
|
||||||
|
|
||||||
print(f"[{self.name}] Fetching trades for date: {target_date}")
|
logger.info(f"[{self.name}] Fetching trades for date: {target_date}")
|
||||||
|
|
||||||
# Download-Links holen
|
# Download-Links holen
|
||||||
all_links = self._get_download_links()
|
all_links = self._get_download_links()
|
||||||
print(f"[{self.name}] Found {len(all_links)} potential download links")
|
logger.info(f"[{self.name}] Found {len(all_links)} potential download links")
|
||||||
|
|
||||||
# Nach Datum filtern
|
# Nach Datum filtern
|
||||||
target_links = self._filter_files_for_date(all_links, target_date)
|
target_links = self._filter_files_for_date(all_links, target_date)
|
||||||
@@ -380,7 +383,7 @@ class StuttgartExchange(BaseExchange):
|
|||||||
# Fallback: Versuche alle Links
|
# Fallback: Versuche alle Links
|
||||||
target_links = all_links
|
target_links = all_links
|
||||||
|
|
||||||
print(f"[{self.name}] Trying {len(target_links)} files for target date")
|
logger.info(f"[{self.name}] Trying {len(target_links)} files for target date")
|
||||||
|
|
||||||
# Dateien herunterladen und parsen
|
# Dateien herunterladen und parsen
|
||||||
successful = 0
|
successful = 0
|
||||||
@@ -389,9 +392,9 @@ class StuttgartExchange(BaseExchange):
|
|||||||
if trades:
|
if trades:
|
||||||
all_trades.extend(trades)
|
all_trades.extend(trades)
|
||||||
successful += 1
|
successful += 1
|
||||||
print(f"[{self.name}] Parsed {len(trades)} trades from {url}")
|
logger.info(f"[{self.name}] Parsed {len(trades)} trades from {url}")
|
||||||
|
|
||||||
print(f"[{self.name}] Successfully processed {successful} files")
|
logger.info(f"[{self.name}] Successfully processed {successful} files")
|
||||||
print(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
logger.info(f"[{self.name}] Total trades fetched: {len(all_trades)}")
|
||||||
|
|
||||||
return all_trades
|
return all_trades
|
||||||
|
|||||||
0
src/utils/__init__.py
Normal file
0
src/utils/__init__.py
Normal file
124
src/utils/validation.py
Normal file
124
src/utils/validation.py
Normal file
@@ -0,0 +1,124 @@
|
|||||||
|
"""
|
||||||
|
Zentrale Validierungs- und Sanitisierungsfunktionen fuer SQL-Queries.
|
||||||
|
|
||||||
|
Da QuestDB's HTTP API keine parametrisierten Queries unterstuetzt,
|
||||||
|
muessen alle Werte vor der Interpolation validiert/bereinigt werden.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import re
|
||||||
|
import datetime
|
||||||
|
|
||||||
|
|
||||||
|
# Gueltige Exchange-Namen (alle registrierten Boersen)
|
||||||
|
VALID_EXCHANGES = {
|
||||||
|
"EIX", "LS", "XETRA", "FRA", "QUOTRIX", "GETTEX", "STU",
|
||||||
|
"DUSA", "DUSB", "DUSC", "DUSD",
|
||||||
|
"HAMA", "HAMB", "HANA", "HANB",
|
||||||
|
"NONE", # Platzhalter fuer leere Analytics-Eintraege
|
||||||
|
}
|
||||||
|
|
||||||
|
# Gueltige Analytics-Tabellennamen
|
||||||
|
VALID_TABLES = {
|
||||||
|
"analytics_custom", "analytics_exchange_daily",
|
||||||
|
"analytics_daily_summary", "analytics_volume_changes",
|
||||||
|
"analytics_stock_trends", "trades", "metadata",
|
||||||
|
}
|
||||||
|
|
||||||
|
# ISIN-Format: 2 Buchstaben Laendercode + 9 alphanumerische Zeichen + 1 Pruefziffer
|
||||||
|
_ISIN_PATTERN = re.compile(r'^[A-Z]{2}[A-Z0-9]{9}[0-9]$')
|
||||||
|
|
||||||
|
# Datumsformat: YYYY-MM-DD
|
||||||
|
_DATE_PATTERN = re.compile(r'^\d{4}-\d{2}-\d{2}$')
|
||||||
|
|
||||||
|
# Gefaehrliche SQL-Fragmente
|
||||||
|
_SQL_DANGEROUS = re.compile(r'(--|/\*|\*/|;)')
|
||||||
|
|
||||||
|
|
||||||
|
def validate_isin(value: str) -> str:
|
||||||
|
"""Validiert einen ISIN-Wert. Wirft ValueError bei ungueltigem Format."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"ISIN muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip().upper()
|
||||||
|
if not _ISIN_PATTERN.match(value):
|
||||||
|
raise ValueError(f"Ungueltiges ISIN-Format: {value!r}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_exchange(value: str) -> str:
|
||||||
|
"""Validiert einen Exchange-Namen gegen die Whitelist."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Exchange muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip().upper()
|
||||||
|
if value not in VALID_EXCHANGES:
|
||||||
|
raise ValueError(f"Unbekannte Exchange: {value!r}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_date(value: str) -> str:
|
||||||
|
"""Validiert ein Datum im Format YYYY-MM-DD."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Datum muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip()
|
||||||
|
if not _DATE_PATTERN.match(value):
|
||||||
|
raise ValueError(f"Ungueltiges Datumsformat: {value!r} (erwartet: YYYY-MM-DD)")
|
||||||
|
# Pruefe ob es ein gueltiges Datum ist
|
||||||
|
try:
|
||||||
|
datetime.date.fromisoformat(value)
|
||||||
|
except ValueError:
|
||||||
|
raise ValueError(f"Ungueltiges Datum: {value!r}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_int_range(value: int, min_val: int = 0, max_val: int = 10000) -> int:
|
||||||
|
"""Validiert einen Integer-Wert innerhalb eines Bereichs."""
|
||||||
|
value = int(value)
|
||||||
|
if value < min_val or value > max_val:
|
||||||
|
raise ValueError(f"Wert {value} ausserhalb des erlaubten Bereichs [{min_val}, {max_val}]")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def sanitize_sql_string(value: str, max_length: int = 200) -> str:
|
||||||
|
"""
|
||||||
|
Bereinigt einen String fuer die Verwendung in SQL-Queries.
|
||||||
|
Escapet Single-Quotes und entfernt gefaehrliche SQL-Fragmente.
|
||||||
|
"""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Wert muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value[:max_length]
|
||||||
|
# Single-Quotes escapen (SQL-Standard)
|
||||||
|
value = value.replace("'", "''")
|
||||||
|
# Gefaehrliche SQL-Fragmente entfernen
|
||||||
|
value = _SQL_DANGEROUS.sub('', value)
|
||||||
|
# Backslashes entfernen
|
||||||
|
value = value.replace('\\', '')
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def validate_isin_list(value: str) -> list:
|
||||||
|
"""Validiert eine komma-separierte Liste von ISINs."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError("ISIN-Liste muss ein String sein")
|
||||||
|
items = [item.strip() for item in value.split(",") if item.strip()]
|
||||||
|
if not items:
|
||||||
|
raise ValueError("ISIN-Liste ist leer")
|
||||||
|
return [validate_isin(item) for item in items]
|
||||||
|
|
||||||
|
|
||||||
|
def validate_exchange_list(value: str) -> list:
|
||||||
|
"""Validiert eine komma-separierte Liste von Exchange-Namen."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError("Exchange-Liste muss ein String sein")
|
||||||
|
items = [item.strip() for item in value.split(",") if item.strip()]
|
||||||
|
if not items:
|
||||||
|
raise ValueError("Exchange-Liste ist leer")
|
||||||
|
return [validate_exchange(item) for item in items]
|
||||||
|
|
||||||
|
|
||||||
|
def validate_table_name(value: str) -> str:
|
||||||
|
"""Validiert einen Tabellennamen gegen die Whitelist."""
|
||||||
|
if not isinstance(value, str):
|
||||||
|
raise ValueError(f"Tabellenname muss ein String sein, nicht {type(value).__name__}")
|
||||||
|
value = value.strip().lower()
|
||||||
|
if value not in VALID_TABLES:
|
||||||
|
raise ValueError(f"Unbekannter Tabellenname: {value!r}")
|
||||||
|
return value
|
||||||
Reference in New Issue
Block a user