Terminmarktmerkmale, Ollama-Erklärungen, schnelleres Lernen

Drei Erweiterungen aus der vorangegangenen Analyse.

Lerngeschwindigkeit
- background_sample_every_n_bars von 10 auf 5. Gemessen stammen nur rund 3 % der
  Beobachtungen aus echten Trades; diese Stichproben sind der wirksamste Hebel.
  Nicht weiter gesenkt, weil benachbarte Kerzen stark korreliert sind und sich die
  Label-Fenster überlappen - mehr Gradientenschritte heißt dort nicht mehr
  Information.

Funding Rate und Open Interest als Merkmale (strategy.derivatives, standardmäßig aus)
- Vier zusätzliche Merkmale vom Perpetual zum jeweiligen Spot-Paar. Gehandelt wird
  weiterhin Spot, die Kennzahlen kommen über einen zweiten ccxt-Client mit
  defaultType=future.
- Die Zuordnung ist lookahead-frei: Für jede Kerze gilt nur der Wert, der zu diesem
  Zeitpunkt bereits veröffentlicht war.
- Fällt eine Quelle aus oder deckt sie weniger als min_coverage ab, bleiben die
  Spalten neutral. Die Modelldimension bleibt dabei stabil.
- Gemessene API-Grenzen bei Binance: Open Interest reicht 30 Tage zurück, 500 Zeilen
  je Abruf; Funding Rate über ein Jahr. Beide Merkmale sind deshalb einzeln
  abschaltbar.

  ERGEBNIS: kein Nutzen. Zwei Backtests mit identischen Kerzen und Seed -
  5m/20 Tage: Rendite -1,79 % auf -1,90 %, Accuracy 51,2 % auf 50,4 %;
  15m/28 Tage: Rendite -2,22 % auf -2,67 %, LogLoss praktisch unverändert.
  Die Anbindung arbeitet einwandfrei (100 % Datenabdeckung), das Modell gewichtet
  die neuen Merkmale aber nur mit 0,01 bis 0,09 gegenüber 0,39 für ema_spread.
  Die Funktion bleibt aus und ist dafür da, das auf anderen Zeiträumen selbst zu
  prüfen - nicht weil sie sich bewährt hätte.

Ollama-Erklärungen (llm, standardmäßig aus)
- Neuer Dashboard-Bereich und POST /control/explain. Das Modell bekommt den Zustand
  als Text und gibt Text zurück; es entscheidet nichts, beeinflusst keine Order und
  wird nie aus dem Handels-Loop heraus aufgerufen. Der System-Prompt untersagt
  Anlageempfehlungen und Kursprognosen.
- Bewusst nicht als Entscheider: nicht reproduzierbar, kaum backtestbar, und es
  würde die Nachvollziehbarkeit des linearen Modells zerstören.
- Beim Test an qwen3.8:27b zeigte sich ein echter Fehler: Reasoning-Modelle legen
  ihre Denkschritte in ein eigenes Antwortfeld und verbrauchten dafür das gesamte
  Token-Budget, response blieb leer. llm.think ist jetzt standardmäßig false, die
  Fehlermeldung nennt Ursache und Ausweg statt nur "leere Antwort", und für ältere
  Ollama-Versionen ohne das Feld gibt es einen Wiederholungsversuch ohne es.

Bewusst nicht enthalten: News- und Google-Trends-Sentiment. Es fehlt eine Quelle mit
Point-in-Time-Historie; ohne die lässt sich das Merkmal nicht backtesten. Nach dem
Ergebnis oben wäre ein unvalidiertes Merkmal der falsche Schritt.

test_stop_loss_bounds_the_worst_trade läuft jetzt mit strategy.name: rules. Er hing
über die geänderte Voreinstellung an zufälligem Modellverhalten, geprüft werden soll
aber die Stop-Logik.

270 Tests (36 neue), ruff sauber. Gegen echte Binance-Daten und eine laufende
Ollama-Instanz geprüft, Dashboard-Bereich im Browser bedient.
This commit is contained in:
Tobias Zimmermann
2026-08-23 14:46:24 +02:00
parent db31bc93f7
commit d36ed142cf
13 changed files with 1403 additions and 23 deletions
+79
View File
@@ -235,6 +235,84 @@ im Paper-Betrieb gereiftes Modell unverändert bleiben soll.
---
## Terminmarktdaten als Zusatzmerkmale
Optional fließen **Funding Rate** und **Open Interest** vom zugehörigen Perpetual
(`BTC/USDT``BTC/USDT:USDT`) als vier weitere Merkmale ein. Gehandelt wird weiterhin Spot.
```yaml
strategy:
derivatives:
enabled: true
```
Die Zuordnung ist lookahead-frei: Für jede Kerze gilt nur der Wert, der zu diesem Zeitpunkt
**bereits veröffentlicht** war. Fällt eine Quelle aus oder deckt sie weniger als
`min_coverage` der Kerzen ab, bleiben die Spalten neutral — die Modelldimension ändert sich
nicht, der Bot läuft weiter.
> **Gemessenes Ergebnis: kein Nutzen.** Zwei Backtests mit identischen Kerzen und Seed:
>
> | Lauf | Rendite | Trefferquote | Modell-Accuracy | LogLoss |
> |---|---|---|---|---|
> | 5m, 20 Tage, ohne | 1,79 % | 18,9 % | 51,2 % | 0,7262 |
> | 5m, 20 Tage, mit | 1,90 % | 21,4 % | 50,4 % | 0,7445 |
> | 15m, 28 Tage, ohne | 2,22 % | 20,7 % | 46,8 % | 0,7856 |
> | 15m, 28 Tage, mit | 2,67 % | 18,8 % | 48,3 % | 0,7831 |
>
> Die Datenanbindung funktioniert (100 % Abdeckung in beiden Läufen), aber das Modell
> gewichtet die neuen Merkmale schwach (0,010,09 gegenüber 0,39 für `ema_spread`). Deshalb
> ist die Funktion **standardmäßig aus**. Sie ist da, damit du es auf deinen Zeiträumen
> selbst prüfen kannst — nicht, weil sie sich bewährt hätte.
**Grenzen der Börsen-API** (gemessen an Binance): Open Interest reicht nur **30 Tage**
zurück, 500 Zeilen je Abruf. Funding Rate reicht über ein Jahr. Längere Backtests deshalb
mit `open_interest: false` fahren.
Ein- und Ausschalten ändert die Anzahl der Merkmale (18 ↔ 22). Ein gespeichertes Modell mit
abweichender Dimension wird verworfen und das Training beginnt von vorn — die Schalter sind
darum als neustartpflichtig markiert.
---
## Erklärungen über ein lokales Sprachmodell
Mit einer laufenden [Ollama](https://ollama.com)-Instanz erscheint im Dashboard der Bereich
**Erklärung**: Auf Knopfdruck fasst ein lokales Modell zusammen, was die Zahlen zeigen —
Merkmalsgewichte, Trefferquote, Risikolage, offene Positionen.
```yaml
llm:
enabled: true
model: llama3.2
```
**Das Modell entscheidet nichts.** Es bekommt den Zustand als Text und gibt Text zurück; es
wird nie aus dem Handels-Loop heraus aufgerufen und beeinflusst keine Order. Bewusst so:
Ein Sprachmodell je Kerze entscheiden zu lassen wäre nicht reproduzierbar, kaum backtestbar
und würde die Nachvollziehbarkeit zerstören, die das lineare Modell heute bietet.
Der System-Prompt untersagt Anlageempfehlungen und Kursprognosen. Beispielausgabe:
> Der Bot befindet sich im simulierten Modus und hat bislang keine Handelsaktivität
> entfaltet […] Da die Trefferquote knapp über dem Zufallswert liegt und keine realen
> Handelsdaten vorliegen, ist die Belastbarkeit des Modells unter realen Marktbedingungen
> noch nicht nachgewiesen.
> **Reasoning-Modelle.** `qwen3`, `deepseek-r1` und Verwandte legen ihre Denkschritte in ein
> eigenes Antwortfeld und verbrauchen dafür das gesamte Token-Budget — die eigentliche
> Antwort bleibt leer. `llm.think: false` (Standard) schaltet das ab; getestet mit
> `qwen3.8:27b`, das damit in rund 20 Sekunden antwortet. Reicht das Budget trotzdem nicht,
> sagt die Fehlermeldung genau das und nennt den Schalter.
**Nicht enthalten:** News- oder Google-Trends-Sentiment als Merkmal. Der Grund ist nicht
technischer Natur — es fehlt eine Quelle mit Point-in-Time-Historie. Ohne die lässt sich ein
solches Merkmal nicht backtesten, und ein unvalidiertes Merkmal in ein Modell zu geben, das
bei 51 % Accuracy steht, macht es eher schlechter. Google Trends kommt zusätzlich mit
täglicher Auflösung und pro Abfrage neu skalierten Werten — für 5-Minuten-Kerzen unbrauchbar.
---
## Risikomanagement
Vor jeder Order greifen mehrere unabhängige Grenzen:
@@ -463,6 +541,7 @@ curl -X POST localhost:8080/control/train/history -H 'Content-Type: application/
| Endpunkt | Methode | Wirkung |
|--------------------------------|---------|--------------------------------------------|
| `/control/explain` | POST | Erklärung des aktuellen Zustands (braucht `llm.enabled`) |
| `/control/config` | GET | Alle Felder mit Wert, Typ, Grenzen und Markierungen |
| `/control/config` | POST | `{"risk.max_open_positions": 5}` — ändern und sichern |
| `/control/config/reset` | POST | `{}` oder `{"paths": [...]}` — Overlay verwerfen |
+39 -1
View File
@@ -95,6 +95,23 @@ strategy:
trend_filter_period: 100 # 0 = Trendfilter aus
min_holding_bars: 3 # Signalausstiege erst danach; Stop/Ziel gelten immer
# Zusatzmerkmale aus dem Terminmarkt. Die Daten stammen vom Perpetual zum jeweiligen
# Spot-Paar (BTC/USDT -> BTC/USDT:USDT); gehandelt wird weiterhin Spot.
#
# In zwei Backtests (5m über 20 Tage, 15m über 28 Tage) brachten sie KEINE Verbesserung
# die Anbindung funktioniert, die Merkmale trugen nichts bei. Deshalb standardmäßig aus.
# Zum Selberprüfen einschalten und gegen einen Lauf ohne vergleichen.
#
# Achtung: Ein- und Ausschalten ändert die Anzahl der Merkmale. Ein gespeichertes Modell
# mit anderer Dimension wird verworfen, das Training beginnt von vorn.
derivatives:
enabled: false
funding_rate: true # was Longs den Shorts zahlen, alle 8h
open_interest: true # offene Kontrakte reicht bei Binance nur 30 Tage zurück
oi_change_bars: 12 # über wie viele Kerzen die OI-Veränderung gemessen wird
max_pages: 40 # Obergrenze für Seitenabrufe je Symbol
min_coverage: 0.9 # darunter gilt die Reihe als unbrauchbar und bleibt neutral
learner:
enabled: true
model_path: /data/models/adaptive.npz
@@ -111,7 +128,12 @@ strategy:
trade_sample_weight: 3.0 # echte Trades zählen dreifach gegenüber Shadow-Labels
# Einstiegssignale sind selten. Zusätzliche Stichproben des Marktzustands verkürzen die
# Aufwärmphase von Wochen auf Tage (0 = aus).
background_sample_every_n_bars: 10
# Der wirksamste Hebel für die Lerngeschwindigkeit: Nur rund 3 % der Beobachtungen
# stammen aus echten Trades, der Rest aus diesen Stichproben. Von 10 auf 5 halbiert
# die Zeit bis zur Einsatzbereitschaft.
# Nicht zu klein wählen benachbarte Kerzen sind stark korreliert, unter dem halben
# label_horizon_bars überlappt praktisch jede Stichprobe mit ihrem Nachbarn.
background_sample_every_n_bars: 5
background_sample_weight: 0.5
# Beim Start ein noch untrainiertes Modell aus der Kurshistorie vorlernen, damit der Bot
# nach Sekunden einsatzbereit ist statt nach Tagen (0 = aus).
@@ -147,6 +169,22 @@ notifications:
notify_on_trade: true
notify_on_risk_halt: true
# ──────────────── Erklärungen über ein lokales Sprachmodell (Ollama) ────────
# Fasst auf Knopfdruck zusammen, was die Zahlen zeigen. Hat KEINEN Einfluss auf den
# Handel es entscheidet nichts, es formuliert nur.
llm:
enabled: false
base_url: http://127.0.0.1:11434
model: llama3.2
timeout_seconds: 120
temperature: 0.2
max_tokens: 700
max_answer_chars: 2000
# Reasoning-Modelle (qwen3, deepseek-r1) verbrauchen sonst das gesamte Token-Budget für
# ihre Denkschritte und liefern eine leere Antwort. Für eine Zustandsbeschreibung
# wird kein Reasoning gebraucht.
think: false
# ───────────────────────────── Backtest-Voreinstellungen ────────────────────
backtest:
bars: 5000
+40 -2
View File
@@ -11,9 +11,11 @@ from .broker import Broker, LiveBroker, PaperBroker
from .config import Config, Mode
from .configstore import ConfigStore
from .data import CcxtDataFeed, DataFeed
from .derivatives import DerivativesProvider
from .engine import TradingEngine
from .exchange import build_exchange, load_market_info
from .features import N_FEATURES
from .features import n_features
from .llm import OllamaClient
from .notify import Notifier
from .portfolio import Portfolio
from .risk import RiskManager
@@ -63,6 +65,8 @@ class Runtime:
notifier: Notifier
server: StatusServer | None
exchange: Any | None
derivatives: DerivativesProvider | None = None
llm: OllamaClient | None = None
async def start_services(self) -> None:
await self.notifier.start()
@@ -73,6 +77,10 @@ class Runtime:
if self.server is not None:
await self.server.close()
await self.notifier.close()
if self.llm is not None:
await self.llm.close()
if self.derivatives is not None:
await self.derivatives.close()
if self.exchange is not None:
await self.exchange.close()
else:
@@ -186,12 +194,38 @@ async def build_runtime(
broker = PaperBroker(paper_config, market_info=market_info)
starting_equity = paper_config.starting_balance
strategy = build_strategy(config.strategy, N_FEATURES, seed=seed, load_model=load_model)
derivatives_provider: DerivativesProvider | None = None
derivatives_exchange = None
if config.strategy.derivatives.enabled:
# Funding Rate und Open Interest gibt es nur am Terminmarkt, also ein zweiter
# Client mit defaultType=future neben dem Spot-Client.
futures_config = config.exchange.model_copy(
update={"options": {**config.exchange.options, "defaultType": "future"}}
)
derivatives_exchange = build_exchange(futures_config, read_only=True)
derivatives_provider = DerivativesProvider(derivatives_exchange, config.strategy.derivatives)
log.info(
"Terminmarktmerkmale aktiv (%s%s) Daten vom Perpetual zu %s",
"Funding Rate" if config.strategy.derivatives.funding_rate else "",
", Open Interest" if config.strategy.derivatives.open_interest else "",
", ".join(config.market.symbols),
)
strategy = build_strategy(
config.strategy,
n_features(config.strategy.derivatives),
seed=seed,
load_model=load_model,
)
learner = getattr(strategy, "learner", None)
if learner is not None and config.mode is Mode.LIVE and config.strategy.learner.freeze_in_live:
learner.frozen = True
log.info("Live-Modus: Online-Lernen eingefroren (freeze_in_live=true)")
llm_client = OllamaClient(config.llm) if config.llm.enabled else None
if llm_client is not None:
log.info("Erklärungen über Ollama aktiv (%s, Modell %s)", config.llm.base_url, config.llm.model)
portfolio = Portfolio(starting_equity=starting_equity, quote_currency=broker.quote_currency)
risk = RiskManager(config.risk)
notifier = Notifier(config.notifications)
@@ -206,6 +240,8 @@ async def build_runtime(
storage=storage,
notifier=notifier,
config_store=config_store,
derivatives=derivatives_provider,
llm=llm_client,
)
server = StatusServer(config.server, engine.status, controller=engine) if serving else None
@@ -230,6 +266,8 @@ async def build_runtime(
notifier=notifier,
server=server,
exchange=exchange,
derivatives=derivatives_provider,
llm=llm_client,
)
+5 -5
View File
@@ -16,7 +16,7 @@ import numpy as np
from .data import format_ts
from .engine import Bar, TradingEngine
from .features import FeatureMatrix, build_feature_matrix
from .features import FeatureMatrix
from .models import Candles, ExitReason
log = logging.getLogger(__name__)
@@ -114,12 +114,12 @@ class BacktestRunner:
self.progress_every = progress_every
self._matrices: dict[str, FeatureMatrix] = {}
def _prepare(self) -> int:
rules = self.engine.config.strategy.rules
async def _prepare(self) -> int:
first_valid = 0
usable: dict[str, Candles] = {}
for symbol, candles in self.series.items():
matrix = build_feature_matrix(candles, rules)
# Über die Engine, damit Terminmarktdaten genauso einfließen wie im Live-Betrieb.
matrix = await self.engine.build_matrix(candles)
if matrix is None:
log.warning(
"%s: nur %d Kerzen zu wenig für die Indikatoren, Symbol wird übersprungen",
@@ -136,7 +136,7 @@ class BacktestRunner:
async def run(self) -> BacktestReport:
start_time = time.perf_counter()
first_valid = self._prepare()
first_valid = await self._prepare()
length = min(len(c) for c in self.series.values())
symbols = list(self.series)
+52 -3
View File
@@ -35,6 +35,10 @@ RESTART_REQUIRED: frozenset[str] = frozenset(
# Wirkt erst beim nächsten Start der aktuelle Handelszustand bleibt, wie er ist.
"trading.autostart",
"strategy.name",
# Ändert die Anzahl der Merkmale und damit die Modelldimension.
"strategy.derivatives.enabled",
"strategy.derivatives.funding_rate",
"strategy.derivatives.open_interest",
"strategy.learner.enabled",
"strategy.learner.model_path",
"strategy.learner.replay_size",
@@ -165,9 +169,15 @@ class LearnerConfig(_Base):
label_horizon_bars: int = Field(default=12, ge=1)
label_target_bps: float = Field(default=30.0, ge=0)
trade_sample_weight: float = Field(default=3.0, gt=0)
# Einstiegssignale sind selten. Zusätzliche Stichproben des Marktzustands beschleunigen
# die Aufwärmphase erheblich (0 = aus).
background_sample_every_n_bars: int = Field(default=10, ge=0)
# Einstiegssignale sind selten nur rund 3 % der Beobachtungen stammen aus echten
# Trades. Regelmäßige Stichproben des Marktzustands sind daher der wirksamste Hebel
# für die Lerngeschwindigkeit (0 = aus).
#
# Achtung beim Verkleinern: Benachbarte Kerzen sind stark korreliert und die
# Label-Fenster überlappen sich. Bei 5 statt 10 verdoppeln sich die Trainingsschritte,
# der Informationsgehalt wächst aber deutlich langsamer. Unter dem halben
# label_horizon_bars überlappt praktisch jede Stichprobe mit ihrem Nachbarn.
background_sample_every_n_bars: int = Field(default=5, ge=0)
background_sample_weight: float = Field(default=0.5, gt=0)
# Beim Start ein noch untrainiertes Modell aus der Kurshistorie vorlernen, statt
# tagelang auf genügend Live-Beobachtungen zu warten (0 = aus).
@@ -176,10 +186,29 @@ class LearnerConfig(_Base):
save_every_n_updates: int = Field(default=50, ge=1)
class DerivativesConfig(_Base):
"""Zusatzmerkmale aus dem Terminmarkt (Funding Rate, Open Interest).
Die Daten stammen vom Perpetual zum jeweiligen Spot-Paar. Gehandelt wird weiterhin Spot.
"""
enabled: bool = False
funding_rate: bool = True
open_interest: bool = True
# Über wie viele Kerzen die Open-Interest-Veränderung gemessen wird.
oi_change_bars: int = Field(default=12, ge=1, le=500)
# Obergrenze für Seitenabrufe je Symbol Binance liefert nur 500 Zeilen je Seite und
# deckt höchstens 30 Tage ab.
max_pages: int = Field(default=40, ge=1, le=500)
# Unter dieser Abdeckung gilt die Reihe als unbrauchbar und wird verworfen.
min_coverage: float = Field(default=0.9, ge=0.0, le=1.0)
class StrategyConfig(_Base):
name: str = "adaptive" # adaptive | rules
rules: RuleConfig = Field(default_factory=RuleConfig)
learner: LearnerConfig = Field(default_factory=LearnerConfig)
derivatives: DerivativesConfig = Field(default_factory=DerivativesConfig)
@field_validator("name")
@classmethod
@@ -230,6 +259,25 @@ class NotificationConfig(_Base):
notify_on_risk_halt: bool = True
class OllamaConfig(_Base):
"""Lokales Sprachmodell zur Erklärung von Entscheidungen.
Nimmt keinerlei Einfluss auf den Handel es formuliert nur, was die Zahlen zeigen.
"""
enabled: bool = False
base_url: str = "http://127.0.0.1:11434"
model: str = "llama3.2"
timeout_seconds: float = Field(default=120.0, gt=0, le=600)
temperature: float = Field(default=0.2, ge=0.0, le=2.0)
max_tokens: int = Field(default=700, ge=32, le=4096)
max_answer_chars: int = Field(default=2000, ge=100, le=20_000)
# Reasoning-Modelle (qwen3, deepseek-r1) legen ihre Denkschritte in ein eigenes Feld und
# verbrauchen dafür das gesamte Token-Budget die eigentliche Antwort bleibt dann leer.
# Für eine Zustandsbeschreibung wird kein Reasoning gebraucht, deshalb standardmäßig aus.
think: bool | None = False
class BacktestConfig(_Base):
start: str | None = None # ISO-8601, z.B. 2024-01-01T00:00:00Z
end: str | None = None
@@ -250,6 +298,7 @@ class Config(_Base):
storage: StorageConfig = Field(default_factory=StorageConfig)
server: ServerConfig = Field(default_factory=ServerConfig)
notifications: NotificationConfig = Field(default_factory=NotificationConfig)
llm: OllamaConfig = Field(default_factory=OllamaConfig)
backtest: BacktestConfig = Field(default_factory=BacktestConfig)
@field_validator("log_level")
+228
View File
@@ -0,0 +1,228 @@
"""Zusatzdaten aus dem Terminmarkt: Funding Rate und Open Interest.
Beides sagt etwas über die Positionierung der Marktteilnehmer, was sich aus reinen
Kerzendaten nicht ablesen lässt:
* **Funding Rate** was Long-Positionen den Short-Positionen zahlen (oder umgekehrt).
Positiv und steigend heißt: Longs sind bereit, für ihre Position zu bezahlen.
* **Open Interest** wie viele Kontrakte offen sind. Zusammen mit der Kursrichtung
unterscheidet das neu aufgebaute Positionen von Glattstellungen.
Gehandelt wird weiterhin Spot; die Kennzahlen stammen vom zugehörigen Perpetual
(``BTC/USDT`` → ``BTC/USDT:USDT``).
Zwei Grenzen der Börsen-API, gemessen an Binance:
* Open Interest reicht nur **30 Tage** zurück, 500 Zeilen je Abruf.
* Funding Rate reicht über ein Jahr zurück, veröffentlicht alle 8 Stunden.
Deshalb sind beide Merkmale einzeln abschaltbar, und fehlende Abdeckung wird gemeldet
statt stillschweigend mit Nullen aufgefüllt.
"""
from __future__ import annotations
import asyncio
import logging
from dataclasses import dataclass
import numpy as np
from .config import DerivativesConfig
from .models import Candles
log = logging.getLogger(__name__)
MAX_PAGE = 500
FUNDING_INTERVAL_MS = 8 * 3_600_000
def perpetual_symbol(spot_symbol: str) -> str:
"""``BTC/USDT`` → ``BTC/USDT:USDT`` die übliche ccxt-Schreibweise für Perpetuals."""
if ":" in spot_symbol:
return spot_symbol
quote = spot_symbol.split("/")[-1]
return f"{spot_symbol}:{quote}"
def forward_fill_to_bars(
bar_timestamps: np.ndarray, source_ts: np.ndarray, source_values: np.ndarray
) -> tuple[np.ndarray, np.ndarray]:
"""Ordnet eine Zeitreihe den Kerzen zu ohne in die Zukunft zu schauen.
Für jede Kerze gilt der letzte Wert, der zum Kerzenzeitpunkt **bereits bekannt** war.
Gibt ``(werte, bekannt)`` zurück; ``bekannt`` markiert Kerzen ohne Vorgängerwert.
"""
bars = np.asarray(bar_timestamps, dtype=np.int64)
values = np.full(bars.size, np.nan, dtype=np.float64)
known = np.zeros(bars.size, dtype=bool)
if source_ts.size == 0:
return values, known
order = np.argsort(source_ts)
src_ts = np.asarray(source_ts, dtype=np.int64)[order]
src_val = np.asarray(source_values, dtype=np.float64)[order]
# searchsorted mit "right" liefert die Anzahl Quellwerte mit ts <= bar_ts.
idx = np.searchsorted(src_ts, bars, side="right") - 1
valid = idx >= 0
values[valid] = src_val[idx[valid]]
known[valid] = True
return values, known
@dataclass(slots=True)
class DerivativeSeries:
"""Auf die Kerzen ausgerichtete Terminmarktdaten."""
symbol: str
timestamp: np.ndarray
funding_rate: np.ndarray # Anteil je 8h-Periode, NaN wenn unbekannt
open_interest: np.ndarray # Kontrakte, NaN wenn unbekannt
funding_coverage: float = 0.0
oi_coverage: float = 0.0
@property
def usable(self) -> bool:
return self.funding_coverage > 0.0 or self.oi_coverage > 0.0
def describe(self) -> str:
return (
f"{self.symbol}: Funding {self.funding_coverage * 100:.0f} %, "
f"Open Interest {self.oi_coverage * 100:.0f} % der Kerzen abgedeckt"
)
class DerivativesProvider:
"""Holt Funding Rate und Open Interest und hält sie je Symbol vor.
Der Zwischenspeicher wächst inkrementell: Im Live-Betrieb wird je neuer Kerze nur
das kurze Stück seit dem letzten bekannten Zeitstempel nachgeladen.
"""
def __init__(self, exchange, config: DerivativesConfig) -> None:
self._exchange = exchange
self.config = config
self._funding: dict[str, tuple[np.ndarray, np.ndarray]] = {}
self._oi: dict[str, tuple[np.ndarray, np.ndarray]] = {}
self.failures: dict[str, str] = {}
async def close(self) -> None:
if self._exchange is not None:
await self._exchange.close()
# ------------------------------------------------------------------ Abruf
async def _fetch_funding(self, symbol: str, since: int, until: int) -> None:
perp = perpetual_symbol(symbol)
known_ts, known_val = self._funding.get(symbol, (np.empty(0, np.int64), np.empty(0)))
cursor = int(known_ts[-1]) + 1 if known_ts.size else since
collected: list[tuple[int, float]] = []
while cursor < until:
try:
rows = await self._exchange.fetch_funding_rate_history(perp, since=cursor, limit=1000)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001 - Zusatzdaten dürfen nie den Bot stoppen
self.failures[f"{symbol}/funding"] = f"{type(exc).__name__}: {exc}"
log.warning("%s: Funding-Rate nicht abrufbar (%s)", perp, exc)
return
if not rows:
break
for row in rows:
rate = row.get("fundingRate")
if row.get("timestamp") is not None and rate is not None:
collected.append((int(row["timestamp"]), float(rate)))
nxt = int(rows[-1]["timestamp"]) + 1
if nxt <= cursor:
break
cursor = nxt
await self._throttle()
if collected:
ts = np.concatenate([known_ts, np.array([c[0] for c in collected], dtype=np.int64)])
val = np.concatenate([known_val, np.array([c[1] for c in collected], dtype=np.float64)])
order = np.argsort(ts)
self._funding[symbol] = (ts[order], val[order])
async def _fetch_open_interest(self, symbol: str, timeframe: str, since: int, until: int) -> None:
perp = perpetual_symbol(symbol)
known_ts, known_val = self._oi.get(symbol, (np.empty(0, np.int64), np.empty(0)))
cursor = int(known_ts[-1]) + 1 if known_ts.size else since
collected: list[tuple[int, float]] = []
pages = 0
while cursor < until and pages < self.config.max_pages:
try:
rows = await self._exchange.fetch_open_interest_history(
perp, timeframe, since=cursor, limit=MAX_PAGE
)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001
# Häufigster Fall: Anfrage älter als das Fenster der Börse (30 Tage).
self.failures[f"{symbol}/open_interest"] = f"{type(exc).__name__}: {exc}"
log.warning("%s: Open Interest ab %d nicht abrufbar (%s)", perp, cursor, exc)
break
if not rows:
break
for row in rows:
amount = row.get("openInterestAmount") or row.get("openInterestValue")
if row.get("timestamp") is not None and amount is not None:
collected.append((int(row["timestamp"]), float(amount)))
nxt = int(rows[-1]["timestamp"]) + 1
if nxt <= cursor:
break
cursor = nxt
pages += 1
await self._throttle()
if collected:
ts = np.concatenate([known_ts, np.array([c[0] for c in collected], dtype=np.int64)])
val = np.concatenate([known_val, np.array([c[1] for c in collected], dtype=np.float64)])
order = np.argsort(ts)
unique = np.concatenate([[True], np.diff(ts[order]) > 0])
self._oi[symbol] = (ts[order][unique], val[order][unique])
async def _throttle(self) -> None:
delay = getattr(self._exchange, "rateLimit", 0) or 0
if delay:
await asyncio.sleep(delay / 1000)
# ------------------------------------------------------------ Bereitstellen
async def series_for(self, candles: Candles) -> DerivativeSeries:
"""Terminmarktdaten passend zu einer Kerzenserie, auf deren Zeitstempel ausgerichtet."""
symbol = candles.symbol
bars = np.asarray(candles.timestamp, dtype=np.int64)
empty = np.full(bars.size, np.nan)
series = DerivativeSeries(symbol=symbol, timestamp=bars, funding_rate=empty.copy(),
open_interest=empty.copy())
if bars.size == 0:
return series
since = int(bars[0]) - FUNDING_INTERVAL_MS # ein Intervall Vorlauf für den ersten Wert
until = int(bars[-1]) + 1
if self.config.funding_rate:
await self._fetch_funding(symbol, since, until)
ts, val = self._funding.get(symbol, (np.empty(0, np.int64), np.empty(0)))
values, known = forward_fill_to_bars(bars, ts, val)
series.funding_rate = values
series.funding_coverage = float(np.mean(known)) if bars.size else 0.0
if self.config.open_interest:
await self._fetch_open_interest(symbol, candles.timeframe, since, until)
ts, val = self._oi.get(symbol, (np.empty(0, np.int64), np.empty(0)))
values, known = forward_fill_to_bars(bars, ts, val)
series.open_interest = values
series.oi_coverage = float(np.mean(known)) if bars.size else 0.0
return series
def snapshot(self) -> dict[str, object]:
return {
"enabled": True,
"funding_symbols": sorted(self._funding),
"oi_symbols": sorted(self._oi),
"failures": dict(self.failures),
}
+78 -4
View File
@@ -20,7 +20,9 @@ from .broker import Broker, InsufficientFunds, OrderRejected, PaperBroker
from .config import Config
from .configstore import ConfigError, ConfigStore
from .data import DataFeed
from .features import FEATURE_NAMES, FeatureSnapshot, build_feature_matrix, required_bars
from .derivatives import DerivativesProvider
from .features import FeatureSnapshot, build_feature_matrix, feature_names, required_bars
from .llm import OllamaClient, build_status_prompt
from .models import Action, Candles, ExitReason, Position, Side, Signal
from .notify import Notifier
from .portfolio import Portfolio
@@ -151,9 +153,14 @@ class TradingEngine:
storage: Storage | NullStorage,
notifier: Notifier | None = None,
config_store: ConfigStore | None = None,
derivatives: DerivativesProvider | None = None,
llm: OllamaClient | None = None,
) -> None:
self.llm = llm
self.config = config
self.config_store = config_store
self.derivatives = derivatives
self._coverage_warned: set[str] = set()
self.broker = broker
self.feed = feed
self.strategy = strategy
@@ -258,7 +265,7 @@ class TradingEngine:
log.warning("%s: Historie nicht abrufbar (%s)", symbol, exc)
job.skipped.append(f"{symbol}: {exc}")
continue
matrix = build_feature_matrix(candles, self.config.strategy.rules)
matrix = await self.build_matrix(candles)
if matrix is None:
log.warning("%s: zu wenig Historie (%d Kerzen)", symbol, len(candles))
job.skipped.append(f"{symbol}: nur {len(candles)} Kerzen")
@@ -462,6 +469,26 @@ class TradingEngine:
self.apply_config(new_config)
return {"accepted": True, "reason": "", **self.config_state()}
# ------------------------------------------------------ Erklärung (LLM)
async def explain(self) -> dict[str, Any]:
"""Den aktuellen Zustand vom Sprachmodell in Worte fassen lassen.
Läuft außerhalb des Handels-Loops und ohne dessen Mutex eine Erklärung darf
den Handel weder blockieren noch beeinflussen.
"""
if not self.config.llm.enabled or self.llm is None:
return {
"ok": False,
"error": "Erklärungen sind deaktiviert (llm.enabled)",
"model": self.config.llm.model,
}
prompt = build_status_prompt(self.status())
result = await self.llm.generate(prompt)
if not result.ok:
log.warning("Erklärung fehlgeschlagen: %s", result.error)
return result.as_dict()
def training_status(self) -> dict[str, Any]:
learner = self.learner
return {
@@ -474,6 +501,44 @@ class TradingEngine:
**self.training.as_dict(),
}
async def build_matrix(self, candles: Candles):
"""Merkmalsmatrix bauen und dabei falls aktiviert Terminmarktdaten einbeziehen."""
derivatives = self.config.strategy.derivatives
series = None
if derivatives.enabled and self.derivatives is not None:
try:
series = await self.derivatives.series_for(candles)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001 - Zusatzdaten dürfen nie den Bot stoppen
log.warning("%s: Terminmarktdaten nicht verfügbar (%s)", candles.symbol, exc)
series = None
else:
if series is not None and not self._coverage_ok(series):
series = None
return build_feature_matrix(candles, self.config.strategy.rules, derivatives, series)
def _coverage_ok(self, series) -> bool:
"""Zu lückenhafte Reihen verwerfen neutrale Spalten sind ehrlicher als Rauschen."""
minimum = self.config.strategy.derivatives.min_coverage
checks = []
if self.config.strategy.derivatives.funding_rate:
checks.append(("Funding", series.funding_coverage))
if self.config.strategy.derivatives.open_interest:
checks.append(("Open Interest", series.oi_coverage))
for label, coverage in checks:
if coverage < minimum:
key = f"{series.symbol}/{label}"
if key not in self._coverage_warned:
self._coverage_warned.add(key)
log.warning(
"%s: %s deckt nur %.0f %% der Kerzen ab (nötig %.0f %%) "
"Terminmarktmerkmale bleiben für dieses Symbol neutral",
series.symbol, label, coverage * 100, minimum * 100,
)
return False
return True
async def _fetch_history(self, symbol: str, bars: int) -> Candles:
"""Längere Historie holen, wenn der Feed das kann sonst das normale Fenster."""
fetch_history = getattr(self.feed, "fetch_history", None)
@@ -543,7 +608,7 @@ class TradingEngine:
continue # noch dieselbe Kerze nichts Neues zu entscheiden
self.last_bar_ts[symbol] = bar.timestamp
matrix = build_feature_matrix(candles, self.config.strategy.rules)
matrix = await self.build_matrix(candles)
snapshot = matrix.snapshot(-1) if matrix is not None else None
if snapshot is None:
log.info(
@@ -793,8 +858,17 @@ class TradingEngine:
"trading": self.trading_control_status(),
"training": self.training_status(),
"feature_weights": (
learner.feature_importance(FEATURE_NAMES) if learner is not None else {}
learner.feature_importance(feature_names(self.config.strategy.derivatives))
if learner is not None
else {}
),
"derivatives": (
self.derivatives.snapshot() if self.derivatives is not None else {"enabled": False}
),
"llm": {
"enabled": self.config.llm.enabled and self.llm is not None,
"model": self.config.llm.model,
},
}
+81 -5
View File
@@ -13,11 +13,11 @@ from dataclasses import dataclass
import numpy as np
from .config import RuleConfig
from .config import DerivativesConfig, RuleConfig
from .indicators import atr, donchian_position, ema, macd, roc, rolling_std, rsi, sma
from .models import Candles
FEATURE_NAMES: tuple[str, ...] = (
BASE_FEATURE_NAMES: tuple[str, ...] = (
"ema_spread", # (EMA_fast - EMA_slow) / Preis [%]
"ema_fast_dist", # (Preis - EMA_fast) / Preis [%]
"trend_dist", # (Preis - EMA_trend) / Preis [%]
@@ -38,11 +38,38 @@ FEATURE_NAMES: tuple[str, ...] = (
"time_cos",
)
N_FEATURES = len(FEATURE_NAMES)
# Zusatzmerkmale aus dem Terminmarkt nur aktiv, wenn strategy.derivatives.enabled.
FUNDING_FEATURE_NAMES: tuple[str, ...] = (
"funding_bps", # Funding Rate der laufenden Periode, in Basispunkten
"funding_trend", # Abweichung vom Mittel der letzten Perioden
)
OI_FEATURE_NAMES: tuple[str, ...] = (
"oi_change", # Veränderung des Open Interest über oi_change_bars [%]
"oi_price_divergence", # OI-Veränderung × Kursrichtung: neue Positionen oder Glattstellung
)
# Rückwärtskompatibler Name für die Basisausstattung.
FEATURE_NAMES = BASE_FEATURE_NAMES
N_FEATURES = len(BASE_FEATURE_NAMES)
MIN_BARS = 140
_CLIP_LIMIT = 8.0
def feature_names(derivatives: DerivativesConfig | None = None) -> tuple[str, ...]:
"""Merkmalsnamen für die gegebene Konfiguration die Anzahl hängt davon ab."""
names = BASE_FEATURE_NAMES
if derivatives is not None and derivatives.enabled:
if derivatives.funding_rate:
names += FUNDING_FEATURE_NAMES
if derivatives.open_interest:
names += OI_FEATURE_NAMES
return names
def n_features(derivatives: DerivativesConfig | None = None) -> int:
return len(feature_names(derivatives))
@dataclass(slots=True)
class FeatureSnapshot:
"""Merkmalsvektor plus Roh-Kennzahlen, die Risiko und Regelwerk zusätzlich brauchen."""
@@ -96,6 +123,7 @@ class FeatureMatrix:
ema_slow: np.ndarray
trend_ema: np.ndarray
first_valid: int # ab hier sind die Zeilen belastbar
names: tuple[str, ...] = BASE_FEATURE_NAMES
def __len__(self) -> int:
return int(self.values.shape[0])
@@ -112,7 +140,7 @@ class FeatureMatrix:
prev = max(idx - 1, 0)
return FeatureSnapshot(
values=self.values[idx].copy(),
names=FEATURE_NAMES,
names=self.names,
index=idx,
price=float(self.price[idx]),
atr=float(self.atr[idx]),
@@ -127,9 +155,54 @@ class FeatureMatrix:
)
def build_feature_matrix(candles: Candles, rules: RuleConfig) -> FeatureMatrix | None:
def _derivative_columns(
close: np.ndarray, series, config: DerivativesConfig
) -> list[np.ndarray]:
"""Spalten aus Funding Rate und Open Interest, in derselben Reihenfolge wie die Namen."""
columns: list[np.ndarray] = []
n = close.size
if config.funding_rate:
rate = np.asarray(getattr(series, "funding_rate", None), dtype=np.float64) \
if series is not None else np.full(n, np.nan)
rate = _clean(rate, 0.0)
# Anteil je 8h → Basispunkte. Typisch ±1 bp, in Extremphasen ±10 bp.
funding_bps = rate * 10_000.0
# Abweichung vom gleitenden Mittel der letzten Perioden: Zuspitzung oder Entspannung.
baseline = _clean(sma(funding_bps, 24), 0.0)
columns += [funding_bps, funding_bps - baseline]
if config.open_interest:
oi = np.asarray(getattr(series, "open_interest", None), dtype=np.float64) \
if series is not None else np.full(n, np.nan)
oi = _clean(oi, 0.0)
lag = min(config.oi_change_bars, max(n - 1, 1))
previous = np.concatenate([np.full(min(lag, n), oi[0] if n else 0.0), oi[:-lag]])[:n]
with np.errstate(divide="ignore", invalid="ignore"):
oi_change = np.where(previous > 0, (oi - previous) / previous * 100.0, 0.0)
prev_close = np.concatenate([np.full(min(lag, n), close[0] if n else 0.0), close[:-lag]])[:n]
with np.errstate(divide="ignore", invalid="ignore"):
price_change = np.where(prev_close > 0, (close - prev_close) / prev_close * 100.0, 0.0)
# Gleiches Vorzeichen = frisches Geld in Richtung des Trends, Gegenzeichen =
# Glattstellungen. Das Produkt fasst beides in einer Zahl zusammen.
divergence = np.sign(oi_change) * np.abs(price_change)
columns += [_clean(oi_change), _clean(divergence)]
return columns
def build_feature_matrix(
candles: Candles,
rules: RuleConfig,
derivatives: DerivativesConfig | None = None,
series=None,
) -> FeatureMatrix | None:
"""Berechnet Indikatoren und Merkmalsvektoren für die gesamte Serie.
``series`` ist eine :class:`~trademind.derivatives.DerivativeSeries` passend zu den
Kerzen; fehlt sie bei aktivierten Terminmarktmerkmalen, werden die Spalten neutral
gefüllt, damit die Modelldimension gleich bleibt.
Gibt ``None`` zurück, wenn die Historie kürzer als ``required_bars`` ist.
"""
n = len(candles)
@@ -201,6 +274,8 @@ def build_feature_matrix(candles: Candles, rules: RuleConfig) -> FeatureMatrix |
np.sin(angle),
np.cos(angle),
]
if derivatives is not None and derivatives.enabled:
columns += _derivative_columns(close, series, derivatives)
values = np.column_stack([_clean(col) for col in columns])
np.clip(values, -_CLIP_LIMIT, _CLIP_LIMIT, out=values)
@@ -214,6 +289,7 @@ def build_feature_matrix(candles: Candles, rules: RuleConfig) -> FeatureMatrix |
ema_slow=_clean(ema_slow, close),
trend_ema=_clean(trend_ema, close),
first_valid=need - 1,
names=feature_names(derivatives),
)
+265
View File
@@ -0,0 +1,265 @@
"""Ollama-Anbindung: erklärt Entscheidungen in Worten.
Bewusst **nicht** als Entscheider. Ein Sprachmodell je Kerze über Käufe entscheiden zu
lassen wäre nicht reproduzierbar, kaum backtestbar und in numerischer Zeitreihenvorhersage
schwach und es würde die Nachvollziehbarkeit zerstören, die das lineare Modell heute
bietet (``/status`` zeigt jedes einzelne Gewicht).
Wozu es taugt: aus Zahlen einen lesbaren Satz machen. Der Bot liefert die tatsächlichen
Entscheidungsgrundlagen Regelbegründung, Modellkonfidenz, Merkmalsgewichte, Risikolage
und das Modell formuliert daraus eine Erklärung. Es bekommt keinerlei Einfluss auf den
Handel und wird nie aus dem Handels-Loop heraus aufgerufen.
"""
from __future__ import annotations
import asyncio
import logging
from dataclasses import dataclass
from typing import Any
import aiohttp
from .config import OllamaConfig
log = logging.getLogger(__name__)
SYSTEM_PROMPT = (
"Du erklärst die Entscheidungen eines Krypto-Trading-Bots auf Deutsch. "
"Beschreibe ausschließlich, was die übergebenen Zahlen zeigen. "
"Gib keine Anlageempfehlung, keine Kursprognose und keine Einschätzung, ob jemand "
"kaufen oder verkaufen sollte. Erfinde keine Zahlen, die nicht dastehen. "
"Wenn die Daten für eine Aussage nicht ausreichen, sage das. "
"Antworte in höchstens drei kurzen Absätzen, sachlich und ohne Aufzählungszeichen."
)
@dataclass(slots=True)
class LlmResult:
ok: bool
text: str = ""
error: str = ""
model: str = ""
duration_seconds: float = 0.0
def as_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"text": self.text,
"error": self.error,
"model": self.model,
"duration_seconds": round(self.duration_seconds, 2),
}
class OllamaClient:
"""Dünner Client für ``/api/generate``. Fehler werden gemeldet, nie geworfen."""
def __init__(self, config: OllamaConfig) -> None:
self.config = config
self._session: aiohttp.ClientSession | None = None
async def _ensure_session(self) -> aiohttp.ClientSession:
if self._session is None or self._session.closed:
timeout = aiohttp.ClientTimeout(total=self.config.timeout_seconds)
self._session = aiohttp.ClientSession(timeout=timeout)
return self._session
async def close(self) -> None:
if self._session is not None and not self._session.closed:
await self._session.close()
self._session = None
async def available_models(self) -> list[str]:
session = await self._ensure_session()
try:
async with session.get(f"{self.config.base_url.rstrip('/')}/api/tags") as response:
if response.status != 200:
return []
payload = await response.json()
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001 - Erreichbarkeit ist optional
log.debug("Ollama nicht erreichbar: %s", exc)
return []
return [m.get("name", "") for m in payload.get("models", []) if m.get("name")]
async def generate(self, prompt: str, system: str = SYSTEM_PROMPT) -> LlmResult:
loop = asyncio.get_running_loop()
started = loop.time()
session = await self._ensure_session()
body: dict[str, Any] = {
"model": self.config.model,
"prompt": prompt,
"system": system,
"stream": False,
"options": {
"temperature": self.config.temperature,
"num_predict": self.config.max_tokens,
},
}
if self.config.think is not None:
body["think"] = self.config.think
url = f"{self.config.base_url.rstrip('/')}/api/generate"
try:
async with session.post(url, json=body) as response:
if response.status == 400 and "think" in body:
# Ältere Ollama-Versionen kennen das Feld nicht ohne es erneut versuchen.
body.pop("think")
async with session.post(url, json=body) as retry:
if retry.status != 200:
detail = (await retry.text())[:200]
return LlmResult(False, error=f"HTTP {retry.status}: {detail}",
model=self.config.model,
duration_seconds=loop.time() - started)
payload = await retry.json()
elif response.status != 200:
detail = (await response.text())[:200]
return LlmResult(False, error=f"HTTP {response.status}: {detail}",
model=self.config.model,
duration_seconds=loop.time() - started)
else:
payload = await response.json()
except TimeoutError:
return LlmResult(
False,
error=f"Zeitüberschreitung nach {self.config.timeout_seconds:g}s "
"größere Modelle brauchen länger, llm.timeout_seconds erhöhen",
model=self.config.model,
duration_seconds=loop.time() - started,
)
except aiohttp.ClientError as exc:
return LlmResult(
False,
error=f"Ollama unter {self.config.base_url} nicht erreichbar ({exc})",
model=self.config.model,
duration_seconds=loop.time() - started,
)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001 - eine Erklärung darf nie den Bot stören
return LlmResult(False, error=f"{type(exc).__name__}: {exc}",
model=self.config.model, duration_seconds=loop.time() - started)
text = str(payload.get("response", "")).strip()
# Manche Modelle betten den Denkprozess als <think>-Block in die Antwort ein.
text = _strip_thinking(text)
if not text:
return LlmResult(
False,
error=_diagnose_empty(payload, self.config),
model=self.config.model,
duration_seconds=loop.time() - started,
)
return LlmResult(True, text=text[: self.config.max_answer_chars],
model=self.config.model, duration_seconds=loop.time() - started)
def _diagnose_empty(payload: dict[str, Any], config: OllamaConfig) -> str:
"""Sagt, *warum* nichts zurückkam „leere Antwort" allein hilft niemandem weiter."""
thinking = str(payload.get("thinking") or "")
reason = payload.get("done_reason")
if thinking and reason == "length":
return (
f"Das Modell hat alle {config.max_tokens} Tokens für seine Denkschritte verbraucht "
f"({len(thinking)} Zeichen im Feld 'thinking'), ohne eine Antwort zu formulieren. "
"llm.think auf false setzen oder llm.max_tokens erhöhen."
)
if thinking:
return (
"Das Modell hat nur Denkschritte geliefert, keine Antwort. "
"llm.think auf false setzen."
)
if reason == "length":
return f"Antwort nach {config.max_tokens} Tokens abgeschnitten llm.max_tokens erhöhen."
return f"Leere Antwort vom Modell (done_reason={reason!r})"
def _strip_thinking(text: str) -> str:
while "<think>" in text and "</think>" in text:
start = text.index("<think>")
end = text.index("</think>") + len("</think>")
text = (text[:start] + text[end:]).strip()
return text
def _fmt(value: Any, digits: int = 2, suffix: str = "") -> str:
if value is None:
return "unbekannt"
if isinstance(value, bool):
return "ja" if value else "nein"
if isinstance(value, (int, float)):
return f"{value:.{digits}f}{suffix}"
return str(value)
def build_status_prompt(status: dict[str, Any], top_weights: int = 8) -> str:
"""Baut aus dem Statusabbild einen Prompt mit den tatsächlichen Entscheidungsgrundlagen."""
portfolio = status.get("portfolio", {})
strategy = status.get("strategy", {})
learner = strategy.get("learner", {}) if isinstance(strategy, dict) else {}
risk = status.get("risk", {})
trading = status.get("trading", {})
weights = status.get("feature_weights", {}) or {}
positions = status.get("positions", []) or []
trades = (status.get("recent_trades", []) or [])[-5:]
ranked = sorted(weights.items(), key=lambda kv: abs(kv[1]), reverse=True)[:top_weights]
lines = [
"Zustand eines Krypto-Trading-Bots. Erkläre in eigenen Worten, was er gerade tut,",
"worauf sein Modell derzeit achtet und wie belastbar das ist.",
"",
f"Betriebsart: {status.get('mode')} ({'simuliert' if trading.get('simulated') else 'echtes Geld'}), "
f"Handel {'aktiv' if trading.get('active') else 'pausiert'}",
f"Börse: {status.get('exchange')}, Symbole: {', '.join(status.get('symbols', []))}, "
f"Kerzenlänge {status.get('timeframe')}",
"",
"PORTFOLIO",
f" Equity {_fmt(portfolio.get('equity'))} {status.get('quote_currency', '')}, "
f"Rendite {_fmt(portfolio.get('total_return_pct'))} %",
f" Abgeschlossene Trades {portfolio.get('trades', 0)}, "
f"Trefferquote {_fmt((portfolio.get('win_rate') or 0) * 100, 1)} %, "
f"Profit-Faktor {_fmt(portfolio.get('profit_factor'))}",
f" Maximaler Drawdown {_fmt(portfolio.get('max_drawdown_pct'))} %, "
f"offene Positionen {portfolio.get('open_positions', 0)}",
"",
"LERNMODELL (logistische Regression über normierte Merkmale)",
f" Beobachtungen {learner.get('samples_seen', 0)}, davon aus echten Trades "
f"{learner.get('trade_samples', 0)}",
f" Trefferquote der Vorhersagen {_fmt((learner.get('online_accuracy') or 0) * 100, 1)} % "
f"(50 % entspricht Raten)",
f" Von {strategy.get('candidates_seen', 0)} Signalen des Regelwerks wurden "
f"{strategy.get('candidates_accepted', 0)} durchgelassen",
"",
"STÄRKSTE MERKMALSGEWICHTE (positiv = spricht für einen Einstieg)",
]
lines += [f" {name}: {value:+.3f}" for name, value in ranked] or [" noch keine"]
lines += ["", "RISIKO", f" Notbremse aktiv: {_fmt(risk.get('halted'))}"]
if risk.get("halt_reason"):
lines.append(f" Grund: {risk['halt_reason']}")
lines.append(
f" Grenzen: höchstens {risk.get('max_open_positions')} Positionen, "
f"{_fmt((risk.get('max_position_pct') or 0) * 100, 0)} % der Equity je Position"
)
if positions:
lines += ["", "OFFENE POSITIONEN"]
for p in positions:
lines.append(
f" {p.get('symbol')}: Einstieg {_fmt(p.get('entry_price'), 6)}, "
f"aktuell {_fmt(p.get('mark_price'), 6)}, "
f"P/L {_fmt(p.get('unrealized_pct'))} %, seit {p.get('bars_held')} Kerzen, "
f"Modellkonfidenz beim Einstieg {_fmt(p.get('confidence'))}"
)
if trades:
lines += ["", "LETZTE TRADES"]
for t in trades:
lines.append(
f" {t.get('symbol')}: {_fmt((t.get('pnl_pct') or 0) * 100)} %, "
f"Ausstieg wegen {t.get('exit_reason')}, {t.get('bars_held')} Kerzen gehalten"
)
return "\n".join(lines)
+45
View File
@@ -31,6 +31,8 @@ class Controller(Protocol):
def trading_control_status(self) -> dict[str, Any]: ...
async def explain(self) -> dict[str, Any]: ...
def config_state(self) -> dict[str, Any]: ...
def update_config(self, patch: dict[str, Any]) -> dict[str, Any]: ...
@@ -130,6 +132,20 @@ _DASHBOARD = r"""<!doctype html>
</div>
</div>
</section>
<section id="explain-section" hidden>
<h2>Erklärung</h2>
<div class="panel">
<div class="row">
<button id="explain-btn">Lage erklären lassen</button>
<span class="hint" id="explain-hint" style="margin:0">
Ein lokales Sprachmodell fasst zusammen, was die Zahlen zeigen. Es hat keinerlei
Einfluss auf den Handel.
</span>
</div>
<div id="explain-out" hidden style="margin-top:12px; white-space:pre-wrap; line-height:1.55"></div>
</div>
</section>
<section id="config-section" hidden>
<h2>Konfiguration</h2>
<div class="panel">
@@ -197,6 +213,7 @@ async function refresh() {
// Der Abschnitt muss aus dem Statuslauf heraus sichtbar werden sein Aufklapp-Knopf
// sitzt darin, er könnte sich sonst nie selbst einblenden.
$("config-section").hidden = !(s.trading || {}).control_enabled;
$("explain-section").hidden = !((s.trading || {}).control_enabled && (s.llm || {}).enabled);
} catch (e) { document.getElementById("sub").textContent = "Status nicht erreichbar: " + e; }
}
@@ -323,6 +340,25 @@ $("token-btn").addEventListener("click", () => {
refresh();
});
// -------------------------------------------------------------- Erklärung
$("explain-btn").addEventListener("click", async () => {
const btn = $("explain-btn"), out = $("explain-out");
btn.disabled = true;
setHint($("explain-hint"), "Das Modell denkt nach das dauert je nach Modellgröße etwas …");
try {
const r = await fetch("control/explain", { method: "POST", headers: headers() });
const d = await r.json().catch(() => ({}));
if (d.ok) {
out.hidden = false;
out.textContent = d.text;
setHint($("explain-hint"), `${d.model}, ${d.duration_seconds}s`);
} else {
setHint($("explain-hint"), d.error || `HTTP ${r.status}`, true);
}
} catch (e) { setHint($("explain-hint"), e.message, true); }
btn.disabled = false;
});
// ----------------------------------------------------------- Konfiguration
let cfgOpen = false, cfgLoaded = null;
@@ -489,6 +525,7 @@ class StatusServer:
web.post("/control/train/live", self._train_live),
web.get("/control/trading", self._trading_state),
web.post("/control/trading", self._set_trading),
web.post("/control/explain", self._explain),
web.get("/control/config", self._config_state),
web.post("/control/config", self._update_config),
web.post("/control/config/reset", self._reset_config),
@@ -631,6 +668,14 @@ class StatusServer:
status = 428 if result.get("requires_confirmation") and payload["enabled"] else 409
return web.json_response(result, status=status, dumps=_dumps)
async def _explain(self, request: web.Request) -> web.Response:
denied = await self._guard(request)
if denied is not None:
return denied
assert self._controller is not None
result = await self._controller.explain()
return web.json_response(result, status=200 if result.get("ok") else 503, dumps=_dumps)
async def _config_state(self, request: web.Request) -> web.Response:
denied = await self._guard(request)
if denied is not None:
+262
View File
@@ -0,0 +1,262 @@
"""Terminmarktdaten: Zuordnung ohne Blick in die Zukunft, Merkmale, Ausfallverhalten."""
from __future__ import annotations
import numpy as np
import pytest
from trademind.config import Config, DerivativesConfig
from trademind.derivatives import (
DerivativeSeries,
DerivativesProvider,
forward_fill_to_bars,
perpetual_symbol,
)
from trademind.features import (
BASE_FEATURE_NAMES,
build_feature_matrix,
feature_names,
n_features,
)
from trademind.models import Candles
from .conftest import make_candles
BAR = 300_000
# ------------------------------------------------------------------ Symbolik
@pytest.mark.parametrize(
("spot", "perp"),
[
("BTC/USDT", "BTC/USDT:USDT"),
("ETH/USDC", "ETH/USDC:USDC"),
("SOL/USDT:USDT", "SOL/USDT:USDT"), # schon ein Perpetual
],
)
def test_perpetual_symbol(spot: str, perp: str):
assert perpetual_symbol(spot) == perp
# ------------------------------------------------- Zuordnung ohne Lookahead
def test_forward_fill_uses_only_past_values():
bars = np.array([1000, 2000, 3000, 4000], dtype=np.int64)
src_ts = np.array([1500, 3500], dtype=np.int64)
src_val = np.array([10.0, 20.0])
values, known = forward_fill_to_bars(bars, src_ts, src_val)
assert not known[0], "vor dem ersten Quellwert darf nichts bekannt sein"
assert values[1] == 10.0 # 1500 <= 2000
assert values[2] == 10.0 # 3500 liegt in der Zukunft von 3000
assert values[3] == 20.0 # 3500 <= 4000
assert list(known) == [False, True, True, True]
def test_value_exactly_on_the_bar_counts_as_known():
values, known = forward_fill_to_bars(
np.array([2000], dtype=np.int64), np.array([2000], dtype=np.int64), np.array([7.0])
)
assert known[0] and values[0] == 7.0
def test_unsorted_source_is_handled():
values, _ = forward_fill_to_bars(
np.array([5000], dtype=np.int64),
np.array([3000, 1000, 2000], dtype=np.int64),
np.array([30.0, 10.0, 20.0]),
)
assert values[0] == 30.0, "der jüngste Wert vor der Kerze zählt"
def test_empty_source_yields_nothing_known():
values, known = forward_fill_to_bars(
np.arange(3, dtype=np.int64), np.empty(0, np.int64), np.empty(0)
)
assert not known.any()
assert np.isnan(values).all()
# --------------------------------------------------------- Merkmalsanzahl
def test_feature_count_depends_on_configuration():
assert n_features(None) == len(BASE_FEATURE_NAMES) == 18
assert n_features(DerivativesConfig(enabled=False)) == 18
assert n_features(DerivativesConfig(enabled=True)) == 22
assert n_features(DerivativesConfig(enabled=True, open_interest=False)) == 20
assert n_features(DerivativesConfig(enabled=True, funding_rate=False)) == 20
def test_feature_names_are_unique_and_ordered():
names = feature_names(DerivativesConfig(enabled=True))
assert names[:18] == BASE_FEATURE_NAMES
assert len(set(names)) == len(names)
assert names[18:] == ("funding_bps", "funding_trend", "oi_change", "oi_price_divergence")
# ------------------------------------------------------ Merkmalsberechnung
def series_for(candles, funding: float = 0.0001, oi_growth: float = 0.0) -> DerivativeSeries:
n = len(candles)
oi = 100_000.0 * (1.0 + oi_growth * np.arange(n) / max(n - 1, 1))
return DerivativeSeries(
symbol=candles.symbol,
timestamp=candles.timestamp,
funding_rate=np.full(n, funding),
open_interest=oi,
funding_coverage=1.0,
oi_coverage=1.0,
)
def test_matrix_gains_columns_when_enabled(candles, rules):
plain = build_feature_matrix(candles, rules)
enriched = build_feature_matrix(
candles, rules, DerivativesConfig(enabled=True), series_for(candles)
)
assert plain.values.shape[1] == 18
assert enriched.values.shape[1] == 22
assert np.allclose(plain.values, enriched.values[:, :18]), "Basismerkmale dürfen sich nicht ändern"
def test_funding_is_converted_to_basis_points(candles, rules):
matrix = build_feature_matrix(
candles, rules, DerivativesConfig(enabled=True), series_for(candles, funding=0.0003)
)
snapshot = matrix.snapshot(-1).as_dict()
assert snapshot["funding_bps"] == pytest.approx(3.0) # 0,03 % = 3 bps
assert snapshot["funding_trend"] == pytest.approx(0.0, abs=1e-9) # konstant, kein Trend
def test_rising_open_interest_shows_up_as_positive_change(candles, rules):
matrix = build_feature_matrix(
candles, rules, DerivativesConfig(enabled=True), series_for(candles, oi_growth=0.5)
)
assert matrix.snapshot(-1).as_dict()["oi_change"] > 0
def test_divergence_sign_follows_open_interest(rules):
up = make_candles(n=400, trend=0.002, noise=0.0002, seed=3)
rising = build_feature_matrix(
up, rules, DerivativesConfig(enabled=True), series_for(up, oi_growth=0.5)
).snapshot(-1).as_dict()
falling = build_feature_matrix(
up, rules, DerivativesConfig(enabled=True), series_for(up, oi_growth=-0.3)
).snapshot(-1).as_dict()
# Steigender Kurs mit steigendem OI = neue Positionen, mit fallendem OI = Glattstellung.
assert rising["oi_price_divergence"] > 0
assert falling["oi_price_divergence"] < 0
def test_missing_series_keeps_the_dimension_stable(candles, rules):
"""Fällt die Datenquelle aus, bleiben die Spalten erhalten neutral gefüllt."""
matrix = build_feature_matrix(candles, rules, DerivativesConfig(enabled=True), None)
assert matrix.values.shape[1] == 22
assert np.isfinite(matrix.values).all()
snapshot = matrix.snapshot(-1).as_dict()
assert snapshot["funding_bps"] == 0.0
assert snapshot["oi_change"] == 0.0
def test_values_stay_within_the_clip_limit(candles, rules):
"""Auch absurde Terminmarktwerte dürfen die Normierung nicht sprengen."""
n = len(candles)
extreme = DerivativeSeries(
symbol=candles.symbol,
timestamp=candles.timestamp,
funding_rate=np.full(n, 0.75), # 7500 bps
open_interest=np.geomspace(1.0, 1e9, n),
funding_coverage=1.0,
oi_coverage=1.0,
)
matrix = build_feature_matrix(candles, rules, DerivativesConfig(enabled=True), extreme)
assert np.abs(matrix.values).max() <= 8.0
assert np.isfinite(matrix.values).all()
# ---------------------------------------------------------- Ausfallverhalten
class FlakyExchange:
"""Börse, die für Funding funktioniert und bei Open Interest scheitert."""
rateLimit = 0
def __init__(self, funding_rows=None):
self.funding_rows = funding_rows or []
self.oi_calls = 0
async def fetch_funding_rate_history(self, symbol, since=None, limit=None):
rows = [r for r in self.funding_rows if since is None or r["timestamp"] >= since]
return rows[:limit] if limit else rows
async def fetch_open_interest_history(self, symbol, timeframe, since=None, limit=None):
self.oi_calls += 1
raise RuntimeError("startTime is invalid")
async def close(self):
return None
async def test_failing_source_is_recorded_not_raised(candles):
start = int(candles.timestamp[0])
rows = [{"timestamp": start - 3_600_000, "fundingRate": 0.0002}]
provider = DerivativesProvider(FlakyExchange(rows), DerivativesConfig(enabled=True))
series = await provider.series_for(candles)
assert series.funding_coverage == 1.0
assert series.oi_coverage == 0.0
assert any("open_interest" in key for key in provider.failures)
assert provider.snapshot()["failures"]
async def test_provider_returns_empty_series_for_empty_candles():
blank = Candles.from_rows("BTC/USDT", "5m", [])
provider = DerivativesProvider(FlakyExchange(), DerivativesConfig(enabled=True))
series = await provider.series_for(blank)
assert series.timestamp.size == 0
assert not series.usable
async def test_incremental_fetch_does_not_refetch_everything(candles):
start = int(candles.timestamp[0])
rows = [{"timestamp": start - 3_600_000 + i * 8 * 3_600_000, "fundingRate": 0.0001}
for i in range(4)]
class Counting(FlakyExchange):
def __init__(self, rows):
super().__init__(rows)
self.funding_calls = 0
async def fetch_funding_rate_history(self, symbol, since=None, limit=None):
self.funding_calls += 1
return await super().fetch_funding_rate_history(symbol, since, limit)
exchange = Counting(rows)
provider = DerivativesProvider(exchange, DerivativesConfig(enabled=True, open_interest=False))
await provider.series_for(candles)
first = exchange.funding_calls
await provider.series_for(candles)
assert exchange.funding_calls - first <= 1, "der zweite Lauf darf nur nachladen"
# -------------------------------------------------------------- Konfiguration
def test_derivatives_switches_require_a_restart():
from trademind.config import requires_restart
for path in ("strategy.derivatives.enabled", "strategy.derivatives.funding_rate",
"strategy.derivatives.open_interest"):
assert requires_restart(path), f"{path} ändert die Modelldimension"
def test_derivatives_are_off_by_default():
assert Config().strategy.derivatives.enabled is False
+8 -3
View File
@@ -95,10 +95,15 @@ async def test_cash_and_equity_stay_consistent(base_config):
async def test_stop_loss_bounds_the_worst_trade(base_config):
# Bewusst ohne Lernmodell: Geprüft wird die Stop-Logik, nicht welche Signale das
# Modell gerade durchlässt. Mit "adaptive" hinge das Ergebnis daran, wie weit das
# Modell aufgewärmt ist das hat mit Stops nichts zu tun.
config = Config.model_validate(
{**base_config.model_dump(), "risk": {**base_config.risk.model_dump(),
"stop_loss_atr_mult": 1.0,
"take_profit_atr_mult": 10.0}}
{**base_config.model_dump(),
"strategy": {"name": "rules"},
"risk": {**base_config.risk.model_dump(),
"stop_loss_atr_mult": 1.0,
"take_profit_atr_mult": 10.0}}
)
engine = build_engine(config)
await engine.prepare()
+221
View File
@@ -0,0 +1,221 @@
"""Ollama-Anbindung: Antwortverarbeitung, Fehlerdiagnose, Prompt-Aufbau."""
from __future__ import annotations
import asyncio
import pytest
from aiohttp import web
from trademind.config import OllamaConfig
from trademind.llm import SYSTEM_PROMPT, OllamaClient, build_status_prompt
STATUS = {
"mode": "paper",
"exchange": "binance",
"symbols": ["BTC/USDT", "ETH/USDT"],
"timeframe": "5m",
"quote_currency": "USDT",
"portfolio": {"equity": 10123.45, "total_return_pct": 1.23, "trades": 7, "win_rate": 0.42,
"profit_factor": 1.1, "max_drawdown_pct": 2.5, "open_positions": 1},
"strategy": {"candidates_seen": 40, "candidates_accepted": 9,
"learner": {"samples_seen": 900, "trade_samples": 7, "online_accuracy": 0.53}},
"risk": {"halted": False, "halt_reason": "", "max_open_positions": 3, "max_position_pct": 0.2},
"trading": {"active": True, "simulated": True},
"feature_weights": {"ema_spread": 0.39, "trend_dist": -0.29, "rsi_norm": 0.01},
"positions": [{"symbol": "BTC/USDT", "entry_price": 70000.0, "mark_price": 70500.0,
"unrealized_pct": 0.71, "bars_held": 4, "confidence": 0.62}],
"recent_trades": [{"symbol": "ETH/USDT", "pnl_pct": -0.004, "exit_reason": "stop_loss",
"bars_held": 9}],
}
# ------------------------------------------------------------------- Prompt
def test_prompt_contains_the_real_numbers():
prompt = build_status_prompt(STATUS)
for needle in ("paper", "binance", "BTC/USDT", "10123.45", "ema_spread", "stop_loss"):
assert needle in prompt, f"{needle} fehlt im Prompt"
def test_prompt_ranks_weights_by_magnitude():
prompt = build_status_prompt(STATUS, top_weights=2)
assert "ema_spread" in prompt and "trend_dist" in prompt
assert "rsi_norm" not in prompt, "das schwächste Gewicht sollte wegfallen"
def test_prompt_survives_a_bare_status():
prompt = build_status_prompt({})
assert "PORTFOLIO" in prompt and "LERNMODELL" in prompt
def test_system_prompt_forbids_advice():
for needle in ("keine Anlageempfehlung", "keine Kursprognose", "Erfinde keine Zahlen"):
assert needle in SYSTEM_PROMPT
# ------------------------------------------------------------ Falsches Ollama
def fake_ollama(handler):
app = web.Application()
app.router.add_post("/api/generate", handler)
app.router.add_get("/api/tags", lambda _: web.json_response({"models": [{"name": "testmodell"}]}))
return app
@pytest.fixture
async def client_for(aiohttp_server):
async def _make(handler, **overrides) -> OllamaClient:
server = await aiohttp_server(fake_ollama(handler))
options = {"model": "testmodell", "timeout_seconds": 5, **overrides}
config = OllamaConfig(
enabled=True, base_url=str(server.make_url("/")).rstrip("/"), **options
)
return OllamaClient(config)
return _make
async def test_plain_answer_is_returned(client_for):
async def handler(request):
assert (await request.json())["stream"] is False
return web.json_response({"response": "Alles ruhig.", "done_reason": "stop"})
client = await client_for(handler)
result = await client.generate("frage")
assert result.ok and result.text == "Alles ruhig."
assert result.model == "testmodell"
await client.close()
async def test_thinking_block_is_stripped(client_for):
async def handler(_):
return web.json_response(
{"response": "<think>erst überlegen</think>Das Ergebnis.", "done_reason": "stop"}
)
client = await client_for(handler)
result = await client.generate("frage")
assert result.ok and result.text == "Das Ergebnis."
await client.close()
async def test_reasoning_model_without_answer_is_diagnosed(client_for):
"""Der reale Fall: qwen3 verbraucht das Token-Budget für 'thinking'."""
async def handler(_):
return web.json_response({"response": "", "thinking": "x" * 1800, "done_reason": "length"})
client = await client_for(handler)
result = await client.generate("frage")
assert not result.ok
assert "Denkschritte" in result.error
assert "llm.think" in result.error, "die Meldung muss den Ausweg nennen"
await client.close()
async def test_truncated_answer_is_diagnosed(client_for):
async def handler(_):
return web.json_response({"response": "", "done_reason": "length"})
client = await client_for(handler)
result = await client.generate("frage")
assert not result.ok and "max_tokens" in result.error
await client.close()
async def test_think_flag_is_sent(client_for):
seen = {}
async def handler(request):
seen.update(await request.json())
return web.json_response({"response": "ok", "done_reason": "stop"})
client = await client_for(handler)
await client.generate("frage")
assert seen["think"] is False, "Reasoning ist standardmäßig aus"
await client.close()
async def test_old_ollama_without_think_field_still_works(client_for):
"""Ältere Versionen lehnen das Feld ab dann ohne es erneut versuchen."""
calls = []
async def handler(request):
body = await request.json()
calls.append("think" in body)
if "think" in body:
return web.json_response({"error": "unknown field think"}, status=400)
return web.json_response({"response": "Klappt doch.", "done_reason": "stop"})
client = await client_for(handler)
result = await client.generate("frage")
assert result.ok and result.text == "Klappt doch."
assert calls == [True, False]
await client.close()
async def test_http_error_is_reported(client_for):
async def handler(_):
return web.json_response({"error": "model not found"}, status=404)
client = await client_for(handler)
result = await client.generate("frage")
assert not result.ok and "404" in result.error
await client.close()
async def test_timeout_is_reported_with_a_hint(client_for):
async def handler(_):
await asyncio.sleep(2)
return web.json_response({"response": "zu spät"})
client = await client_for(handler, timeout_seconds=0.2)
result = await client.generate("frage")
assert not result.ok
assert "Zeitüberschreitung" in result.error and "timeout_seconds" in result.error
await client.close()
async def test_unreachable_server_is_reported(aiohttp_server):
"""Server starten, wieder beenden, dann anfragen die Verbindung wird abgelehnt."""
async def handler(_):
return web.json_response({"response": "nie erreicht"})
server = await aiohttp_server(fake_ollama(handler))
url = str(server.make_url("/")).rstrip("/")
await server.close()
client = OllamaClient(OllamaConfig(enabled=True, base_url=url, timeout_seconds=5))
result = await client.generate("frage")
assert not result.ok
assert "nicht erreichbar" in result.error and url in result.error
await client.close()
async def test_answer_is_capped(client_for):
async def handler(_):
return web.json_response({"response": "y" * 5000, "done_reason": "stop"})
client = await client_for(handler, max_answer_chars=100)
result = await client.generate("frage")
assert result.ok and len(result.text) == 100
await client.close()
async def test_model_listing(client_for):
async def handler(_):
return web.json_response({"response": "ok"})
client = await client_for(handler)
assert await client.available_models() == ["testmodell"]
await client.close()
def test_llm_is_off_by_default():
from trademind.config import Config
assert Config().llm.enabled is False