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