#!/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 time 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, trader=None) -> dict: """Ergänzt eine stats_overview-Periode um Netto-Werte nach Quellensteuer. ⚠⚠ SEIT 2026-08-12 AUS DEN **ECHTEN** BUCHUNGEN, nicht mehr geschätzt. Die alte Rechnung war `net_pnl = total_pnl − gross_win × 26,375 %` und unterstellte, die WHT sei endgültig verloren. Sie wird aber täglich per `Tax settlement` erstattet — gemessen über 90 Tage zu **99 %** (einbehalten −2.625,89 €, erstattet +2.608,83 €, verblieben −17,06 €). Die Schätzung sagte −2.619,93 € und machte damit aus **+153 € rund −2.450 €**. ✅ Die Schätzung der EINBEHALTENEN Summe war fast exakt — falsch war allein die Annahme, sie bleibe weg. ⚠ Fällt der Broker aus, wird auf die alte Schätzung zurückgefallen; das Feld `wht_quelle` sagt, welche Zahl gerade drinsteht („echt" / „geschätzt"). ⚠ KURZE ZEITRÄUME sind naturgemäss pessimistisch: die Erstattung kommt erst am Folgetag. Für „today" fehlt sie also meist — deshalb `wht_lag`. """ 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 tp = stats.get("total_pnl", 0.0) or 0.0 echt = None if trader is not None: try: echt = trader.steuer_buchungen(int(stats.get("since") or 0), int(stats.get("until") or time.time())) except Exception: echt = None if echt and echt.get("n"): wht = -float(echt["netto"]) # netto ist negativ → Last positiv stats["wht_quelle"] = "echt" stats["wht_brutto"] = echt["wht"] stats["wht_erstattet"] = echt["tax"] stats["wht_lag"] = bool(echt["tax"] == 0 and echt["wht"] < 0) else: wht = gw * rate stats["wht_quelle"] = "geschätzt" stats["wht_pct"] = round(rate * 100, 3) stats["wht"] = round(wht, 2) stats["net_pnl"] = round(tp - wht, 2) stats["net_profit_factor"] = ((gw - wht) / gl) if gl > 0 else None stats["net_avg_win"] = (stats.get("avg_win", 0.0) or 0.0) * ( 1 - (wht / gw if gw > 0 else 0.0)) 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, getattr(eng, "trader", None)) 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/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/circuit") async def toggle_circuit(x_auth_token: str | None = Header(default=None)): """Tagesverlust-Stopp (Circuit Breaker) an/aus (User-Wunsch 2026-08-19).""" _require_auth(x_auth_token) eng = app.state.engine state = eng.set_circuit_breaker(not (eng._cb_limit_pct > 0)) # Herkunft mitloggen wie bei den anderen Schaltern -- sonst ist beim # Nachvollziehen eines Zustandswechsels nicht feststellbar, ob Dashboard # oder Code geschaltet hat (Fix-Muster vom 02.08.). log.info(f"[WEB] Circuit-Breaker {'AN' if state else 'AUS'}") return {"ok": True, "enabled": state, "limit_pct": eng._cb_limit_pct} @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/manualmargin") async def set_manual_margin(body: dict, x_auth_token: str | None = Header(default=None)): """Feste Einsatz-Margin in Kontowährung setzen (User-Wunsch 2026-08-04). `{"value": 0}` = automatisch (wie bisher `margin_buffer_pct` % der freien Margin). ⚠ Wirkt auf ALLE neuen Positionen — auch die autonomen. Die freie Margin bleibt Obergrenze; `calc_lots` deckelt einen zu hohen Wunschwert und loggt das.""" _require_auth(x_auth_token) try: v = float(body.get("value") or 0) except (TypeError, ValueError): return {"ok": False, "error": "ungültiger Wert"} eng = app.state.engine val = eng.set_manual_margin(v) get_logger("server").info( f"[WEB] Einsatz-Margin: {'automatisch' if val <= 0 else f'{val:.2f}'}") return {"ok": True, "value": val} @app.post("/api/marginpct") async def set_margin_pct(body: dict, x_auth_token: str | None = Header(default=None)): """Einsatz als PROZENT der freien Margin setzen (User-Wunsch 2026-08-05). Gegenstück zu `/api/manualmargin` (fester Betrag). ⚠ Ein gesetzter FESTER Betrag hat Vorrang; der Prozentsatz bleibt dann nur noch die Obergrenze. Wird auf 1–99 % geklemmt — 100 % ließe keinen Puffer für Spread/Swap.""" _require_auth(x_auth_token) try: v = float(body.get("value") or 0) except (TypeError, ValueError): return {"ok": False, "error": "ungültiger Wert"} if v <= 0: return {"ok": False, "error": "Prozentsatz muss > 0 sein"} val = app.state.engine.set_margin_pct(v) get_logger("server").info(f"[WEB] Einsatz-Prozentsatz: {val:.0f} % der freien Margin") return {"ok": True, "value": val} # ⚠⚠ `/api/autosignal` ENTFERNT am 2026-08-12 (User-Entscheidung). Der autonome # Entry auf die Empfehlung war seit dem 01.08. aus und auf BEIDEN # Ausführungs-Achsen gemessen negativ — Markt-Einstieg in allen 8 Feldern, # ruhende Order mit Touch-Fill in allen 6 und dabei sogar schlechter als # Markt. Ursache strukturell: das Bestätigungs-Level wandert. # Sicherung: `.removed_backup/auto_signal_2026-08-12.py` + Git-Historie. @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)