commit e4c2012cf88caf5f892d41b89fa55e7ac9be2df2 Author: savsis Date: Thu Sep 10 17:45:36 2026 +0500 core: env-driven config, sqlite schema, xray/node management, subscription links Co-Authored-By: Claude Sonnet 5 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}