Initial commit: Oil Trading Bot (MT5, WTI)

Headless FastAPI-Backend (server.py + core/engine.py) mit Mobile-PWA (web/),
Strategie-/Backtest-Suite und Doku. Secrets, DB, Logs und Laufzeit-State sind
via .gitignore ausgeschlossen; Config-Vorlage: oil_widget_config.ini.example.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Axel Hocks
2026-07-24 08:29:23 +02:00
co-authored by Claude Opus 4.8
commit 75d28827e8
104 changed files with 21059 additions and 0 deletions
+394
View File
@@ -0,0 +1,394 @@
#!/usr/bin/env python3
"""
server.py — Web-Backend für den Oil Trading Agent
===================================================
FastAPI-Server, der die headless `TradingEngine` (core/engine.py) fährt und
ihren Zustand per REST bereitstellt. Grundlage für die Mobile-Web-Oberfläche.
Läuft auf demselben Windows-Rechner wie das MT5-Terminal. Das Tkinter-Widget
ist NICHT nötig — Server und Widget teilen sich nur die core/-Logik. Betreibe
entweder das Widget ODER den Server gegen dasselbe MT5-Terminal.
Start:
pip install fastapi "uvicorn[standard]"
python server.py
→ http://localhost:8000/api/snapshot
Aktuell (Schritt 1): nur Read-only-Endpoints.
GET / → Status-Kurzinfo
GET /api/health → läuft die Engine, ist MT5 verbunden?
GET /api/snapshot → kompletter Live-Zustand (Markt/Position/Analyse/Agent)
"""
from __future__ import annotations
import asyncio
import contextlib
import secrets
import sys
from pathlib import Path
try:
import uvicorn
from fastapi import (FastAPI, WebSocket, WebSocketDisconnect,
Header, HTTPException, Body)
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import JSONResponse
from fastapi.staticfiles import StaticFiles
except ImportError:
print('pip install fastapi "uvicorn[standard]"')
sys.exit(1)
from core.config import save_config
from core.engine import TradingEngine
from core.logger import get_logger
log = get_logger("server")
WEB_DIR = Path(__file__).parent / "web"
WS_PUSH_S = 1.0 # Push-Intervall des Live-Snapshots
# Im Heim-VPN (WireGuard) erreichen Handy/Frontend den Server direkt; CORS
# offen lassen ist hier vertretbar. Bei öffentlichem Betrieb einschränken.
HOST = "0.0.0.0"
PORT = 8000
def _ensure_token(engine: TradingEngine) -> str:
"""App-Token für Trade-Aktionen sicherstellen — bei Bedarf generieren."""
cfg = engine.cfg
if "web" not in cfg:
cfg["web"] = {}
token = cfg["web"].get("api_token", "").strip()
if not token:
token = secrets.token_urlsafe(16)
cfg["web"]["api_token"] = token
try:
save_config(cfg)
except Exception as e:
log.warning(f"Token speichern: {e}")
log.warning("=" * 54)
log.warning(f" APP-TOKEN (für Trade-Steuerung am Handy):\n {token}")
log.warning(" Einmalig im Handy-UI eingeben. Liegt in [web] api_token.")
log.warning("=" * 54)
return token
@contextlib.asynccontextmanager
async def lifespan(app: FastAPI):
engine = TradingEngine()
app.state.token = _ensure_token(engine)
if not engine.start():
log.error(f"Engine-Start fehlgeschlagen: {engine._error}")
# Server läuft trotzdem weiter, damit /api/health den Fehler zeigt.
app.state.engine = engine
try:
yield
finally:
engine.stop()
log.info("Engine gestoppt")
def _require_auth(x_auth_token: str | None):
"""Prüft den App-Token (konstante-Zeit-Vergleich).
Ist die Token-Pflicht abgeschaltet ([web] require_token=false), entfällt
die Prüfung — der Zugang ist dann anderweitig abzusichern (z.B. WireGuard)."""
if not getattr(app.state.engine, "require_token", True):
return
expected = app.state.token
if not x_auth_token or not secrets.compare_digest(x_auth_token, expected):
raise HTTPException(status_code=401, detail="Ungültiger Token")
app = FastAPI(title="Oil Trading Agent", version="1.0", lifespan=lifespan)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"], allow_methods=["*"], allow_headers=["*"],
)
@app.middleware("http")
async def _no_cache(request, call_next):
"""UI-Dateien NIE cachen (no-store) — sonst läuft am Handy alte JS/HTML
weiter. no-store statt no-cache, damit der Browser jedesmal frisch lädt."""
resp = await call_next(request)
path = request.url.path
if path == "/" or path.endswith((".js", ".css", ".html", ".json")):
resp.headers["Cache-Control"] = "no-store, no-cache, must-revalidate, max-age=0"
resp.headers["Pragma"] = "no-cache"
resp.headers["Expires"] = "0"
return resp
@app.get("/api/health")
def health():
eng = app.state.engine
return {"running": eng._running,
"connected": eng._connected,
"error": eng._error,
"symbol": eng.data.symbol,
"wave_tf": eng._wave_tf_label}
@app.get("/api/snapshot")
def snapshot():
eng = app.state.engine
return JSONResponse(content=_json_safe(eng.snapshot_cached()))
@app.get("/api/logs")
def logs(lines: int = 200):
"""Letzte N Zeilen der Logdatei (nur das Dateiende lesen — effizient)."""
from core.logger import LOG_FILE
import os
lines = max(1, min(lines, 1000))
try:
size = os.path.getsize(LOG_FILE)
with open(LOG_FILE, "rb") as f:
f.seek(max(0, size - 256 * 1024)) # nur die letzten 256 KB
data = f.read().decode("utf-8", "ignore")
tail = data.splitlines()[-lines:]
return {"lines": tail}
except Exception as e:
return {"error": str(e), "lines": []}
def _add_net(stats: dict, rate: float) -> dict:
"""Ergänzt eine stats_overview-Periode um Netto-Werte nach Quellensteuer.
WHT wird nur auf die Gewinnseite (gross_win) erhoben, Verluste bleiben voll."""
if not stats or stats.get("error"):
return stats
gw = stats.get("gross_win", 0.0) or 0.0
gl = stats.get("gross_loss", 0.0) or 0.0
wht = gw * rate
net_gw = gw - wht
stats["wht_pct"] = round(rate * 100, 3)
stats["wht"] = wht
stats["net_pnl"] = (stats.get("total_pnl", 0.0) or 0.0) - wht
stats["net_profit_factor"] = (net_gw / gl) if gl > 0 else None
stats["net_avg_win"] = (stats.get("avg_win", 0.0) or 0.0) * (1 - rate)
return stats
def _add_net_setup(s: dict, rate: float) -> dict:
"""Netto-Werte nach WHT für eine Setup-Zeile (nur Gewinnseite besteuert)."""
gw = s.get("gross_win", 0.0) or 0.0
gl = s.get("gross_loss", 0.0) or 0.0
s["net_pnl"] = (s.get("total_pnl", 0.0) or 0.0) - gw * rate
s["net_profit_factor"] = (gw * (1 - rate) / gl) if gl > 0 else None
return s
@app.get("/api/stats")
def stats():
eng = app.state.engine
try:
rate = float(eng.cfg["trading"].get("wht_pct", "0")) / 100.0
except Exception:
rate = 0.0
out: dict = {"wht_pct": round(rate * 100, 3)}
for period in ("today", "week", "all"):
try:
out[period] = _add_net(eng.history.stats_overview(period), rate)
except Exception as e:
out[period] = {"error": str(e)}
try:
out["setups"] = [_add_net_setup(s, rate)
for s in eng.history.setup_stats("all")]
except Exception as e:
out["setups"] = [{"error": str(e)}]
try:
out["n_total"] = eng.history.n_total_trades()
except Exception:
out["n_total"] = None
return JSONResponse(content=_json_safe(out))
@app.get("/api/news")
def news():
eng = app.state.engine
headlines, last_fetch, error = eng.news.snapshot()
return JSONResponse(content=_json_safe({
"headlines": headlines[:30],
"last_fetch": last_fetch,
"error": error,
"sentiment": eng.news.sentiment_snapshot(),
}))
@app.get("/api/bars")
def bars(tf: str = "M5", n: int = 90):
"""OHLC-Bars + EMA12/50 + Trend je Timeframe (Charts-Tab). Read-only, kein Token."""
return JSONResponse(content=_json_safe(app.state.engine.get_bars(tf, n)))
@app.get("/api/squeeze_monitor")
def squeeze_monitor():
"""Auto-Squeeze-B4-Monitor: Live-Trades (echt) + Mechanik (Entry-Edge isoliert)
gegen die Backtest-Erwartung. Read-only, kein Token. Cache-TTL 60 s im Engine."""
return JSONResponse(content=_json_safe(app.state.engine.squeeze_monitor()))
@app.get("/api/pbreak_accuracy")
def pbreak_accuracy(period: str = "all"):
"""Trefferquote der Abprall/Durchbruch-Prognose (P(break)-Modell), ausgewertet
gegen candles_m1. period: today|week|all. Read-only, kein Token."""
return JSONResponse(content=_json_safe(app.state.engine.history.pbreak_accuracy(period)))
@app.websocket("/ws")
async def ws_snapshot(websocket: WebSocket):
"""Pusht den Live-Snapshot im Sekundentakt. Mehrere Geräte parallel ok."""
await websocket.accept()
eng = app.state.engine
loop = asyncio.get_event_loop()
try:
while True:
# snapshot() ist blockierend (Locks + sqlite) → im Threadpool holen,
# damit der Event-Loop frei bleibt. Cache teilt EINE Berechnung über
# alle WS-Clients (TTL) → DB-Last unabhängig von der Geräteanzahl.
snap = await loop.run_in_executor(None, eng.snapshot_cached)
await websocket.send_json(_json_safe(snap))
await asyncio.sleep(WS_PUSH_S)
except WebSocketDisconnect:
pass
except Exception as e:
log.warning(f"WebSocket beendet: {e}")
# ── Trade-Steuerung (token-geschützt) ─────────────────────────────────────────
@app.post("/api/auth/check")
async def auth_check(x_auth_token: str | None = Header(default=None)):
"""Frontend prüft hiermit, ob der eingegebene Token gültig ist."""
_require_auth(x_auth_token)
return {"ok": True}
@app.post("/api/order")
async def order(payload: dict = Body(...),
x_auth_token: str | None = Header(default=None)):
_require_auth(x_auth_token)
side = (payload.get("side") or "").lower()
eng = app.state.engine
loop = asyncio.get_event_loop()
if side == "long":
fn = eng.open_long
elif side == "short":
fn = eng.open_short
elif side == "close":
fn = eng.close
else:
raise HTTPException(status_code=400, detail="side muss long|short|close sein")
log.info(f"[WEB] Order angefordert: {side}")
err = await loop.run_in_executor(None, fn) # blockierender MT5-Call
if err:
log.warning(f"[WEB] Order {side} fehlgeschlagen: {err}")
return JSONResponse(status_code=200, content={"ok": False, "error": str(err)})
return {"ok": True, "side": side}
@app.post("/api/emergency")
async def emergency(payload: dict = Body(...),
x_auth_token: str | None = Header(default=None)):
"""Notfall-Verlust-Schwelle setzen/löschen. `value`: Betrag (Kontowährung)
oder null/0 = aus. Der Server schließt die Position, sobald P&L ≤ value."""
_require_auth(x_auth_token)
cur = app.state.engine.set_emergency(payload.get("value"))
log.info(f"[WEB] Notfall-Stop gesetzt: {cur}")
return {"ok": True, "emergency_loss": cur}
@app.post("/api/takeprofit")
async def takeprofit(payload: dict = Body(...),
x_auth_token: str | None = Header(default=None)):
"""Gewinn-Ziel setzen/löschen. `value`: Betrag (Kontowährung) oder null/0 = aus.
Der Server schließt die Position, sobald P&L ≥ +value."""
_require_auth(x_auth_token)
cur = app.state.engine.set_takeprofit(payload.get("value"))
log.info(f"[WEB] Gewinn-Ziel gesetzt: {cur}")
return {"ok": True, "takeprofit": cur}
@app.post("/api/sltp")
async def set_sltp(payload: dict = Body(...),
x_auth_token: str | None = Header(default=None)):
"""Manuelles SL/TP der offenen Position setzen (deaktiviert das Trailing).
`sl`/`tp`: Preis oder null (= unverändert). Broker-Prüfung → error bei Ablehnung."""
_require_auth(x_auth_token)
res = await asyncio.get_event_loop().run_in_executor(
None, lambda: app.state.engine.set_sltp(payload.get("sl"), payload.get("tp")))
log.info(f"[WEB] SL/TP: {res}")
return res
@app.post("/api/srclose")
async def toggle_srclose(x_auth_token: str | None = Header(default=None)):
"""Automatischen S/R-Close (P(break)-Regel) an/aus."""
_require_auth(x_auth_token)
eng = app.state.engine
state = eng.set_sr_autoclose(not eng._auto_sr_close)
return {"ok": True, "enabled": state}
@app.post("/api/autosqueeze")
async def toggle_autosqueeze(x_auth_token: str | None = Header(default=None)):
"""Autonomen Squeeze-Entry (echte Order auf Ausbruch) an/aus."""
_require_auth(x_auth_token)
eng = app.state.engine
state = eng.set_auto_squeeze(not eng._auto_squeeze)
return {"ok": True, "enabled": state}
@app.post("/api/srclose_min")
async def srclose_min(payload: dict = Body(...),
x_auth_token: str | None = Header(default=None)):
"""Mindestgewinn (EUR) für den S/R-Auto-Close setzen. 0/leer = aus."""
_require_auth(x_auth_token)
get_logger("server").info(f"[WEB] Mindestgewinn angefordert: {payload.get('value')!r}")
cur = app.state.engine.set_sr_close_min_gain(payload.get("value"))
return {"ok": True, "min_gain": cur}
@app.post("/api/trail")
async def toggle_trail(x_auth_token: str | None = Header(default=None)):
_require_auth(x_auth_token)
state = await asyncio.get_event_loop().run_in_executor(
None, app.state.engine.toggle_trail)
return {"ok": True, "enabled": state}
def _json_safe(obj):
"""Macht den Snapshot robust JSON-fähig (numpy-Floats, datetime, Tupel)."""
import datetime as _dt
if isinstance(obj, dict):
return {str(k): _json_safe(v) for k, v in obj.items()}
if isinstance(obj, (list, tuple)):
return [_json_safe(v) for v in obj]
if isinstance(obj, _dt.datetime):
return obj.isoformat()
# numpy-Skalare → Python-Zahl
if hasattr(obj, "item") and not isinstance(obj, (str, bytes)):
try:
return obj.item()
except Exception:
return str(obj)
return obj
# Frontend zuletzt mounten, damit /api/* und /ws Vorrang haben.
# html=True liefert web/index.html unter "/" aus.
if WEB_DIR.is_dir():
app.mount("/", StaticFiles(directory=str(WEB_DIR), html=True), name="web")
else:
log.warning(f"Web-Verzeichnis fehlt: {WEB_DIR}")
if __name__ == "__main__":
print("=" * 58)
print(" Oil Trading Agent — Web-Backend (FastAPI)")
print(f" Mobile-UI: http://localhost:{PORT}/")
print(f" JSON-API: http://localhost:{PORT}/api/snapshot")
print("=" * 58)
# access_log=False: Uvicorns Zugriffs-Log (eine Zeile je HTTP-Request) schweigt
# jetzt — seit dem 4-s-REST-Polling (2026-07-22) + WS-Push flutete das die Konsole
# (GET /api/snapshot alle paar Sekunden). Betrifft NUR das Access-Log, nicht die
# eigenen oil.*-Logger (core/logger.py, eigene Hierarchie) oder Uvicorn-Fehler.
uvicorn.run(app, host=HOST, port=PORT, log_level="info", access_log=False)