core: env-driven config, sqlite schema, xray/node management, subscription links
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
commit
e4c2012cf8
8 changed files with 1182 additions and 0 deletions
31
.env.example
Normal file
31
.env.example
Normal file
|
|
@ -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
|
||||||
9
.gitignore
vendored
Normal file
9
.gitignore
vendored
Normal file
|
|
@ -0,0 +1,9 @@
|
||||||
|
.env
|
||||||
|
*.db
|
||||||
|
*.db-journal
|
||||||
|
*.db-wal
|
||||||
|
*.db-shm
|
||||||
|
__pycache__/
|
||||||
|
*.pyc
|
||||||
|
venv/
|
||||||
|
.claude/
|
||||||
100
config.py
Normal file
100
config.py
Normal file
|
|
@ -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}
|
||||||
344
db.py
Normal file
344
db.py
Normal file
|
|
@ -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,
|
||||||
|
}
|
||||||
112
links.py
Normal file
112
links.py
Normal file
|
|
@ -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()
|
||||||
374
nodeprov.py
Normal file
374
nodeprov.py
Normal file
|
|
@ -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
|
||||||
4
requirements.txt
Normal file
4
requirements.txt
Normal file
|
|
@ -0,0 +1,4 @@
|
||||||
|
aiogram==3.15.0
|
||||||
|
fastapi==0.115.6
|
||||||
|
uvicorn[standard]==0.32.1
|
||||||
|
paramiko==3.5.0
|
||||||
208
xray_manager.py
Normal file
208
xray_manager.py
Normal file
|
|
@ -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}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue