From e4c2012cf88caf5f892d41b89fa55e7ac9be2df2 Mon Sep 17 00:00:00 2001 From: savsis Date: Thu, 10 Sep 2026 17:45:36 +0500 Subject: [PATCH] core: env-driven config, sqlite schema, xray/node management, subscription links Co-Authored-By: Claude Sonnet 5 --- .env.example | 31 ++++ .gitignore | 9 ++ config.py | 100 +++++++++++++ db.py | 344 +++++++++++++++++++++++++++++++++++++++++++ links.py | 112 ++++++++++++++ nodeprov.py | 374 +++++++++++++++++++++++++++++++++++++++++++++++ requirements.txt | 4 + xray_manager.py | 208 ++++++++++++++++++++++++++ 8 files changed, 1182 insertions(+) create mode 100644 .env.example create mode 100644 .gitignore create mode 100644 config.py create mode 100644 db.py create mode 100644 links.py create mode 100644 nodeprov.py create mode 100644 requirements.txt create mode 100644 xray_manager.py diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..1d28c89 --- /dev/null +++ b/.env.example @@ -0,0 +1,31 @@ +# Copy this to .env and fill in real values. Never commit .env. + +# From @BotFather +BOT_TOKEN=123456789:AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA +BOT_USERNAME=YourBot_robot + +# Telegram user IDs allowed into the bot's admin menu and the web panel, comma-separated +ADMIN_IDS=111111111,222222222 + +# Web panel login password (change it any time with: mbs pass) +ADMIN_PANEL_PASSWORD=change-me + +# Domains — point all three A records at this server's IP (see README) +PANEL_DOMAIN=panel.example.com +SUB_DOMAIN=sub.example.com +SITE_DOMAIN=example.com + +# Reality identity for the local node (this same box). Generate with: +# /usr/local/bin/xray x25519 +# XRAY_PUBLIC_KEY is the "Password (PublicKey)" line; keep the matching +# private key only inside /usr/local/etc/xray/config.json (never here). +XRAY_PUBLIC_KEY= +REALITY_SNI=www.wildberries.ru + +# Any 16-hex-char string per transport, e.g.: openssl rand -hex 8 +XRAY_SHORT_ID_TCP= +XRAY_SHORT_ID_GRPC= +XRAY_SHORT_ID_XHTTP= + +# Public hostname clients connect to for the local node (A record -> this server) +DE1_ADDRESS=de1.example.com diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..5516894 --- /dev/null +++ b/.gitignore @@ -0,0 +1,9 @@ +.env +*.db +*.db-journal +*.db-wal +*.db-shm +__pycache__/ +*.pyc +venv/ +.claude/ diff --git a/config.py b/config.py new file mode 100644 index 0000000..7d66490 --- /dev/null +++ b/config.py @@ -0,0 +1,100 @@ +import os + + +def _load_dotenv(path): + if not os.path.exists(path): + return + with open(path, encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line or line.startswith("#") or "=" not in line: + continue + key, _, value = line.partition("=") + os.environ.setdefault(key.strip(), value.strip()) + + +BASE_DIR = os.path.dirname(os.path.abspath(__file__)) +_load_dotenv(os.path.join(BASE_DIR, ".env")) + + +def env(key, default=None, required=False): + val = os.environ.get(key, default) + if required and not val: + raise RuntimeError(f"missing required env var: {key} — copy .env.example to .env and fill it in") + return val + + +BOT_TOKEN = env("BOT_TOKEN", required=True) +BOT_USERNAME = env("BOT_USERNAME", required=True) +ADMIN_IDS = {int(x) for x in env("ADMIN_IDS", "").split(",") if x.strip()} +ADMIN_PANEL_PASSWORD = env("ADMIN_PANEL_PASSWORD", required=True) +PANEL_DOMAIN = env("PANEL_DOMAIN", required=True) +SUB_DOMAIN = env("SUB_DOMAIN", required=True) +SITE_DOMAIN = env("SITE_DOMAIN", required=True) + +DB_PATH = os.path.join(BASE_DIR, "mbs.db") +XRAY_CONFIG_PATH = "/usr/local/etc/xray/config.json" + +XRAY_PUBLIC_KEY = env("XRAY_PUBLIC_KEY", required=True) +REALITY_SNI = env("REALITY_SNI", "www.wildberries.ru") +XRAY_SHORT_ID_TCP = env("XRAY_SHORT_ID_TCP", required=True) +XRAY_SHORT_ID_GRPC = env("XRAY_SHORT_ID_GRPC", required=True) +XRAY_SHORT_ID_XHTTP = env("XRAY_SHORT_ID_XHTTP", required=True) +DE1_ADDRESS = env("DE1_ADDRESS", f"de1.{SITE_DOMAIN}") + +DE1_TRANSPORTS = [ + { + "tag": "vless-tcp-reality", + "label": "TCP + Reality (основной)", + "network": "tcp", + "security": "reality", + "address": DE1_ADDRESS, + "port": 443, + "public_key": XRAY_PUBLIC_KEY, + "short_id": XRAY_SHORT_ID_TCP, + "sni": REALITY_SNI, + "flow": "xtls-rprx-vision", + }, + { + "tag": "vless-grpc-reality", + "label": "gRPC + Reality", + "network": "grpc", + "security": "reality", + "address": DE1_ADDRESS, + "port": 2053, + "public_key": XRAY_PUBLIC_KEY, + "short_id": XRAY_SHORT_ID_GRPC, + "sni": REALITY_SNI, + "service_name": "mbs-grpc", + }, + { + "tag": "vless-xhttp-reality", + "label": "XHTTP + Reality", + "network": "xhttp", + "security": "reality", + "address": DE1_ADDRESS, + "port": 2087, + "public_key": XRAY_PUBLIC_KEY, + "short_id": XRAY_SHORT_ID_XHTTP, + "sni": REALITY_SNI, + "path": "/mbs-xh", + }, + { + "tag": "vless-ws-tls", + "label": "WebSocket + TLS", + "network": "ws", + "security": "tls", + "address": DE1_ADDRESS, + "port": 8880, + "path": "/mbs-ws", + }, +] + +PLANS = [ + {"code": "7d", "label": "7 дней", "days": 7}, + {"code": "1m", "label": "1 месяц", "days": 30}, + {"code": "3m", "label": "3 месяца", "days": 90}, + {"code": "6m", "label": "6 месяцев", "days": 180}, + {"code": "1y", "label": "1 год", "days": 365}, +] +PLANS_BY_CODE = {p["code"]: p for p in PLANS} diff --git a/db.py b/db.py new file mode 100644 index 0000000..0a1a215 --- /dev/null +++ b/db.py @@ -0,0 +1,344 @@ +import sqlite3 +import secrets +import datetime +import contextlib + +from config import DB_PATH + +SCHEMA = """ +CREATE TABLE IF NOT EXISTS users ( + tg_id INTEGER PRIMARY KEY, + token TEXT UNIQUE NOT NULL, + username TEXT, + created_at TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS subscriptions ( + uuid TEXT PRIMARY KEY, + tg_id INTEGER NOT NULL, + node TEXT NOT NULL, + plan TEXT NOT NULL, + created_at TEXT NOT NULL, + expires_at TEXT NOT NULL, + active INTEGER NOT NULL DEFAULT 1, + source TEXT NOT NULL DEFAULT 'bot' +); + +CREATE TABLE IF NOT EXISTS gift_codes ( + code TEXT PRIMARY KEY, + node TEXT NOT NULL, + plan TEXT NOT NULL, + created_by INTEGER NOT NULL, + created_at TEXT NOT NULL, + used_by INTEGER, + used_at TEXT +); + +CREATE TABLE IF NOT EXISTS admin_sessions ( + token TEXT PRIMARY KEY, + created_at TEXT NOT NULL, + expires_at TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS nodes ( + code TEXT PRIMARY KEY, + label TEXT NOT NULL, + kind TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'active', + address TEXT, + port INTEGER NOT NULL DEFAULT 443, + public_key TEXT, + private_key TEXT, + short_id TEXT, + sni TEXT, + flow TEXT, + shared_uuid TEXT, + provision_token TEXT, + transports_json TEXT, + hysteria_enabled INTEGER NOT NULL DEFAULT 0, + hysteria_port INTEGER, + hysteria_password TEXT, + hysteria_obfs_password TEXT, + enabled INTEGER NOT NULL DEFAULT 1, + created_at TEXT NOT NULL +); +""" + +_NEW_NODE_COLUMNS = { + "transports_json": "TEXT", + "hysteria_enabled": "INTEGER NOT NULL DEFAULT 0", + "hysteria_port": "INTEGER", + "hysteria_password": "TEXT", + "hysteria_obfs_password": "TEXT", +} + + +def _migrate(): + with get_conn() as conn: + cols = {r["name"] for r in conn.execute("PRAGMA table_info(nodes)").fetchall()} + for name, decl in _NEW_NODE_COLUMNS.items(): + if name not in cols: + conn.execute(f"ALTER TABLE nodes ADD COLUMN {name} {decl}") + + +def now_iso(): + return datetime.datetime.utcnow().isoformat() + + +@contextlib.contextmanager +def get_conn(): + conn = sqlite3.connect(DB_PATH, timeout=10) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA foreign_keys = ON") + conn.execute("PRAGMA journal_mode = WAL") + try: + yield conn + conn.commit() + finally: + conn.close() + + +def init_db(): + with get_conn() as conn: + conn.executescript(SCHEMA) + _migrate() + _seed_local_node() + + +def _seed_local_node(): + from config import XRAY_PUBLIC_KEY, XRAY_SHORT_ID_TCP, REALITY_SNI, DE1_ADDRESS + + with get_conn() as conn: + row = conn.execute("SELECT 1 FROM nodes WHERE code='de1'").fetchone() + if row: + return + conn.execute( + "INSERT INTO nodes (code, label, kind, address, port, public_key, short_id, sni, flow, enabled, created_at) " + "VALUES ('de1', ?, 'local', ?, 443, ?, ?, ?, 'xtls-rprx-vision', 1, ?)", + ("Локальная нода (de1)", DE1_ADDRESS, XRAY_PUBLIC_KEY, XRAY_SHORT_ID_TCP, REALITY_SNI, now_iso()), + ) + + +def list_nodes(enabled_only: bool = False): + q = "SELECT * FROM nodes" + if enabled_only: + q += " WHERE enabled=1" + q += " ORDER BY (code='de1') DESC, created_at ASC" + with get_conn() as conn: + rows = conn.execute(q).fetchall() + return [dict(r) for r in rows] + + +def get_node(code: str): + with get_conn() as conn: + row = conn.execute("SELECT * FROM nodes WHERE code=?", (code,)).fetchone() + return dict(row) if row else None + + +def create_node(code, label, kind, address, port, public_key, short_id, sni, flow, shared_uuid=None): + with get_conn() as conn: + conn.execute( + "INSERT INTO nodes (code, label, kind, status, address, port, public_key, short_id, sni, flow, shared_uuid, enabled, created_at) " + "VALUES (?,?,?, 'active', ?,?,?,?,?,?,?,1,?)", + (code, label, kind, address, port, public_key, short_id, sni, flow, shared_uuid, now_iso()), + ) + return get_node(code) + + +def create_pending_node(label, address, port, sni, private_key, public_key, short_id, transports_json=None, + hysteria_port=None, hysteria_password=None, hysteria_obfs_password=None): + import json as jsonmod + + code = "n" + secrets.token_hex(4) + token = secrets.token_urlsafe(24) + hysteria_enabled = 1 if hysteria_password else 0 + with get_conn() as conn: + conn.execute( + "INSERT INTO nodes (code, label, kind, status, address, port, public_key, private_key, short_id, sni, flow, " + "provision_token, transports_json, hysteria_enabled, hysteria_port, hysteria_password, hysteria_obfs_password, enabled, created_at) " + "VALUES (?,?, 'managed', 'pending', ?,?,?,?,?,?, 'xtls-rprx-vision', ?, ?, ?, ?, ?, ?, 0, ?)", + (code, label, address, port, public_key, private_key, short_id, sni, token, + transports_json if transports_json is not None else jsonmod.dumps([]), + hysteria_enabled, hysteria_port, hysteria_password, hysteria_obfs_password, now_iso()), + ) + return get_node(code), token + + +def get_node_by_token(token: str): + with get_conn() as conn: + row = conn.execute("SELECT * FROM nodes WHERE provision_token=?", (token,)).fetchone() + return dict(row) if row else None + + +def activate_node(code: str): + with get_conn() as conn: + conn.execute("UPDATE nodes SET status='active', enabled=1 WHERE code=?", (code,)) + return get_node(code) + + +def update_node(code: str, **fields): + if not fields: + return get_node(code) + cols = ", ".join(f"{k}=?" for k in fields) + with get_conn() as conn: + conn.execute(f"UPDATE nodes SET {cols} WHERE code=?", (*fields.values(), code)) + return get_node(code) + + +def delete_node(code: str): + if code == "de1": + raise ValueError("cannot delete the local node") + with get_conn() as conn: + conn.execute("DELETE FROM nodes WHERE code=?", (code,)) + + +def get_or_create_user(tg_id: int, username: str | None): + with get_conn() as conn: + row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone() + if row: + if username and row["username"] != username: + conn.execute("UPDATE users SET username=? WHERE tg_id=?", (username, tg_id)) + return dict(row) + token = secrets.token_hex(16) + conn.execute( + "INSERT INTO users (tg_id, token, username, created_at) VALUES (?,?,?,?)", + (tg_id, token, username, now_iso()), + ) + row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone() + return dict(row) + + +def get_user_by_token(token: str): + with get_conn() as conn: + row = conn.execute("SELECT * FROM users WHERE token=?", (token,)).fetchone() + return dict(row) if row else None + + +def create_subscription(tg_id: int, node: str, plan_days: int, plan_code: str, source: str = "bot", client_uuid: str | None = None): + import uuid as uuidlib + + cid = client_uuid or str(uuidlib.uuid4()) + created = datetime.datetime.utcnow() + expires = created + datetime.timedelta(days=plan_days) + with get_conn() as conn: + conn.execute( + "INSERT INTO subscriptions (uuid, tg_id, node, plan, created_at, expires_at, active, source) " + "VALUES (?,?,?,?,?,?,1,?)", + (cid, tg_id, node, plan_code, created.isoformat(), expires.isoformat(), source), + ) + return {"uuid": cid, "tg_id": tg_id, "node": node, "plan": plan_code, "expires_at": expires.isoformat()} + + +def list_active_subscriptions(tg_id: int | None = None, node: str | None = None): + q = "SELECT * FROM subscriptions WHERE active=1 AND expires_at > ?" + params = [now_iso()] + if tg_id is not None: + q += " AND tg_id=?" + params.append(tg_id) + if node is not None: + q += " AND node=?" + params.append(node) + with get_conn() as conn: + rows = conn.execute(q, params).fetchall() + return [dict(r) for r in rows] + + +def deactivate_expired(): + with get_conn() as conn: + expired = conn.execute( + "SELECT * FROM subscriptions WHERE active=1 AND expires_at <= ?", (now_iso(),) + ).fetchall() + conn.execute("UPDATE subscriptions SET active=0 WHERE active=1 AND expires_at <= ?", (now_iso(),)) + return [dict(r) for r in expired] + + +def create_gift_code(node: str, plan_code: str, created_by: int): + code = secrets.token_hex(6) + with get_conn() as conn: + conn.execute( + "INSERT INTO gift_codes (code, node, plan, created_by, created_at) VALUES (?,?,?,?,?)", + (code, node, plan_code, created_by, now_iso()), + ) + return code + + +def redeem_gift_code(code: str, tg_id: int): + with get_conn() as conn: + row = conn.execute("SELECT * FROM gift_codes WHERE code=?", (code,)).fetchone() + if not row: + return None, "not_found" + if row["used_by"] is not None: + return None, "already_used" + conn.execute( + "UPDATE gift_codes SET used_by=?, used_at=? WHERE code=?", + (tg_id, now_iso(), code), + ) + return dict(row), None + + +def create_admin_session(hours: int = 168): + token = secrets.token_urlsafe(32) + expires = datetime.datetime.utcnow() + datetime.timedelta(hours=hours) + with get_conn() as conn: + conn.execute( + "INSERT INTO admin_sessions (token, created_at, expires_at) VALUES (?,?,?)", + (token, now_iso(), expires.isoformat()), + ) + return token + + +def validate_admin_session(token: str) -> bool: + if not token: + return False + with get_conn() as conn: + row = conn.execute( + "SELECT * FROM admin_sessions WHERE token=? AND expires_at>?", (token, now_iso()) + ).fetchone() + return row is not None + + +def delete_admin_session(token: str): + with get_conn() as conn: + conn.execute("DELETE FROM admin_sessions WHERE token=?", (token,)) + + +def list_all_subscriptions(limit: int = 200): + with get_conn() as conn: + rows = conn.execute( + "SELECT s.*, u.username FROM subscriptions s " + "LEFT JOIN users u ON u.tg_id = s.tg_id " + "ORDER BY s.created_at DESC LIMIT ?", + (limit,), + ).fetchall() + return [dict(r) for r in rows] + + +def revoke_subscription(client_uuid: str): + with get_conn() as conn: + conn.execute("UPDATE subscriptions SET active=0 WHERE uuid=?", (client_uuid,)) + + +def list_gift_codes(limit: int = 200): + with get_conn() as conn: + rows = conn.execute( + "SELECT * FROM gift_codes ORDER BY created_at DESC LIMIT ?", (limit,) + ).fetchall() + return [dict(r) for r in rows] + + +def stats(): + with get_conn() as conn: + users_n = conn.execute("SELECT COUNT(*) c FROM users").fetchone()["c"] + active_n = conn.execute( + "SELECT COUNT(*) c FROM subscriptions WHERE active=1 AND expires_at>?", (now_iso(),) + ).fetchone()["c"] + total_subs = conn.execute("SELECT COUNT(*) c FROM subscriptions").fetchone()["c"] + gifts_created = conn.execute("SELECT COUNT(*) c FROM gift_codes").fetchone()["c"] + gifts_used = conn.execute("SELECT COUNT(*) c FROM gift_codes WHERE used_by IS NOT NULL").fetchone()["c"] + return { + "users": users_n, + "active_subscriptions": active_n, + "total_subscriptions": total_subs, + "gifts_created": gifts_created, + "gifts_used": gifts_used, + } diff --git a/links.py b/links.py new file mode 100644 index 0000000..3551439 --- /dev/null +++ b/links.py @@ -0,0 +1,112 @@ +import base64 +import json +import re +import urllib.parse + +from config import DE1_TRANSPORTS + +TRANSPORT_NAMES = { + "tcp": "TCP", + "grpc": "gRPC", + "xhttp": "XHTTP", + "ws": "WS", +} + + +def display_name(node_label: str) -> str: + return re.sub(r"\s*\([^)]*\)\s*$", "", node_label).strip() + + +def _tcp_reality_uri(client_uuid: str, address: str, port: int, public_key: str, short_id: str, sni: str, flow: str, remark: str) -> str: + params = { + "encryption": "none", "type": "tcp", "security": "reality", + "sni": sni, "fp": "chrome", "pbk": public_key, "sid": short_id, "flow": flow or "xtls-rprx-vision", + } + qs = urllib.parse.urlencode(params) + frag = urllib.parse.quote(remark) + return f"vless://{client_uuid}@{address}:{port}?{qs}#{frag}" + + +def _transport_uris(client_uuid: str, transports: list[dict], base_name: str) -> list[str]: + uris = [] + multi = len(transports) > 1 + for t in transports: + params = {"encryption": "none", "type": t["network"]} + if t["security"] == "reality": + params.update({ + "security": "reality", "sni": t["sni"], "fp": "chrome", + "pbk": t["public_key"], "sid": t["short_id"], + }) + if t["network"] == "tcp": + params["flow"] = t.get("flow") or "xtls-rprx-vision" + elif t["network"] == "grpc": + params["serviceName"] = t["service_name"] + params["mode"] = "gun" + elif t["network"] == "xhttp": + params["path"] = t["path"] + params["mode"] = "auto" + elif t["security"] == "tls": + params["security"] = "tls" + params["sni"] = t["address"] + params["fp"] = "chrome" + params["alpn"] = "http/1.1" + if t["network"] == "ws": + params["path"] = t["path"] + params["host"] = t["address"] + qs = urllib.parse.urlencode(params) + name = f"{base_name} ({TRANSPORT_NAMES.get(t['network'], t['network'])})" if multi else base_name + frag = urllib.parse.quote(name) + uris.append(f"vless://{client_uuid}@{t['address']}:{t['port']}?{qs}#{frag}") + return uris + + +def hysteria_uri_for_node(node: dict, base_name: str) -> str | None: + if not node or not node.get("hysteria_enabled") or not node.get("hysteria_password"): + return None + params = { + "obfs": "salamander", "obfs-password": node["hysteria_obfs_password"], + "insecure": "1", "sni": node["sni"], + } + qs = urllib.parse.urlencode(params) + password = urllib.parse.quote(node["hysteria_password"], safe="") + frag = urllib.parse.quote(f"{base_name} (Hysteria2)") + return f"hysteria2://{password}@{node['address']}:{node['hysteria_port']}/?{qs}#{frag}" + + +def vless_uris_for_node(client_uuid: str, node: dict, base_name: str) -> list[str]: + if not node or not node.get("enabled"): + return [] + if node["code"] == "de1": + return _transport_uris(client_uuid, DE1_TRANSPORTS, base_name) + if node["kind"] == "managed" and node.get("transports_json"): + transports = json.loads(node["transports_json"]) + if transports: + return _transport_uris(client_uuid, transports, base_name) + uuid_to_use = node["shared_uuid"] if node["kind"] == "external" and node.get("shared_uuid") else client_uuid + return [_tcp_reality_uri( + uuid_to_use, node["address"], node["port"], node["public_key"], + node["short_id"], node["sni"], node.get("flow"), base_name, + )] + + +def build_subscription_text(subs: list[dict]) -> str: + import db + + best_by_node = {} + for s in subs: + cur = best_by_node.get(s["node"]) + if cur is None or s["expires_at"] > cur["expires_at"]: + best_by_node[s["node"]] = s + + lines = [] + for node_code, s in best_by_node.items(): + node = db.get_node(node_code) + if not node: + continue + base_name = display_name(node["label"]) + lines.extend(vless_uris_for_node(s["uuid"], node, base_name)) + hy = hysteria_uri_for_node(node, base_name) + if hy: + lines.append(hy) + raw = "\n".join(lines) + return base64.b64encode(raw.encode()).decode() diff --git a/nodeprov.py b/nodeprov.py new file mode 100644 index 0000000..947adf0 --- /dev/null +++ b/nodeprov.py @@ -0,0 +1,374 @@ +import json +import secrets +import socket +import subprocess + +import paramiko + +from config import PANEL_DOMAIN + +MGMT_KEY_PATH = "/root/.ssh/mbs_nodes_ed25519" +LOCAL_TAGS = {"vless-tcp-reality", "vless-grpc-reality", "vless-xhttp-reality", "vless-ws-tls"} +TAG_FLOW = {"vless-tcp-reality": "xtls-rprx-vision"} + +ONE_COMMAND_TEMPLATE = "bash <(curl -Ls https://{panel}/install/{token}.sh)" + +CERTBOT_SNIPPET = """echo "issuing a real TLS cert for {address} (needed for WS+TLS)..." +command -v certbot >/dev/null 2>&1 || apt-get install -y certbot +ss -ltnp | grep -q ':80 ' && {{ echo "something is already on port 80, stop it first"; exit 1; }} +certbot certonly --standalone --non-interactive --agree-tos --register-unsafely-without-email -d {address} +mkdir -p /etc/letsencrypt/renewal-hooks/deploy +cat > /etc/letsencrypt/renewal-hooks/deploy/mbs-restart-xray.sh << 'HOOK' +#!/bin/bash +systemctl restart xray || true +HOOK +chmod +x /etc/letsencrypt/renewal-hooks/deploy/mbs-restart-xray.sh +""" + +HYSTERIA_SNIPPET = """echo "installing Hysteria2..." +bash <(curl -fsSL https://get.hy2.sh/) || true +mkdir -p /etc/hysteria +openssl req -x509 -nodes -newkey ec -pkeyopt ec_paramgen_curve:prime256v1 -keyout /etc/hysteria/server.key -out /etc/hysteria/server.crt -subj "/CN={address}" -days 3650 +cat > /etc/hysteria/config.yaml << 'HYCFG' +listen: :{hysteria_port} +tls: + cert: /etc/hysteria/server.crt + key: /etc/hysteria/server.key +auth: + type: password + password: "{hysteria_password}" +obfs: + type: salamander + salamander: + password: "{hysteria_obfs_password}" +masquerade: + type: proxy + proxy: + url: https://{sni}/ + rewriteHost: true +HYCFG +command -v ufw >/dev/null 2>&1 && ufw allow {hysteria_port}/udp || true +systemctl enable --now hysteria-server.service +sleep 1 +echo "hysteria status: $(systemctl is-active hysteria-server.service)" +""" + +SELF_INSTALL_SCRIPT = """#!/bin/bash +set -e +echo "== MBS Panel node install ==" +export DEBIAN_FRONTEND=noninteractive + +mkdir -p /root/.ssh +chmod 700 /root/.ssh +curl -Ls https://{panel}/mgmt-pubkey.txt >> /root/.ssh/authorized_keys +sort -u /root/.ssh/authorized_keys -o /root/.ssh/authorized_keys +chmod 600 /root/.ssh/authorized_keys +echo "management key installed" + +if [ ! -f /usr/local/bin/xray ]; then + bash -c "$(curl -Ls https://github.com/XTLS/Xray-install/raw/main/install-release.sh)" @ install +fi + +{certbot_block} +mkdir -p /usr/local/etc/xray +cat > /usr/local/etc/xray/config.json << 'XRAYCFG' +{config_json} +XRAYCFG + +command -v ufw >/dev/null 2>&1 && {{ ufw allow 22/tcp || true; {ufw_rules} }} +systemctl enable xray >/dev/null 2>&1 || true +systemctl restart xray +sleep 1 +STATUS=$(systemctl is-active xray) +echo "xray status: $STATUS" + +{hysteria_block} +MY_IP=$(curl -s https://api.ipify.org || echo unknown) +curl -s -X POST "https://{panel}/nodes/register/{token}" -H "Content-Type: application/json" -d "{{\\"ip\\":\\"$MY_IP\\",\\"status\\":\\"$STATUS\\"}}" >/dev/null + +if [ "$STATUS" = "active" ]; then + echo "Done. Node registered — check the panel." +else + echo "xray did not start, check: journalctl -u xray -n 50" +fi +""" + + +def ensure_mgmt_key() -> str: + try: + with open(MGMT_KEY_PATH + ".pub") as f: + return f.read().strip() + except FileNotFoundError: + subprocess.run( + ["ssh-keygen", "-t", "ed25519", "-f", MGMT_KEY_PATH, "-N", "", "-C", "mbs-node-mgmt"], + check=True, capture_output=True, + ) + with open(MGMT_KEY_PATH + ".pub") as f: + return f.read().strip() + + +def generate_reality_keys(): + out = subprocess.run(["/usr/local/bin/xray", "x25519"], check=True, capture_output=True, text=True).stdout + private_key = None + public_key = None + for line in out.splitlines(): + if line.startswith("PrivateKey:"): + private_key = line.split(":", 1)[1].strip() + elif line.startswith("Password (PublicKey):") or line.startswith("PublicKey:"): + public_key = line.split(":", 1)[1].strip() + return private_key, public_key + + +def generate_short_id() -> str: + return secrets.token_hex(8) + + +def generate_hysteria_credentials(): + return secrets.token_urlsafe(16), secrets.token_urlsafe(12) + + +def build_transports(address, tcp_port, sni, public_key, include_ws=False): + transports = [ + { + "tag": "vless-tcp-reality", "label": "TCP + Reality (основной)", + "network": "tcp", "security": "reality", "address": address, "port": tcp_port, + "public_key": public_key, "short_id": generate_short_id(), "sni": sni, + "flow": "xtls-rprx-vision", + }, + { + "tag": "vless-grpc-reality", "label": "gRPC + Reality", + "network": "grpc", "security": "reality", "address": address, "port": 2053, + "public_key": public_key, "short_id": generate_short_id(), "sni": sni, + "service_name": "mbs-grpc", + }, + { + "tag": "vless-xhttp-reality", "label": "XHTTP + Reality", + "network": "xhttp", "security": "reality", "address": address, "port": 2087, + "public_key": public_key, "short_id": generate_short_id(), "sni": sni, + "path": "/mbs-xh", + }, + ] + if include_ws: + transports.append({ + "tag": "vless-ws-tls", "label": "WebSocket + TLS (реальный серт)", + "network": "ws", "security": "tls", "address": address, "port": 8880, + "path": "/mbs-ws", + }) + return transports + + +def one_command(token: str) -> str: + return ONE_COMMAND_TEMPLATE.format(panel=PANEL_DOMAIN, token=token) + + +def _build_config_json(transports, private_key, address): + inbounds = [{ + "tag": "api", "listen": "127.0.0.1", "port": 10085, + "protocol": "dokodemo-door", "settings": {"address": "127.0.0.1"}, + }] + for t in transports: + ib = { + "tag": t["tag"], "listen": "0.0.0.0", "port": t["port"], + "protocol": "vless", + "settings": {"clients": [], "decryption": "none"}, + "sniffing": {"enabled": True, "destOverride": ["http", "tls"]}, + } + if t["security"] == "reality": + ib["streamSettings"] = {"network": t["network"], "security": "reality", "realitySettings": { + "show": False, "dest": f"{t['sni']}:443", "xver": 0, + "serverNames": [t["sni"]], "privateKey": private_key, "shortIds": [t["short_id"]], + }} + if t["network"] == "grpc": + ib["streamSettings"]["grpcSettings"] = {"serviceName": t["service_name"]} + elif t["network"] == "xhttp": + ib["streamSettings"]["xhttpSettings"] = {"path": t["path"], "mode": "auto"} + elif t["security"] == "tls": + ib["streamSettings"] = {"network": "ws", "security": "tls", "wsSettings": {"path": t["path"]}, + "tlsSettings": {"certificates": [{ + "certificateFile": f"/etc/letsencrypt/live/{address}/fullchain.pem", + "keyFile": f"/etc/letsencrypt/live/{address}/privkey.pem", + }]}} + inbounds.append(ib) + + cfg = { + "log": {"loglevel": "warning"}, + "api": {"tag": "api", "services": ["HandlerService", "LoggerService", "StatsService"]}, + "stats": {}, + "policy": { + "levels": {"0": {"statsUserUplink": True, "statsUserDownlink": True}}, + "system": {"statsInboundUplink": True, "statsInboundDownlink": True, "statsOutboundUplink": True, "statsOutboundDownlink": True}, + }, + "routing": {"rules": [{"type": "field", "inboundTag": ["api"], "outboundTag": "api"}]}, + "inbounds": inbounds, + "outbounds": [{"protocol": "freedom", "tag": "direct"}, {"protocol": "blackhole", "tag": "block"}], + } + return json.dumps(cfg, indent=2) + + +def render_install_script(node: dict) -> str: + transports = json.loads(node["transports_json"]) + config_json = _build_config_json(transports, node["private_key"], node["address"]) + by_tag = {t["tag"]: t for t in transports} + has_ws = "vless-ws-tls" in by_tag + + ufw_rules = " ".join(f"ufw allow {t['port']}/tcp || true;" for t in transports) + certbot_block = CERTBOT_SNIPPET.format(address=node["address"]) if has_ws else "" + hysteria_block = "" + if node.get("hysteria_enabled"): + hysteria_block = HYSTERIA_SNIPPET.format( + address=node["address"], hysteria_port=node["hysteria_port"], + hysteria_password=node["hysteria_password"], hysteria_obfs_password=node["hysteria_obfs_password"], + sni=node["sni"], + ) + + return SELF_INSTALL_SCRIPT.format( + panel=PANEL_DOMAIN, token=node["provision_token"], config_json=config_json, + certbot_block=certbot_block, hysteria_block=hysteria_block, ufw_rules=ufw_rules, + ) + + +def _mgmt_connect(address: str, ssh_port: int = 22) -> paramiko.SSHClient: + client = paramiko.SSHClient() + client.set_missing_host_key_policy(paramiko.AutoAddPolicy()) + key = paramiko.Ed25519Key.from_private_key_file(MGMT_KEY_PATH) + client.connect(address, port=ssh_port, username="root", pkey=key, timeout=15, banner_timeout=15, auth_timeout=15) + return client + + +def _remote_edit_clients(node: dict, mutate_fn): + client = _mgmt_connect(node["address"]) + try: + sftp = client.open_sftp() + with sftp.open("/usr/local/etc/xray/config.json") as f: + cfg = json.loads(f.read().decode()) + changed = False + for ib in cfg["inbounds"]: + if ib.get("tag") not in LOCAL_TAGS: + continue + clients = ib["settings"]["clients"] + new_clients = mutate_fn(clients, ib["tag"]) + if new_clients is not None: + ib["settings"]["clients"] = new_clients + changed = True + if changed: + data = json.dumps(cfg, indent=2).encode() + with sftp.open("/usr/local/etc/xray/config.json", "wb") as f: + f.write(data) + sftp.close() + client.exec_command("systemctl restart xray")[1].channel.recv_exit_status() + else: + sftp.close() + finally: + client.close() + + +def remote_add_client(node: dict, client_uuid: str, email: str): + def mutate(clients, tag): + if any(c["id"] == client_uuid for c in clients): + return None + entry = {"id": client_uuid, "email": email} + flow = TAG_FLOW.get(tag) + if flow: + entry["flow"] = flow + clients.append(entry) + return clients + _remote_edit_clients(node, mutate) + + +def remote_remove_client(node: dict, client_uuid: str): + def mutate(clients, tag): + new = [c for c in clients if c["id"] != client_uuid] + return new if len(new) != len(clients) else None + _remote_edit_clients(node, mutate) + + +def remote_sync(node: dict, active_subs: list[dict]): + active_by_id = {s["uuid"]: s for s in active_subs} + + def mutate(clients, tag): + current_ids = {c["id"] for c in clients} + if current_ids == set(active_by_id.keys()): + return None + new_clients = [c for c in clients if c["id"] in active_by_id] + existing_ids = {c["id"] for c in new_clients} + flow = TAG_FLOW.get(tag) + for cid in active_by_id: + if cid not in existing_ids: + entry = {"id": cid, "email": cid} + if flow: + entry["flow"] = flow + new_clients.append(entry) + return new_clients + _remote_edit_clients(node, mutate) + + +def remote_query_stats(node: dict) -> dict: + try: + client = _mgmt_connect(node["address"]) + try: + _, stdout, _ = client.exec_command( + "/usr/local/bin/xray api statsquery --server=127.0.0.1:10085 -pattern 'user>>>'", timeout=10, + ) + out = stdout.read().decode(errors="replace") + finally: + client.close() + data = json.loads(out) if out.strip() else {"stat": []} + except Exception: + return {} + result = {} + for entry in data.get("stat") or []: + parts = entry.get("name", "").split(">>>") + if len(parts) != 4 or parts[0] != "user": + continue + _, email, _, direction = parts + row = result.setdefault(email, {"up": 0, "down": 0}) + value = entry.get("value", 0) + if direction == "uplink": + row["up"] += value + elif direction == "downlink": + row["down"] += value + return result + + +STATUS_CMD = ( + "echo LOAD:$(cut -d' ' -f1 /proc/loadavg); " + "echo MEMTOTAL:$(grep MemTotal /proc/meminfo | awk '{print $2}'); " + "echo MEMAVAIL:$(grep MemAvailable /proc/meminfo | awk '{print $2}'); " + "echo UPTIME:$(cut -d'.' -f1 /proc/uptime); " + "echo XRAY:$(systemctl is-active xray)" +) + + +def remote_node_status(node: dict) -> dict: + try: + client = _mgmt_connect(node["address"]) + try: + _, stdout, _ = client.exec_command(STATUS_CMD, timeout=10) + out = stdout.read().decode(errors="replace") + finally: + client.close() + vals = {} + for line in out.splitlines(): + if ":" in line: + k, v = line.split(":", 1) + vals[k] = v.strip() + mem_total = int(vals["MEMTOTAL"]) // 1024 if vals.get("MEMTOTAL", "").isdigit() else None + mem_avail = int(vals["MEMAVAIL"]) // 1024 if vals.get("MEMAVAIL", "").isdigit() else None + return { + "ok": True, + "load1": float(vals["LOAD"]) if vals.get("LOAD") else None, + "mem_used_mb": (mem_total - mem_avail) if (mem_total is not None and mem_avail is not None) else None, + "mem_total_mb": mem_total, + "uptime_s": int(vals["UPTIME"]) if vals.get("UPTIME", "").lstrip("-").isdigit() else None, + "xray_active": vals.get("XRAY") == "active", + } + except Exception as e: + return {"ok": False, "error": str(e)} + + +def check_node_alive(address: str, port: int, timeout: float = 5.0) -> bool: + try: + with socket.create_connection((address, port), timeout=timeout): + return True + except OSError: + return False diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..6903ff9 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,4 @@ +aiogram==3.15.0 +fastapi==0.115.6 +uvicorn[standard]==0.32.1 +paramiko==3.5.0 diff --git a/xray_manager.py b/xray_manager.py new file mode 100644 index 0000000..54e4474 --- /dev/null +++ b/xray_manager.py @@ -0,0 +1,208 @@ +import json +import os +import subprocess +import fcntl +import contextlib + +from config import XRAY_CONFIG_PATH, DE1_TRANSPORTS + +_LOCK_PATH = XRAY_CONFIG_PATH + ".lock" +_LOCAL_TAGS = {t["tag"] for t in DE1_TRANSPORTS} +_TAG_FLOW = {t["tag"]: t.get("flow") for t in DE1_TRANSPORTS} + + +@contextlib.contextmanager +def _locked(): + with open(_LOCK_PATH, "w") as lf: + fcntl.flock(lf, fcntl.LOCK_EX) + try: + yield + finally: + fcntl.flock(lf, fcntl.LOCK_UN) + + +def _load(): + with open(XRAY_CONFIG_PATH, "r", encoding="utf-8") as f: + return json.load(f) + + +def _save(cfg): + tmp = XRAY_CONFIG_PATH + ".tmp" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(cfg, f, indent=2) + os.replace(tmp, XRAY_CONFIG_PATH) + + +def _reload_xray(): + subprocess.run(["systemctl", "restart", "xray"], check=True, timeout=20) + + +def _local_inbounds(cfg): + return [ib for ib in cfg["inbounds"] if ib.get("tag") in _LOCAL_TAGS] + + +def add_client(client_uuid: str, email: str): + with _locked(): + cfg = _load() + changed = False + for ib in _local_inbounds(cfg): + clients = ib["settings"]["clients"] + if any(c["id"] == client_uuid for c in clients): + continue + entry = {"id": client_uuid, "email": email} + flow = _TAG_FLOW.get(ib["tag"]) + if flow: + entry["flow"] = flow + clients.append(entry) + changed = True + if changed: + _save(cfg) + _reload_xray() + + +def remove_client(client_uuid: str): + with _locked(): + cfg = _load() + changed = False + for ib in _local_inbounds(cfg): + clients = ib["settings"]["clients"] + new_clients = [c for c in clients if c["id"] != client_uuid] + if len(new_clients) != len(clients): + ib["settings"]["clients"] = new_clients + changed = True + if changed: + _save(cfg) + _reload_xray() + + +def list_client_ids(): + with _locked(): + cfg = _load() + ids = set() + for ib in _local_inbounds(cfg): + ids |= {c["id"] for c in ib["settings"]["clients"]} + return ids + + +def sync_from_db(): + import db as dbmod + + expired = dbmod.deactivate_expired() + active = dbmod.list_active_subscriptions(node="de1") + active_by_id = {s["uuid"]: s for s in active} + + with _locked(): + cfg = _load() + changed = False + for ib in _local_inbounds(cfg): + clients = ib["settings"]["clients"] + current_ids = {c["id"] for c in clients} + if current_ids == set(active_by_id.keys()): + continue + new_clients = [c for c in clients if c["id"] in active_by_id] + existing_ids = {c["id"] for c in new_clients} + flow = _TAG_FLOW.get(ib["tag"]) + for cid, sub in active_by_id.items(): + if cid not in existing_ids: + entry = {"id": cid, "email": cid} + if flow: + entry["flow"] = flow + new_clients.append(entry) + ib["settings"]["clients"] = new_clients + changed = True + + if changed: + _save(cfg) + _reload_xray() + + return {"removed_expired": len(expired), "active_now": len(active_by_id), "reloaded": changed} + + +def add_client_to_node(node: dict, client_uuid: str, email: str): + if node["kind"] == "local": + add_client(client_uuid, email) + elif node["kind"] == "managed": + import nodeprov + nodeprov.remote_add_client(node, client_uuid, email) + + +def remove_client_from_node(node: dict, client_uuid: str): + if node["kind"] == "local": + remove_client(client_uuid) + elif node["kind"] == "managed": + import nodeprov + nodeprov.remote_remove_client(node, client_uuid) + + +def query_stats() -> dict: + try: + out = subprocess.run( + ["/usr/local/bin/xray", "api", "statsquery", "--server=127.0.0.1:10085", "-pattern", "user>>>"], + capture_output=True, text=True, timeout=10, + ).stdout + data = json.loads(out) if out.strip() else {"stat": []} + except Exception: + return {} + result = {} + for entry in data.get("stat") or []: + name = entry.get("name", "") + value = entry.get("value", 0) + parts = name.split(">>>") + if len(parts) != 4 or parts[0] != "user": + continue + _, email, _, direction = parts + row = result.setdefault(email, {"up": 0, "down": 0}) + if direction == "uplink": + row["up"] += value + elif direction == "downlink": + row["down"] += value + return result + + +def local_node_status() -> dict: + try: + with open("/proc/loadavg") as f: + load1 = float(f.read().split()[0]) + except Exception: + load1 = None + mem_total = mem_avail = None + try: + with open("/proc/meminfo") as f: + for line in f: + if line.startswith("MemTotal:"): + mem_total = int(line.split()[1]) // 1024 + elif line.startswith("MemAvailable:"): + mem_avail = int(line.split()[1]) // 1024 + except Exception: + pass + uptime_s = None + try: + with open("/proc/uptime") as f: + uptime_s = int(float(f.read().split()[0])) + except Exception: + pass + xray_active = subprocess.run(["systemctl", "is-active", "xray"], capture_output=True, text=True).stdout.strip() + return { + "ok": True, "load1": load1, + "mem_used_mb": (mem_total - mem_avail) if (mem_total and mem_avail is not None) else None, + "mem_total_mb": mem_total, "uptime_s": uptime_s, "xray_active": xray_active == "active", + } + + +def sync_all(): + import db as dbmod + import nodeprov + + expired = dbmod.deactivate_expired() + results = {} + for node in dbmod.list_nodes(): + if node["kind"] == "local": + results[node["code"]] = sync_from_db() + elif node["kind"] == "managed": + active = dbmod.list_active_subscriptions(node=node["code"]) + try: + nodeprov.remote_sync(node, active) + results[node["code"]] = {"active_now": len(active), "ok": True} + except Exception as e: + results[node["code"]] = {"active_now": len(active), "ok": False, "error": str(e)} + return {"removed_expired": len(expired), "nodes": results}