feat: StreamScope remote control center #1

Merged
sascha merged 8 commits from feat/streamscope-control-center-20260907 into main 2026-09-07 09:04:21 +02:00
8 changed files with 405 additions and 0 deletions

View file

@ -0,0 +1,153 @@
#!/usr/bin/env python3
"""Restricted network-profile helper for the edge host."""
from __future__ import annotations
import json, os, re, shlex, sqlite3, subprocess, sys, time
DB = "/var/lib/streamscope/control.db"
DEV = "ens3"
CLASSES = {"ch": "1:200", "vpn": "1:300", "default": "1:999"}
SERVICES = ("netqos", "streamscope-worker", "streamscope-admin", "promtail")
def run(argv, timeout=30):
try:
p = subprocess.run(argv, capture_output=True, text=True, timeout=timeout)
return p.returncode, p.stdout.strip(), p.stderr.strip()
except Exception as exc:
return 125, "", str(exc)
def db():
c = sqlite3.connect(DB, timeout=15)
c.row_factory = sqlite3.Row
c.execute("PRAGMA busy_timeout=15000")
return c
def rate_map():
rc, out, _ = run(["sudo", "tc", "class", "show", "dev", DEV])
found = {}
for line in out.splitlines():
m = re.search(r"^class htb (1:\w+).*?\brate\s+([0-9.]+)([KMG])bit", line, re.I)
if not m:
continue
value = float(m.group(2)) * {"K": .001, "M": 1, "G": 1000}[m.group(3).upper()]
target = next((k for k, v in CLASSES.items() if v == m.group(1)), m.group(1))
found[target] = value
return found, rc
def state():
c = db()
limits = {r["target"]: r["limit_mbit"] for r in c.execute("SELECT target,limit_mbit FROM global_limits")}
settings = {r["key"]: r["value"] for r in c.execute("SELECT key,value FROM settings")}
active_sessions = c.execute("SELECT count(*) FROM session_history WHERE active=1").fetchone()[0]
recent = [dict(r) for r in c.execute("SELECT ts,user_name,kind,detail FROM multisession_events ORDER BY id DESC LIMIT 8")]
active_blocks = c.execute("SELECT count(*) FROM multisession_blocks WHERE active=1 AND expires_at>strftime('%s','now')").fetchone()[0]
c.close()
tc, tc_rc = rate_map()
ipsets = {}
for name in ("geo-a-v4", "geo-a-v6", "cloud-v4", "cloud-v6", "persist-v4", "persist-v6"):
rc, out, _ = run(["sudo", "ipset", "list", name, "-t"])
m = re.search(r"Number of entries:\s*(\d+)", out)
ipsets[name] = int(m.group(1)) if m else 0
services = {}
for name in SERVICES:
rc, out, _ = run(["systemctl", "is-active", name])
services[name] = out or "unknown"
rc, rules, _ = run(["sudo", "iptables", "-t", "mangle", "-S"])
return {
"ok": tc_rc == 0 and all(v == "active" for k, v in services.items() if k != "netqos"),
"ts": int(time.time()), "limits": limits, "tc": tc, "ipsets": ipsets,
"services": services, "active_sessions": active_sessions, "active_blocks": active_blocks,
"policy": {
"enabled": settings.get("multisession_enabled", "0") == "1",
"mode": settings.get("multisession_mode", "monitor"),
"min_sessions": int(settings.get("multisession_min_sessions", "2")),
"block_minutes": int(settings.get("multisession_block_minutes", "60")),
"persist_seconds": int(settings.get("multisession_persist_seconds", "60")),
},
"mark_rules": sum(1 for line in rules.splitlines() if "--set-xmark" in line or "--set-mark" in line),
"recent_events": recent,
}
def set_limits(ch, vpn):
values = {"ch": float(ch), "vpn": float(vpn)}
if any(not 0.1 <= value <= 100 for value in values.values()):
raise ValueError("rate outside 0.1..100 Mbit")
old, _ = rate_map()
changed = []
try:
for target, value in values.items():
cls = CLASSES[target]
rc, out, err = run(["sudo", "tc", "class", "change", "dev", DEV, "parent", "1:", "classid", cls, "htb", "rate", f"{value:g}mbit", "ceil", f"{value:g}mbit"])
if rc:
raise RuntimeError(err or out or f"tc failed for {target}")
changed.append(target)
c = db()
for target, value in values.items():
c.execute("UPDATE global_limits SET limit_mbit=?,updated_at=datetime('now') WHERE target=?", (value, target))
c.execute("UPDATE limit_profiles SET active=0")
c.commit(); c.close()
except Exception:
for target in changed:
if target in old:
value = old[target]; cls = CLASSES[target]
run(["sudo", "tc", "class", "change", "dev", DEV, "parent", "1:", "classid", cls, "htb", "rate", f"{value:g}mbit", "ceil", f"{value:g}mbit"])
raise
return {"ok": True, "limits": values, "tc": rate_map()[0]}
def set_policy(mode, min_sessions=None, block_minutes=None):
if mode not in {"off", "monitor", "block"}:
raise ValueError("invalid policy mode")
c = db()
c.execute("UPDATE settings SET value=? WHERE key='multisession_enabled'", ("0" if mode == "off" else "1",))
if mode != "off":
c.execute("UPDATE settings SET value=? WHERE key='multisession_mode'", (mode,))
if min_sessions is not None:
value = int(min_sessions)
if not 2 <= value <= 8: raise ValueError("min sessions outside 2..8")
c.execute("UPDATE settings SET value=? WHERE key='multisession_min_sessions'", (str(value),))
if block_minutes is not None:
value = int(block_minutes)
if not 1 <= value <= 1440: raise ValueError("block minutes outside 1..1440")
c.execute("UPDATE settings SET value=? WHERE key='multisession_block_minutes'", (str(value),))
c.commit(); c.close()
rc, out, err = run(["sudo", "systemctl", "restart", "streamscope-worker"], timeout=20)
if rc: raise RuntimeError(err or out or "worker restart failed")
return {"ok": True, "policy": state()["policy"]}
def background_limiter(action):
if action not in {"start", "stop", "refresh"}: raise ValueError("invalid limiter action")
unit = "netprofile-rebuild" if action in {"start", "refresh"} else "netprofile-clear"
rc, out, err = run(["sudo", "systemd-run", "--unit", unit, "--collect", "/usr/lib/systemd/system/netqos/link-optimizer.sh", action], timeout=15)
if rc and "already exists" not in (out + err): raise RuntimeError(err or out)
return {"ok": True, "started": action}
def restart(name):
mapping = {"worker": "streamscope-worker", "admin": "streamscope-admin"}
if name not in mapping: raise ValueError("invalid service")
rc, out, err = run(["sudo", "systemctl", "restart", mapping[name]], timeout=20)
if rc: raise RuntimeError(err or out)
return {"ok": True, "restarted": name}
def main():
raw = os.environ.get("SSH_ORIGINAL_COMMAND", "state").strip()
argv = shlex.split(raw)
if not argv: argv = ["state"]
if argv == ["state"]:
result = state()
elif len(argv) == 3 and argv[0] == "set-limits":
result = set_limits(argv[1], argv[2])
elif argv[0] == "set-policy" and 2 <= len(argv) <= 4:
result = set_policy(argv[1], argv[2] if len(argv) > 2 else None, argv[3] if len(argv) > 3 else None)
elif len(argv) == 2 and argv[0] == "limiter":
result = background_limiter(argv[1])
elif len(argv) == 2 and argv[0] == "restart":
result = restart(argv[1])
else:
raise ValueError("unsupported command")
print(json.dumps(result, ensure_ascii=False, separators=(",", ":")))
if __name__ == "__main__":
try:
main()
except Exception as exc:
print(json.dumps({"ok": False, "error": str(exc)[:300]}, ensure_ascii=False, separators=(",", ":")))
sys.exit(1)

123
stream-control/app.py Normal file
View file

@ -0,0 +1,123 @@
"""Network profile control surface; backend access is restricted to a forced SSH command."""
from __future__ import annotations
import json, os, secrets, subprocess, time
from flask import Flask, flash, jsonify, redirect, render_template, request, session, url_for
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)
PROFILES = {
"strict": {"label": "Strikt", "ch": 0.5, "vpn": 0.5, "color": "#ff667d"},
"standard": {"label": "Standard", "ch": 1.0, "vpn": 1.0, "color": "#a88dff"},
"relaxed": {"label": "Entspannt", "ch": 5.0, "vpn": 5.0, "color": "#32d6a3"},
"open": {"label": "Off-Peak", "ch": 10.0, "vpn": 10.0, "color": "#26c6da"},
}
_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()
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
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": {}}
_CACHE.update(at=now, data=data)
return data
def csrf_token():
token = session.get("csrf")
if not token:
token = session["csrf"] = secrets.token_urlsafe(32)
return token
app.jinja_env.globals["csrf_token"] = csrf_token
@app.before_request
def protect_post():
if request.method == "POST" and not secrets.compare_digest(request.form.get("csrf", ""), session.get("csrf", "invalid")):
return "invalid request token", 400
@app.get("/")
def dashboard():
state = get_state()
return render_template("dashboard.html", state=state, profiles=PROFILES)
@app.post("/profile/<name>")
def apply_profile(name):
profile = PROFILES.get(name)
if not profile:
flash("Unbekanntes Profil", "error")
return redirect(url_for("dashboard"))
result = rpc(f"set-limits {profile['ch']} {profile['vpn']}")
_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"))
@app.post("/limits")
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
except (KeyError, ValueError):
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}")
_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"))
@app.post("/policy")
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
except ValueError:
flash("Ungültige Policy-Werte", "error")
return redirect(url_for("dashboard"))
result = rpc(f"set-policy {mode} {min_sessions} {block_minutes}")
_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/<action>")
def maintenance(action):
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
result = rpc(command)
_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())
@app.get("/health")
def health():
state = get_state()
return jsonify({"ok": bool(state.get("ok")), "backend": "reachable" if state else "unreachable"}), 200 if state.get("ok") else 503

View file

@ -0,0 +1,22 @@
server {
listen 80;
listen [::]:80;
server_name center.sascha-lutz.de;
auth_basic "StreamScope Control";
auth_basic_user_file /etc/nginx/.center.htpasswd;
location / {
proxy_pass http://127.0.0.1:9092;
proxy_http_version 1.1;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
proxy_read_timeout 60s;
proxy_connect_timeout 10s;
add_header X-Content-Type-Options nosniff always;
add_header X-Frame-Options DENY always;
add_header Referrer-Policy no-referrer always;
}
}

File diff suppressed because one or more lines are too long

View file

@ -0,0 +1 @@
<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 128 128"><defs><linearGradient id="g" x1="0" y1="0" x2="1" y2="1"><stop stop-color="#9b82ff"/><stop offset="1" stop-color="#5634dc"/></linearGradient></defs><rect width="128" height="128" rx="30" fill="#0a1321"/><circle cx="64" cy="64" r="43" fill="url(#g)"/><path d="M82 39a29 29 0 1 0 9 20h-16a14 14 0 1 1-5-10z" fill="#0d1728"/><circle cx="79" cy="44" r="9" fill="#d7ceff"/></svg>

After

Width:  |  Height:  |  Size: 436 B

View file

@ -0,0 +1,21 @@
[Unit]
Description=Network Profile Cache
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=root
WorkingDirectory=/app-config/netprofile-cache
EnvironmentFile=/app-config/netprofile-cache/runtime.env
ExecStart=/usr/bin/python3 -m gunicorn --bind 127.0.0.1:9092 --workers 1 --threads 4 --timeout 45 --access-logfile - app:app
Restart=always
RestartSec=5
NoNewPrivileges=true
PrivateTmp=true
ProtectHome=read-only
ProtectSystem=strict
ReadWritePaths=/app-config/netprofile-cache
[Install]
WantedBy=multi-user.target

View file

@ -0,0 +1,22 @@
<!doctype html>
<html lang="de" data-theme="dark">
<head>
<meta charset="utf-8"><meta name="viewport" content="width=device-width,initial-scale=1">
<meta name="theme-color" content="#0d1728"><title>StreamScope Control</title>
<link rel="stylesheet" href="{{ url_for('static',filename='admin.css') }}">
<link rel="icon" href="{{ url_for('static',filename='icon.svg') }}">
</head>
<body>
<aside class="sidebar">
<a class="brand" href="/"><span class="brand-mark">◈</span><span>Stream<span>Scope</span></span></a>
<p class="nav-label">CONTROL</p>
<nav><a class="nav-item active" href="/"><i>⌁</i>Traffic Policy</a><a class="nav-item" href="/api/status"><i>{ }</i>Status API</a></nav>
<div class="sidebar-foot"><span class="status-dot {{ 'ok' if state.get('ok') else 'down' }}"></span><small>Edge {{ 'verbunden' if state.get('ok') else 'nicht erreichbar' }}</small></div>
</aside>
<header class="mobile-head"><a class="brand" href="/"><span class="brand-mark">◈</span><span>Stream<span>Scope</span></span></a></header>
<main>
<div class="topbar"><div><p class="eyebrow">REMOTE CONTROL</p><h1>TV & Session Policy</h1></div><div class="top-actions"><span class="status-pill {{ 'ok' if state.get('ok') else 'down' }}">{{ 'OVH ONLINE' if state.get('ok') else 'OVH OFFLINE' }}</span></div></div>
{% with messages=get_flashed_messages(with_categories=true) %}{% for category,message in messages %}<div class="notice {{ category }}"><strong>{{ '✓' if category=='success' else '!' }}</strong><p>{{ message }}</p></div>{% endfor %}{% endwith %}
{% block content %}{% endblock %}
</main>
</body></html>

View file

@ -0,0 +1,58 @@
{% extends 'base.html' %}{% block content %}
{% set limits=state.get('limits',{}) %}{% set tc=state.get('tc',{}) %}{% set policy=state.get('policy',{}) %}
<section class="kpi-grid">
<article class="kpi"><span>REGION-PROFIL</span><strong>{{ limits.get('ch','—') }} Mbit</strong><small>tc: {{ tc.get('ch','—') }} Mbit</small></article>
<article class="kpi"><span>CLOUD-PROFIL</span><strong>{{ limits.get('vpn','—') }} Mbit</strong><small>tc: {{ tc.get('vpn','—') }} Mbit</small></article>
<article class="kpi"><span>SESSIONS</span><strong>{{ state.get('active_sessions','—') }}</strong><small>{{ state.get('active_blocks',0) }} aktive Sperren</small></article>
<article class="kpi"><span>POLICY</span><strong>{{ ('AUS' if not policy.get('enabled') else policy.get('mode','—')|upper) }}</strong><small>ab {{ policy.get('min_sessions','—') }} Streams</small></article>
</section>
<section class="panel"><div class="panel-head"><div><p class="eyebrow">SCHNELLWAHL</p><h2>Traffic-Profile</h2></div><small>Wirkt sofort auf die aktiven tc-Klassen am OVH</small></div>
<div class="profile-grid control-profiles">
{% for key,p in profiles.items() %}
<form method="post" action="{{ url_for('apply_profile',name=key) }}" class="profile-card" style="--accent:{{ p.color }}">
<input type="hidden" name="csrf" value="{{ csrf_token() }}"><span class="profile-dot"></span><h3>{{ p.label }}</h3><strong>{{ p.ch }} / {{ p.vpn }} Mbit</strong><small>Region / Cloud</small><button type="submit">Aktivieren</button>
</form>
{% endfor %}
</div>
</section>
<section class="grid two">
<article class="panel"><div class="panel-head"><div><p class="eyebrow">MANUELL</p><h2>Bandbreiten</h2></div></div>
<form method="post" action="{{ url_for('custom_limits') }}" class="form-grid compact-form">
<input type="hidden" name="csrf" value="{{ csrf_token() }}">
<label>Region Mbit<input type="number" name="ch" min="0.1" max="100" step="0.1" value="{{ limits.get('ch',1) }}" required></label>
<label>Cloud Mbit<input type="number" name="vpn" min="0.1" max="100" step="0.1" value="{{ limits.get('vpn',1) }}" required></label>
<button class="full" type="submit">Limits anwenden</button>
</form>
</article>
<article class="panel"><div class="panel-head"><div><p class="eyebrow">SESSION POLICY</p><h2>Mehrfachnutzung</h2></div></div>
<form method="post" action="{{ url_for('policy') }}" class="form-grid compact-form">
<input type="hidden" name="csrf" value="{{ csrf_token() }}">
<label>Modus<select name="mode"><option value="off" {{ 'selected' if not policy.get('enabled') }}>Aus</option><option value="monitor" {{ 'selected' if policy.get('enabled') and policy.get('mode')=='monitor' }}>Nur beobachten</option><option value="block" {{ 'selected' if policy.get('enabled') and policy.get('mode')=='block' }}>Blockieren</option></select></label>
<label>Ab Streams<input type="number" name="min_sessions" min="2" max="8" value="{{ policy.get('min_sessions',2) }}"></label>
<label>Sperre Minuten<input type="number" name="block_minutes" min="1" max="1440" value="{{ policy.get('block_minutes',10) }}"></label>
<button class="full" type="submit">Policy speichern</button>
</form>
</article>
</section>
<section class="grid two">
<article class="panel"><div class="panel-head"><div><p class="eyebrow">RUNTIME</p><h2>Dienste & Regeln</h2></div></div>
<div class="table-wrap"><table><tbody>
{% for name,status in state.get('services',{}).items() %}<tr><td>{{ name }}</td><td><span class="badge {{ 'active' if status=='active' else 'danger' }}">{{ status }}</span></td></tr>{% endfor %}
<tr><td>Markierungsregeln</td><td><strong>{{ state.get('mark_rules','—') }}</strong></td></tr>
<tr><td>Region-Netze</td><td><strong>{{ state.get('ipsets',{}).get('geo-a-v4',0)+state.get('ipsets',{}).get('geo-a-v6',0) }}</strong></td></tr>
<tr><td>Cloud-Netze</td><td><strong>{{ state.get('ipsets',{}).get('cloud-v4',0)+state.get('ipsets',{}).get('cloud-v6',0) }}</strong></td></tr>
</tbody></table></div>
<div class="action-row">
{% for action,label,cls in [('refresh','Netzlisten aktualisieren','ghost'),('worker','Worker neu starten','ghost'),('admin','Admin neu starten','ghost')] %}<form method="post" action="{{ url_for('maintenance',action=action) }}"><input type="hidden" name="csrf" value="{{ csrf_token() }}"><button class="{{ cls }}">{{ label }}</button></form>{% endfor %}
</div>
</article>
<article class="panel"><div class="panel-head"><div><p class="eyebrow">EREIGNISSE</p><h2>Letzte Policy-Aktionen</h2></div></div>
<div class="event-list">{% for e in state.get('recent_events',[]) %}<div class="event"><span class="event-kind">{{ e.kind }}</span><div><strong>{{ e.user_name or 'System' }}</strong><small>{{ e.detail }}</small></div></div>{% else %}<p class="muted">Keine Ereignisse vorhanden.</p>{% endfor %}</div>
</article>
</section>
{% if state.get('error') %}<div class="notice error"><strong>!</strong><p>{{ state.error }}</p></div>{% endif %}
{% endblock %}