Erster vollstaendiger Durchlauf von docs/review-prompt.md. KEINE Strategie-Aenderung - alles Telemetrie, Doku und toter Code. 1) analyze_divergence.py kannte keine Epochen und war damit selbst driftanfaellig. D0 mischte Vorhersagen des alten und des am 31.07. nachtrainierten P(break)-Modells und meldete dessen Fehlkalibrierung als aktuellen Alarm (das neue Modell hat n=0, Markt seit Fr zu). B las die Prae-Migrations-NULLs von block_reason als blinden Fleck. C druckte bei 0 Zeilen ein "OK", obwohl es fehlende Daten waren. Neu: _EPOCHS + _clamp(). 2) Konsens-Pfeil AR;K: der MQL5-Export rechnete im ~5-s-Takt ein komplettes zweites _verdict(), obwohl der Indikator die Zeile seit v1.33 per Default verwirft. Neu [trading] export_consensus_arrow (Default false). Verifiziert ueber die exportierte CSV: AR;K weg, AR;L und AR;S bleiben. 3) /api/autosqueeze loggt jetzt die Herkunft ([WEB] ...) wie /api/autosignal. Vorher war ein Zustandswechsel nicht als Nutzeraktion belegbar. 4) CLAUDE.md: die Reversal-Kennzahl "OR +0,185 / PF 1,35 / 70 %" stand unkorrigiert an der Fundstelle, die Widerlegung 2000 Zeilen weiter im Legacy-Recheck. Korrektur an die Fundstelle geholt. 5) core/notify.py: zwei tote "import datetime" entfernt (beide Funktionen nutzen _time), funktional nachgetestet. Geprueft und sauber: 0 fehlende Frontend-IDs von 86, nur 2 Config-Schluessel ohne Leser (beide dokumentiert dormant), Snapshot-Median 13 ms und alle DB-Abfragen <13 ms -> keine Performance-Massnahme, 124 Datei- und 83 Funktionsreferenzen in CLAUDE.md stimmen. Zwischenverdacht zurueckgezogen: "block_reason erklaert nur 33 % der WARTEN" war ein Migrations-Artefakt; seit 01.08. 100 % Abdeckung. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
414 lines
16 KiB
Python
414 lines
16 KiB
Python
#!/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)
|
||
# ⚠ Herkunft mitloggen (Fix 2026-08-02, Review-Durchgang 1): `set_auto_squeeze`
|
||
# loggt nur „Auto-Squeeze-Entry AN/AUS" ohne Quelle — beim Nachvollziehen eines
|
||
# Zustandswechsels war dadurch nicht feststellbar, ob der Dashboard-Button oder
|
||
# Code dahinterstand (real: BRK stand auf AUS und es kostete mehrere Schritte,
|
||
# das als Nutzeraktion zu belegen). `/api/autosignal` macht das längst richtig —
|
||
# bei einem Schalter für autonome ECHTGELD-Einstiege gehört die Herkunft ins Log.
|
||
get_logger("server").info(f"[WEB] Auto-Squeeze-Entry {'AN' if state else 'AUS'}")
|
||
return {"ok": True, "enabled": state}
|
||
|
||
|
||
@app.post("/api/autosignal")
|
||
async def toggle_autosignal(x_auth_token: str | None = Header(default=None)):
|
||
"""Autonomen Entry auf die EMPFEHLUNG an/aus (User-Wunsch 2026-07-30).
|
||
⚠ Gemessen NICHT tragfähig (backtest_auto_signal.py: H1 negativ) — bewusst
|
||
vom User gewollt, Default aus, mit 15-Min-Adverse-Schutz."""
|
||
_require_auth(x_auth_token)
|
||
eng = app.state.engine
|
||
state = eng.set_auto_signal(not eng._auto_signal)
|
||
get_logger("server").info(f"[WEB] Auto-Signal-Entry {'AN' if state else 'AUS'}")
|
||
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)
|