From f9ee4f593385a90bd29a39837a571e6553cf87e1 Mon Sep 17 00:00:00 2001 From: savsis Date: Fri, 11 Sep 2026 22:20:54 +0500 Subject: [PATCH] =?UTF-8?q?fix:=2018-point=20audit=20pass=20=E2=80=94=20pa?= =?UTF-8?q?yment=20races,=20hwid=20limit=20bugs,=20blocking=20SSH/HTTP=20i?= =?UTF-8?q?n=20event=20loops,=20N+1=20queries,=20ssh=20host-key=20pinning,?= =?UTF-8?q?=20dead=20code?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit payments: _grant_paid_subscription now validates plan/node exist before marking a payment paid instead of after (was leaving charged-but-ungranted payments with no error trail); mark_payment_paid is now a single atomic UPDATE ... WHERE status='pending' instead of check-then-act, closing a double-grant race between webhooks and the periodic reconciler; yookassa webhook now re-verifies payment status server-side via the API instead of trusting the posted body (platega already had HMAC verification). hwid: 'user["hwid_limit"] or FALLBACK' treated an explicit 0 (admin fully blocking a user) as unset — now an explicit None check. Device count-check and insert are now one atomic transaction (db.add_device_if_under_limit) instead of two raceable statements. perf: payment webhooks and _grant_paid_subscription's SSH/HTTP calls now run via asyncio.to_thread instead of blocking the event loop; same for bot.py's periodic_sync/reconcile_pending_payments and the manual admin sync button. Admin endpoints (traffic/subscriptions/payments/gift-codes/ user-card) now resolve node labels from one db.list_nodes() call instead of a fresh db.get_node() per row. revoke/reset-traffic use a direct PK lookup instead of scanning up to 5000 rows. Dashboard now asks the API for 8 rows instead of fetching 200 and slicing client-side. security: mbs.db (and -wal/-shm) now chmod 600 right after creation — it held session tokens and subscription bearer tokens world-readable by default. Node SSH connections now pin host keys via a persisted known_hosts file (TOFU) instead of accepting any key on every connection. delete_node now refuses to delete a node with active subscriptions instead of silently orphaning their xray clients. deadcode: removed unused xray_manager.list_client_ids and admin.html's superseded staggerReveal (rows animate via rowAttr() inline now). Also guards gift-code redemption against a plan/node deleted after the code was created (was an unhandled KeyError/TypeError crash). Co-Authored-By: Claude Sonnet 5 --- admin.html | 8 +------ api.py | 64 ++++++++++++++++++++++++++++--------------------- bot.py | 25 ++++++++++--------- db.py | 50 ++++++++++++++++++++++++++++++++++---- nodeprov.py | 6 +++++ xray_manager.py | 9 ------- 6 files changed, 103 insertions(+), 59 deletions(-) diff --git a/admin.html b/admin.html index 5e6cd1f..a0cbf9d 100644 --- a/admin.html +++ b/admin.html @@ -873,11 +873,6 @@ function statusBadge(active, daysLeft) { if (daysLeft <= 2) return '' + daysLeft + ' дн.'; return '' + daysLeft + ' дн.'; } -function staggerReveal(container) { - const items = container.querySelectorAll(".reveal"); - items.forEach((el, i) => { el.style.animationDelay = (i * 0.04) + "s"; }); -} - async function loadDashboard() { const stats = await api("/admin/api/stats"); const grid = document.getElementById("stat-grid"); @@ -888,9 +883,8 @@ async function loadDashboard() { ["Гифт-коды (созд./исп.)", stats.gifts_created + " / " + stats.gifts_used, "var(--pink)", ICONS.gift], ].map(([l, v, c, ic], i) => statCard(l, v, c, ic, i * 0.05)).join(""); - const subs = await api("/admin/api/subscriptions"); + const recent = await api("/admin/api/subscriptions?limit=8"); const body = document.getElementById("recent-subs-body"); - const recent = subs.slice(0, 8); body.innerHTML = recent.length ? recent.map((s, i) => ` ${s.username ? "@" + esc(s.username) : "tg" + s.tg_id}${esc(s.node_label)}${esc(s.plan_label)} ${fmtDate(s.expires_at)}${statusBadge(s.active, s.days_left)} diff --git a/api.py b/api.py index f73b654..9f5cac3 100644 --- a/api.py +++ b/api.py @@ -1,3 +1,4 @@ +import asyncio import datetime import json import os @@ -241,16 +242,15 @@ def get_subscription(token: str, request: Request): hwid = request.headers.get("x-hwid", "") if not HWID_RE.match(hwid): raise HTTPException(404, "hwid required") - if not db.get_device(user["tg_id"], hwid): - limit = user["hwid_limit"] or HWID_FALLBACK_LIMIT - if db.count_devices(user["tg_id"]) >= limit: - raise HTTPException(404, "device limit reached", headers={"x-hwid-max-devices-reached": "true"}) - db.add_device( - user["tg_id"], hwid, - request.headers.get("x-device-os"), - request.headers.get("x-device-model"), - ua, - ) + limit = user["hwid_limit"] if user["hwid_limit"] is not None else HWID_FALLBACK_LIMIT + _, allowed = db.add_device_if_under_limit( + user["tg_id"], hwid, limit, + request.headers.get("x-device-os"), + request.headers.get("x-device-model"), + ua, + ) + if not allowed: + raise HTTPException(404, "device limit reached", headers={"x-hwid-max-devices-reached": "true"}) content = links.build_subscription_text(subs) return Response(content=content, media_type="text/plain") @@ -305,13 +305,16 @@ def _tg_send_message(tg_id: int, text: str): def _grant_paid_subscription(payment_id: str): - payment = db.mark_payment_paid(payment_id) - if not payment: + payment = db.get_payment(payment_id) + if not payment or payment["status"] == "paid": return plan = PLANS_BY_CODE.get(payment["plan"]) node = db.get_node(payment["node"]) if not plan or not node: return + payment = db.mark_payment_paid(payment_id) + if not payment: + return sub = db.create_subscription(payment["tg_id"], payment["node"], plan["days"], payment["plan"], source="payment") xray_manager.add_client_to_node(node, sub["uuid"], email=sub["uuid"]) user = db.get_or_create_user(payment["tg_id"], None) @@ -343,9 +346,10 @@ def _check_and_reconcile_payment(payment: dict) -> str: @app.get("/admin/api/payments") def admin_list_payments(request: Request): require_admin(request) + nodes_by_code = {n["code"]: n for n in db.list_nodes()} out = [] for p in db.list_payments(): - node = db.get_node(p["node"]) + node = nodes_by_code.get(p["node"]) plan = PLANS_BY_CODE.get(p["plan"]) out.append({ **p, @@ -374,12 +378,16 @@ async def yookassa_webhook(request: Request): if not payments.verify_yookassa_notification(body): raise HTTPException(400, "unexpected event") obj = body.get("object", {}) - if obj.get("status") != "succeeded": - return {"ok": True} payment_id = (obj.get("metadata") or {}).get("payment_id") if not payment_id: raise HTTPException(400, "missing payment_id") - _grant_paid_subscription(payment_id) + payment = db.get_payment(payment_id) + if not payment or not payment.get("external_id"): + raise HTTPException(400, "unknown payment") + status = await asyncio.to_thread(payments.check_yookassa_payment, payment["external_id"]) + if status not in payments.PAID_STATUSES: + return {"ok": True} + await asyncio.to_thread(_grant_paid_subscription, payment_id) return {"ok": True} @@ -394,7 +402,7 @@ async def platega_webhook(request: Request): payment_id = body.get("id") or body.get("paymentId") if status not in ("succeeded", "success", "paid") or not payment_id: return {"ok": True} - _grant_paid_subscription(payment_id) + await asyncio.to_thread(_grant_paid_subscription, payment_id) return {"ok": True} @@ -519,12 +527,13 @@ def admin_traffic(request: Request): total_down = sum(v["down"] for v in all_stats.values()) subs = db.list_all_subscriptions(limit=5000) + nodes_by_code = {n["code"]: n for n in db.list_nodes()} per_sub = [] for s in subs: st = all_stats.get(s["uuid"]) if not st: continue - node = db.get_node(s["node"]) + node = nodes_by_code.get(s["node"]) per_sub.append({ "uuid": s["uuid"], "username": ("@" + s["username"]) if s.get("username") else f"tg{s['tg_id']}", @@ -544,12 +553,13 @@ def admin_traffic(request: Request): @app.get("/admin/api/subscriptions") -def admin_subscriptions(request: Request): +def admin_subscriptions(request: Request, limit: int = 200): require_admin(request) - subs = db.list_all_subscriptions() + subs = db.list_all_subscriptions(limit=limit) + nodes_by_code = {n["code"]: n for n in db.list_nodes()} out = [] for s in subs: - node = db.get_node(s["node"]) + node = nodes_by_code.get(s["node"]) plan = PLANS_BY_CODE.get(s["plan"]) out.append({ **s, @@ -563,8 +573,7 @@ def admin_subscriptions(request: Request): @app.post("/admin/api/subscriptions/{uuid}/revoke") def admin_revoke_subscription(uuid: str, request: Request): require_admin(request) - subs = db.list_all_subscriptions(limit=5000) - sub = next((s for s in subs if s["uuid"] == uuid), None) + sub = db.get_subscription(uuid) if not sub: raise HTTPException(404, "not found") node = db.get_node(sub["node"]) @@ -577,8 +586,7 @@ def admin_revoke_subscription(uuid: str, request: Request): @app.post("/admin/api/subscriptions/{uuid}/reset-traffic") def admin_reset_traffic(uuid: str, request: Request): require_admin(request) - subs = db.list_all_subscriptions(limit=5000) - sub = next((s for s in subs if s["uuid"] == uuid), None) + sub = db.get_subscription(uuid) if not sub: raise HTTPException(404, "not found") node = db.get_node(sub["node"]) @@ -595,9 +603,10 @@ def admin_user_card(tg_id: int, request: Request): if not user: raise HTTPException(404, "not found") subs = db.list_subscriptions_for_user(tg_id) + nodes_by_code = {n["code"]: n for n in db.list_nodes()} out_subs = [] for s in subs: - node = db.get_node(s["node"]) + node = nodes_by_code.get(s["node"]) plan = PLANS_BY_CODE.get(s["plan"]) out_subs.append({ **s, @@ -662,9 +671,10 @@ def admin_set_hwid_limit(tg_id: int, request: Request, body: dict = Body(...)): def admin_gift_codes(request: Request): require_admin(request) codes = db.list_gift_codes() + nodes_by_code = {n["code"]: n for n in db.list_nodes()} out = [] for c in codes: - node = db.get_node(c["node"]) + node = nodes_by_code.get(c["node"]) plan = PLANS_BY_CODE.get(c["plan"]) out.append({ **c, diff --git a/bot.py b/bot.py index 3d674ec..3e95e6e 100644 --- a/bot.py +++ b/bot.py @@ -96,10 +96,13 @@ async def start_deeplink(message: Message, command: CommandObject): if err == "already_used": await message.answer("Этот код уже был использован.") return await send_main_menu(message) - plan = PLANS_BY_CODE[gift["plan"]] - sub = db.create_subscription(message.from_user.id, gift["node"], plan["days"], plan["code"], source="gift", ) + plan = PLANS_BY_CODE.get(gift["plan"]) gift_node = db.get_node(gift["node"]) - xray_manager.add_client_to_node(gift_node, sub["uuid"], email=sub["uuid"]) + if not plan or not gift_node: + await message.answer("Этот подарок больше недоступен.") + return await send_main_menu(message) + sub = db.create_subscription(message.from_user.id, gift["node"], plan["days"], plan["code"], source="gift", ) + await asyncio.to_thread(xray_manager.add_client_to_node, gift_node, sub["uuid"], email=sub["uuid"]) await message.answer( f"Подарок активирован\n\n" f"Сервер: {gift_node['label']}\n" @@ -171,7 +174,7 @@ async def cb_plan(cb: CallbackQuery): user = db.get_or_create_user(cb.from_user.id, cb.from_user.username) sub = db.create_subscription(cb.from_user.id, node_code, plan["days"], plan_code, source="bot") node_row = db.get_node(node_code) - xray_manager.add_client_to_node(node_row, sub["uuid"], email=sub["uuid"]) + await asyncio.to_thread(xray_manager.add_client_to_node, node_row, sub["uuid"], email=sub["uuid"]) kb = connect_kb(user["token"], extra_rows=[ [InlineKeyboardButton(text="Моя подписка", callback_data="menu:mysub")], [InlineKeyboardButton(text="В меню", callback_data="menu:main")], @@ -309,7 +312,7 @@ async def cb_admin_stats(cb: CallbackQuery): async def cb_admin_sync(cb: CallbackQuery): if not is_admin(cb.from_user.id): return await cb.answer("Нет доступа", show_alert=True) - result = xray_manager.sync_all() + result = await asyncio.to_thread(xray_manager.sync_all) kb = InlineKeyboardMarkup(inline_keyboard=[[InlineKeyboardButton(text="В админку", callback_data="menu:admin")]]) await cb.message.edit_text( f"Синхронизация xray выполнена.\nАктивно клиентов: {result['active_now']}\n" @@ -326,19 +329,19 @@ async def reconcile_pending_payments(): if payment["status"] != "pending" or not payment.get("external_id"): continue try: - status = payments.check_payment_status(payment["provider"], payment["external_id"]) + status = await asyncio.to_thread(payments.check_payment_status, payment["provider"], payment["external_id"]) except Exception: continue if status in payments.PAID_STATUSES: - granted = db.mark_payment_paid(payment["id"]) - if not granted: - continue plan = PLANS_BY_CODE.get(payment["plan"]) node_row = db.get_node(payment["node"]) if not plan or not node_row: continue + granted = db.mark_payment_paid(payment["id"]) + if not granted: + continue sub = db.create_subscription(payment["tg_id"], payment["node"], plan["days"], payment["plan"], source="payment") - xray_manager.add_client_to_node(node_row, sub["uuid"], email=sub["uuid"]) + await asyncio.to_thread(xray_manager.add_client_to_node, node_row, sub["uuid"], email=sub["uuid"]) user = db.get_or_create_user(payment["tg_id"], None) try: await bot.send_message( @@ -357,7 +360,7 @@ async def reconcile_pending_payments(): async def periodic_sync(): while True: try: - xray_manager.sync_all() + await asyncio.to_thread(xray_manager.sync_all) except Exception: log.exception("periodic sync failed") try: diff --git a/db.py b/db.py index b16ecfd..fd59685 100644 --- a/db.py +++ b/db.py @@ -1,3 +1,4 @@ +import os import sqlite3 import secrets import datetime @@ -145,6 +146,13 @@ def init_db(): conn.executescript(SCHEMA) _migrate() _seed_local_node() + for suffix in ("", "-wal", "-shm"): + path = DB_PATH + suffix + if os.path.exists(path): + try: + os.chmod(path, 0o600) + except OSError: + pass def _seed_local_node(): @@ -231,6 +239,12 @@ def delete_node(code: str): if code == "de1": raise ValueError("cannot delete the local node") with get_conn() as conn: + active = conn.execute( + "SELECT COUNT(*) c FROM subscriptions WHERE node=? AND active=1 AND expires_at>?", + (code, now_iso()), + ).fetchone()["c"] + if active: + raise ValueError(f"node has {active} active subscriptions, revoke them first") conn.execute("DELETE FROM nodes WHERE code=?", (code,)) @@ -360,6 +374,16 @@ def revoke_subscription(client_uuid: str): conn.execute("UPDATE subscriptions SET active=0 WHERE uuid=?", (client_uuid,)) +def get_subscription(client_uuid: str): + with get_conn() as conn: + row = conn.execute( + "SELECT s.*, u.username FROM subscriptions s " + "LEFT JOIN users u ON u.tg_id = s.tg_id WHERE s.uuid=?", + (client_uuid,), + ).fetchone() + return dict(row) if row else None + + def get_user(tg_id: int): with get_conn() as conn: row = conn.execute("SELECT * FROM users WHERE tg_id=?", (tg_id,)).fetchone() @@ -427,13 +451,12 @@ def set_payment_external(payment_id: str, external_id: str, pay_url: str): def mark_payment_paid(payment_id: str): with get_conn() as conn: - row = conn.execute("SELECT status FROM payments WHERE id=?", (payment_id,)).fetchone() - if not row or row["status"] == "paid": - return None - conn.execute( - "UPDATE payments SET status='paid', paid_at=? WHERE id=?", + cur = conn.execute( + "UPDATE payments SET status='paid', paid_at=? WHERE id=? AND status='pending'", (now_iso(), payment_id), ) + if cur.rowcount == 0: + return None return get_payment(payment_id) @@ -483,6 +506,23 @@ def add_device(tg_id: int, hwid: str, device_os: str | None, device_model: str | return get_device(tg_id, hwid) +def add_device_if_under_limit(tg_id: int, hwid: str, limit: int, device_os: str | None, device_model: str | None, user_agent: str | None): + with get_conn() as conn: + conn.execute("BEGIN IMMEDIATE") + existing = conn.execute("SELECT * FROM devices WHERE tg_id=? AND hwid=?", (tg_id, hwid)).fetchone() + if existing: + return dict(existing), True + count = conn.execute("SELECT COUNT(*) c FROM devices WHERE tg_id=?", (tg_id,)).fetchone()["c"] + if count >= limit: + return None, False + conn.execute( + "INSERT INTO devices (tg_id, hwid, device_os, device_model, user_agent, first_seen) VALUES (?,?,?,?,?,?)", + (tg_id, hwid, device_os, device_model, user_agent, now_iso()), + ) + row = conn.execute("SELECT * FROM devices WHERE tg_id=? AND hwid=?", (tg_id, hwid)).fetchone() + return dict(row), True + + def delete_device(device_id: int): with get_conn() as conn: conn.execute("DELETE FROM devices WHERE id=?", (device_id,)) diff --git a/nodeprov.py b/nodeprov.py index 32f547b..9b4fb31 100644 --- a/nodeprov.py +++ b/nodeprov.py @@ -8,6 +8,7 @@ import paramiko from config import PANEL_DOMAIN MGMT_KEY_PATH = "/root/.ssh/mbs_nodes_ed25519" +MGMT_KNOWN_HOSTS_PATH = "/root/.ssh/mbs_nodes_known_hosts" LOCAL_TAGS = {"vless-tcp-reality", "vless-grpc-reality", "vless-xhttp-reality", "vless-ws-tls"} TAG_FLOW = {"vless-tcp-reality": "xtls-rprx-vision"} @@ -238,9 +239,14 @@ def render_install_script(node: dict) -> str: def _mgmt_connect(address: str, ssh_port: int = 22) -> paramiko.SSHClient: client = paramiko.SSHClient() + try: + client.load_host_keys(MGMT_KNOWN_HOSTS_PATH) + except IOError: + pass 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) + client.save_host_keys(MGMT_KNOWN_HOSTS_PATH) return client diff --git a/xray_manager.py b/xray_manager.py index 0c75c20..4c1d04e 100644 --- a/xray_manager.py +++ b/xray_manager.py @@ -75,15 +75,6 @@ def remove_client(client_uuid: str): _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