diff --git a/stream-control/app.py b/stream-control/app.py index 5dc1203..c321389 100644 --- a/stream-control/app.py +++ b/stream-control/app.py @@ -1,8 +1,13 @@ """Network profile control surface; backend access is restricted to a forced SSH command.""" from __future__ import annotations -import json, os, secrets, subprocess, time +import json, os, secrets, sqlite3, subprocess, time +from datetime import datetime +from zoneinfo import ZoneInfo from flask import Flask, flash, jsonify, redirect, render_template, request, session, url_for +BASE_DIR = os.path.dirname(os.path.abspath(__file__)) +AUDIT_DB = os.path.join(BASE_DIR, "control-history.db") +TZ = ZoneInfo("Europe/Berlin") app = Flask(__name__) app.secret_key = os.environ["SESSION_SECRET"] app.config.update(SESSION_COOKIE_HTTPONLY=True, SESSION_COOKIE_SAMESITE="Strict", SESSION_COOKIE_SECURE=True) @@ -15,28 +20,42 @@ PROFILES = { } _CACHE = {"at": 0.0, "data": None} -def rpc(command: str, timeout: int = 25): - p = subprocess.run( - ["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8", "edge-maint", command], - capture_output=True, text=True, timeout=timeout, - ) - raw = p.stdout.strip() +def init_audit(): + with sqlite3.connect(AUDIT_DB) as db: + db.execute("CREATE TABLE IF NOT EXISTS audit(id INTEGER PRIMARY KEY AUTOINCREMENT,ts INTEGER NOT NULL,action TEXT NOT NULL,detail TEXT,ok INTEGER NOT NULL,error TEXT,remote_ip TEXT)") + db.execute("CREATE INDEX IF NOT EXISTS idx_audit_ts ON audit(ts DESC)") + +def audit(action, detail, result): try: - data = json.loads(raw) if raw else {} - except json.JSONDecodeError: - data = {"ok": False, "error": (p.stderr or raw or "invalid response")[:300]} - if p.returncode and data.get("ok") is not False: - data = {"ok": False, "error": (p.stderr or raw or f"rpc exit {p.returncode}")[:300]} - return data + remote = (request.headers.get("X-Forwarded-For") or request.remote_addr or "").split(",")[0].strip() + with sqlite3.connect(AUDIT_DB) as db: + db.execute("INSERT INTO audit(ts,action,detail,ok,error,remote_ip) VALUES(?,?,?,?,?,?)", (int(time.time()), action, detail[:500], 1 if result.get("ok") else 0, str(result.get("error", ""))[:500], remote[:80])) + db.execute("DELETE FROM audit WHERE id NOT IN (SELECT id FROM audit ORDER BY id DESC LIMIT 2000)") + except Exception: + app.logger.exception("audit write failed") + +init_audit() + +def rpc(command: str, timeout: int = 25): + try: + p = subprocess.run(["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8", "edge-maint", command], capture_output=True, text=True, timeout=timeout) + raw = p.stdout.strip() + try: + data = json.loads(raw) if raw else {} + except json.JSONDecodeError: + data = {"ok": False, "error": (p.stderr or raw or "invalid response")[:300]} + if p.returncode and data.get("ok") is not False: + data = {"ok": False, "error": (p.stderr or raw or f"rpc exit {p.returncode}")[:300]} + return data + except Exception as exc: + return {"ok": False, "error": str(exc)[:300]} def get_state(force=False): now = time.monotonic() if not force and _CACHE["data"] is not None and now - _CACHE["at"] < 3: return _CACHE["data"] - try: - data = rpc("state") - except Exception as exc: - data = {"ok": False, "error": str(exc)[:300], "limits": {}, "tc": {}, "services": {}, "policy": {}, "ipsets": {}} + data = rpc("state") + data.setdefault("limits", {}); data.setdefault("tc", {}); data.setdefault("services", {}); data.setdefault("policy", {}); data.setdefault("ipsets", {}) _CACHE.update(at=now, data=data) return data @@ -48,6 +67,26 @@ def csrf_token(): app.jinja_env.globals["csrf_token"] = csrf_token +@app.template_filter("dt") +def fmt_datetime(value): + try: return datetime.fromtimestamp(int(value), TZ).strftime("%d.%m.%Y %H:%M:%S") + except (TypeError, ValueError, OSError): return "—" + +@app.template_filter("bytes") +def fmt_bytes(value): + try: number = float(value or 0) + except (TypeError, ValueError): return "—" + for unit in ("B", "KiB", "MiB", "GiB", "TiB"): + if number < 1024 or unit == "TiB": return f"{number:.1f} {unit}" + number /= 1024 + +@app.template_filter("duration") +def fmt_duration(value): + try: seconds = max(0, int(value)) + except (TypeError, ValueError): return "—" + days, seconds = divmod(seconds, 86400); hours, seconds = divmod(seconds, 3600); minutes, _ = divmod(seconds, 60) + return (f"{days}d {hours}h" if days else f"{hours}h {minutes}m" if hours else f"{minutes}m") + @app.before_request def protect_post(): if request.method == "POST" and not secrets.compare_digest(request.form.get("csrf", ""), session.get("csrf", "invalid")): @@ -55,16 +94,25 @@ def protect_post(): @app.get("/") def dashboard(): - state = get_state() - return render_template("dashboard.html", state=state, profiles=PROFILES) + return render_template("dashboard.html", state=get_state(), profiles=PROFILES, page="dashboard") + +@app.get("/history") +def history_page(): + remote = rpc("history", timeout=35) + with sqlite3.connect(AUDIT_DB) as db: + db.row_factory = sqlite3.Row + local_audit = [dict(row) for row in db.execute("SELECT ts,action,detail,ok,error,remote_ip FROM audit ORDER BY id DESC LIMIT 250")] + history = remote.get("history", {}) if remote.get("ok") else {} + for name in ("sessions", "ips", "events", "blocks"): history.setdefault(name, []) + return render_template("history.html", state=get_state(), history=history, audit=local_audit, history_error=remote.get("error"), page="history") @app.post("/profile/") def apply_profile(name): profile = PROFILES.get(name) if not profile: - flash("Unbekanntes Profil", "error") - return redirect(url_for("dashboard")) + flash("Unbekanntes Profil", "error"); return redirect(url_for("dashboard")) result = rpc(f"set-limits {profile['ch']} {profile['vpn']}") + audit("profile", f"{profile['label']}: Region {profile['ch']} / Cloud {profile['vpn']} Mbit", result) _CACHE["data"] = None flash(f"Profil {profile['label']} aktiviert" if result.get("ok") else f"Fehler: {result.get('error','unbekannt')}", "success" if result.get("ok") else "error") return redirect(url_for("dashboard")) @@ -73,12 +121,11 @@ def apply_profile(name): def custom_limits(): try: ch = float(request.form["ch"]); vpn = float(request.form["vpn"]) - if not 0.1 <= ch <= 100 or not 0.1 <= vpn <= 100: - raise ValueError + if not 0.1 <= ch <= 100 or not 0.1 <= vpn <= 100: raise ValueError except (KeyError, ValueError): - flash("Limits müssen zwischen 0,1 und 100 Mbit liegen", "error") - return redirect(url_for("dashboard")) + flash("Limits müssen zwischen 0,1 und 100 Mbit liegen", "error"); return redirect(url_for("dashboard")) result = rpc(f"set-limits {ch:g} {vpn:g}") + audit("limits", f"Region {ch:g} / Cloud {vpn:g} Mbit", result) _CACHE["data"] = None flash("Benutzerdefinierte Limits gesetzt" if result.get("ok") else f"Fehler: {result.get('error','unbekannt')}", "success" if result.get("ok") else "error") return redirect(url_for("dashboard")) @@ -87,35 +134,29 @@ def custom_limits(): def policy(): mode = request.form.get("mode", "") try: - min_sessions = int(request.form.get("min_sessions", "2")) - block_minutes = int(request.form.get("block_minutes", "10")) - if mode not in {"off", "monitor", "block"} or not 2 <= min_sessions <= 8 or not 1 <= block_minutes <= 1440: - raise ValueError + min_sessions = int(request.form.get("min_sessions", "2")); block_minutes = int(request.form.get("block_minutes", "10")) + if mode not in {"off", "monitor", "block"} or not 2 <= min_sessions <= 8 or not 1 <= block_minutes <= 1440: raise ValueError except ValueError: - flash("Ungültige Policy-Werte", "error") - return redirect(url_for("dashboard")) + flash("Ungültige Policy-Werte", "error"); return redirect(url_for("dashboard")) result = rpc(f"set-policy {mode} {min_sessions} {block_minutes}") + audit("policy", f"Modus {mode}, ab {min_sessions} Streams, {block_minutes} Minuten", result) _CACHE["data"] = None flash("Session-Policy aktualisiert" if result.get("ok") else f"Fehler: {result.get('error','unbekannt')}", "success" if result.get("ok") else "error") return redirect(url_for("dashboard")) @app.post("/maintenance/") def maintenance(action): - mapping = { - "refresh": "limiter refresh", "start": "limiter start", "stop": "limiter stop", - "worker": "restart worker", "admin": "restart admin", - } + mapping = {"refresh": "limiter refresh", "start": "limiter start", "stop": "limiter stop", "worker": "restart worker", "admin": "restart admin"} command = mapping.get(action) - if not command: - return "unknown action", 404 + if not command: return "unknown action", 404 result = rpc(command) + audit("maintenance", action, result) _CACHE["data"] = None flash(f"Aktion {action} gestartet" if result.get("ok") else f"Fehler: {result.get('error','unbekannt')}", "success" if result.get("ok") else "error") return redirect(url_for("dashboard")) @app.get("/api/status") -def api_status(): - return jsonify(get_state()) +def api_status(): return jsonify(get_state()) @app.get("/health") def health():