diff --git a/assets/schema.sql b/assets/schema.sql
index b0af1ccc..d557cc17 100644
--- a/assets/schema.sql
+++ b/assets/schema.sql
@@ -158,6 +158,17 @@ CREATE TABLE IF NOT EXISTS servers (
UNIQUE (cluster_name, server_name)
);
+DO $$
+BEGIN
+ IF NOT EXISTS (
+ SELECT 1 FROM information_schema.columns
+ WHERE table_name = 'servers' AND column_name = 'tariff_group'
+ ) THEN
+ ALTER TABLE servers ADD COLUMN tariff_group TEXT;
+ END IF;
+END$$;
+
+
CREATE TABLE IF NOT EXISTS gifts (
gift_id TEXT PRIMARY KEY NOT NULL,
sender_tg_id BIGINT NOT NULL REFERENCES users (tg_id),
diff --git a/database.py b/database.py
index c0d4796c..b9ffb737 100644
--- a/database.py
+++ b/database.py
@@ -1,4 +1,5 @@
import json
+import re
from datetime import datetime
from typing import Any
@@ -56,18 +57,41 @@ async def delete_blocked_user(tg_id: int | list[int], conn: asyncpg.Connection):
async def init_db(file_path: str = "assets/schema.sql"):
- with open(file_path) as file:
+ """
+ Инициализация базы данных: создаёт таблицы и применяет DO-блоки.
+ Поддерживает логгирование как успешных, так и неудачных операций.
+ """
+ with open(file_path, encoding="utf-8") as file:
sql_content = file.read()
conn = await asyncpg.connect(DATABASE_URL)
try:
- await conn.execute(sql_content)
- except Exception as e:
- logger.error(f"Error while executing SQL statement: {e}")
+ blocks = re.split(r';\s*(?=CREATE|DO|\Z)', sql_content, flags=re.IGNORECASE)
+
+ for stmt in blocks:
+ stmt_clean = stmt.strip()
+ if stmt_clean.lower().startswith("create table"):
+ try:
+ await conn.execute(stmt_clean)
+ table_name_match = re.search(r'create table if not exists\s+(\w+)', stmt_clean, re.IGNORECASE)
+ table_name = table_name_match.group(1) if table_name_match else "неизвестно"
+ logger.info(f"[DB INIT] Таблица успешно создана: {table_name}")
+ except Exception as e:
+ logger.error(f"[DB INIT] Ошибка в CREATE TABLE:\n{stmt_clean}\nОшибка: {e}")
+
+ for stmt in blocks:
+ stmt_clean = stmt.strip()
+ if stmt_clean.lower().startswith("do"):
+ try:
+ await conn.execute(stmt_clean)
+ logger.info("[DB INIT] DO-блок выполнен успешно")
+ except Exception as e:
+ logger.error(f"[DB INIT] Ошибка в DO блоке:\n{stmt_clean}\nОшибка: {e}")
+
finally:
- logger.info("Tables created successfully")
await conn.close()
+ logger.info("[DB INIT] Инициализация базы данных завершена")
async def check_unique_server_name(server_name: str, session: Any, cluster_name: str | None = None) -> bool:
@@ -284,6 +308,7 @@ async def store_key(
server_id: str,
session: Any,
remnawave_link: str = None,
+ tariff_id: int | None = None,
):
"""
Сохраняет информацию о ключе в базу данных, если ключ ещё не существует.
@@ -301,8 +326,12 @@ async def store_key(
await session.execute(
"""
- INSERT INTO keys (tg_id, client_id, email, created_at, expiry_time, key, server_id, remnawave_link)
- VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
+ INSERT INTO keys (
+ tg_id, client_id, email, created_at, expiry_time,
+ key, server_id, remnawave_link, tariff_id
+ )
+ VALUES ($1, $2, $3, $4, $5,
+ $6, $7, $8, $9)
""",
tg_id,
client_id,
@@ -312,6 +341,7 @@ async def store_key(
key,
server_id,
remnawave_link,
+ tariff_id,
)
logger.info(f"✅ Ключ сохранён: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
diff --git a/handlers/admin/clusters/clusters_handler.py b/handlers/admin/clusters/clusters_handler.py
index 70063fbf..1a007976 100644
--- a/handlers/admin/clusters/clusters_handler.py
+++ b/handlers/admin/clusters/clusters_handler.py
@@ -488,7 +488,7 @@ async def handle_sync_cluster(callback_query: types.CallbackQuery, callback_data
try:
query_keys = """
- SELECT tg_id, client_id, email, expiry_time, remnawave_link
+ SELECT tg_id, client_id, email, expiry_time, remnawave_link, tariff_id
FROM keys
WHERE server_id = $1
"""
@@ -518,6 +518,7 @@ async def handle_sync_cluster(callback_query: types.CallbackQuery, callback_data
key["client_id"],
key["email"],
key["expiry_time"],
+ plan=key.get("tariff_id"),
session=session,
remnawave_link=key.get("remnawave_link"),
)
diff --git a/handlers/admin/stats/stats_handler.py b/handlers/admin/stats/stats_handler.py
index ad394753..659e662c 100644
--- a/handlers/admin/stats/stats_handler.py
+++ b/handlers/admin/stats/stats_handler.py
@@ -29,96 +29,81 @@ async def handle_stats(callback_query: CallbackQuery, session: Any):
total_users = await session.fetchval("SELECT COUNT(*) FROM users")
total_keys = await session.fetchval("SELECT COUNT(*) FROM keys")
total_referrals = await session.fetchval("SELECT COUNT(*) FROM referrals")
+ users_updated_today = await session.fetchval("SELECT COUNT(*) FROM users WHERE updated_at >= CURRENT_DATE")
- total_payments_today = int(
- await session.fetchval("SELECT COALESCE(SUM(amount), 0) FROM payments WHERE created_at >= CURRENT_DATE")
- )
- total_payments_yesterday = int(
- await session.fetchval("""
- SELECT COALESCE(SUM(amount), 0)
- FROM payments
- WHERE created_at >= CURRENT_DATE - interval '1 day'
- AND created_at < CURRENT_DATE
- """)
- )
- total_payments_week = int(
- await session.fetchval(
- "SELECT COALESCE(SUM(amount), 0) FROM payments WHERE created_at >= date_trunc('week', CURRENT_DATE)"
- )
- )
- total_payments_month = int(
- await session.fetchval(
- "SELECT COALESCE(SUM(amount), 0) FROM payments WHERE created_at >= date_trunc('month', CURRENT_DATE)"
- )
- )
- total_payments_last_month = int(
- await session.fetchval("""
- SELECT COALESCE(SUM(amount), 0)
- FROM payments
- WHERE created_at >= date_trunc('month', CURRENT_DATE - interval '1 month')
- AND created_at < date_trunc('month', CURRENT_DATE)
- """)
- )
+ total_payments_today = int(await session.fetchval("SELECT COALESCE(SUM(amount), 0) FROM payments WHERE created_at >= CURRENT_DATE"))
+ total_payments_yesterday = int(await session.fetchval("""
+ SELECT COALESCE(SUM(amount), 0) FROM payments
+ WHERE created_at >= CURRENT_DATE - interval '1 day' AND created_at < CURRENT_DATE
+ """))
+ total_payments_week = int(await session.fetchval("SELECT COALESCE(SUM(amount), 0) FROM payments WHERE created_at >= date_trunc('week', CURRENT_DATE)"))
+ total_payments_month = int(await session.fetchval("SELECT COALESCE(SUM(amount), 0) FROM payments WHERE created_at >= date_trunc('month', CURRENT_DATE)"))
+ total_payments_last_month = int(await session.fetchval("""
+ SELECT COALESCE(SUM(amount), 0) FROM payments
+ WHERE created_at >= date_trunc('month', CURRENT_DATE - interval '1 month') AND created_at < date_trunc('month', CURRENT_DATE)
+ """))
total_payments_all_time = int(await session.fetchval("SELECT COALESCE(SUM(amount), 0) FROM payments"))
registrations_today = await session.fetchval("SELECT COUNT(*) FROM users WHERE created_at >= CURRENT_DATE")
registrations_yesterday = await session.fetchval("""
SELECT COUNT(*) FROM users
- WHERE created_at >= CURRENT_DATE - interval '1 day'
- AND created_at < CURRENT_DATE
+ WHERE created_at >= CURRENT_DATE - interval '1 day' AND created_at < CURRENT_DATE
""")
- registrations_week = await session.fetchval(
- "SELECT COUNT(*) FROM users WHERE created_at >= date_trunc('week', CURRENT_DATE)"
- )
- registrations_month = await session.fetchval(
- "SELECT COUNT(*) FROM users WHERE created_at >= date_trunc('month', CURRENT_DATE)"
- )
+ registrations_week = await session.fetchval("SELECT COUNT(*) FROM users WHERE created_at >= date_trunc('week', CURRENT_DATE)")
+ registrations_month = await session.fetchval("SELECT COUNT(*) FROM users WHERE created_at >= date_trunc('month', CURRENT_DATE)")
registrations_last_month = await session.fetchval("""
SELECT COUNT(*) FROM users
- WHERE created_at >= date_trunc('month', CURRENT_DATE - interval '1 month')
- AND created_at < date_trunc('month', CURRENT_DATE)
+ WHERE created_at >= date_trunc('month', CURRENT_DATE - interval '1 month') AND created_at < date_trunc('month', CURRENT_DATE)
""")
- all_keys = await session.fetch("SELECT created_at, expiry_time FROM keys")
-
- def count_subscriptions_by_duration(keys):
- periods = {"trial": 0, "1": 0, "3": 0, "6": 0, "12": 0}
- for key in keys:
- try:
- duration_days = (key["expiry_time"] - key["created_at"]) / (1000 * 60 * 60 * 24)
-
- if duration_days <= 29:
- periods["trial"] += 1
- elif duration_days <= 89:
- periods["1"] += 1
- elif duration_days <= 179:
- periods["3"] += 1
- elif duration_days <= 359:
- periods["6"] += 1
- else:
- periods["12"] += 1
- except Exception as e:
- logger.error(f"Error processing key duration: {e}")
- continue
- return periods
-
- subs_all_time = count_subscriptions_by_duration(all_keys)
-
- users_updated_today = await session.fetchval("SELECT COUNT(*) FROM users WHERE updated_at >= CURRENT_DATE")
-
- active_keys = await session.fetchval(
- "SELECT COUNT(*) FROM keys WHERE expiry_time > $1",
- int(datetime.utcnow().timestamp() * 1000),
- )
+ active_keys = await session.fetchval("SELECT COUNT(*) FROM keys WHERE expiry_time > $1", int(datetime.utcnow().timestamp() * 1000))
expired_keys = total_keys - active_keys
+ tariffs = await session.fetch("SELECT id, name, duration_days FROM tariffs WHERE is_active = TRUE")
+ tariff_map = {t["id"]: t["name"] for t in tariffs}
+ durations = [(t["id"], t["name"], t["duration_days"]) for t in tariffs]
+
+ tariff_counter: dict[str, int] = {}
+
+ keys_with_tariffs = await session.fetch("SELECT tariff_id FROM keys WHERE tariff_id IS NOT NULL")
+ for row in keys_with_tariffs:
+ name = tariff_map.get(row["tariff_id"], "Неизвестно")
+ tariff_counter[name] = tariff_counter.get(name, 0) + 1
+
+ keys_without_tariffs = await session.fetch("SELECT created_at, expiry_time FROM keys WHERE tariff_id IS NULL")
+ for row in keys_without_tariffs:
+ duration_days = (row["expiry_time"] - row["created_at"]) / (1000 * 60 * 60 * 24)
+ if durations:
+ closest = min(durations, key=lambda t: abs(t[2] - duration_days))
+ name = closest[1]
+ else:
+ name = "Неизвестно"
+ tariff_counter[name] = tariff_counter.get(name, 0) + 1
+
+ tariff_order = {t["name"]: t["id"] for t in sorted(tariffs, key=lambda t: t["id"])}
+ tariff_stats_text = "\n".join(
+ f" • {name}: {tariff_counter[name]}"
+ for name in sorted(tariff_counter.keys(), key=lambda name: tariff_order.get(name, float('inf')))
+ )
+
+ if not tariff_stats_text:
+ tariff_stats_text = " • Нет активных тарифов"
+
+
hot_leads_count = await session.fetchval("""
SELECT COUNT(DISTINCT u.tg_id)
FROM users u
JOIN payments p ON u.tg_id = p.tg_id
LEFT JOIN keys k ON u.tg_id = k.tg_id
- WHERE p.status = 'success'
- AND k.tg_id IS NULL
+ WHERE p.status = 'success' AND k.tg_id IS NULL
+ """)
+
+ trial_only_count = await session.fetchval("""
+ SELECT COUNT(DISTINCT k.tg_id)
+ FROM keys k
+ LEFT JOIN tariffs t ON k.tariff_id = t.id
+ LEFT JOIN payments p ON k.tg_id = p.tg_id
+ WHERE p.id IS NULL
""")
moscow_tz = pytz.timezone("Europe/Moscow")
@@ -141,30 +126,27 @@ async def handle_stats(callback_query: CallbackQuery, session: Any):
f"├ 📦 Всего сгенерировано: {total_keys}\n"
f"├ ✅ Активных: {active_keys}\n"
f"├ ❌ Просроченных: {expired_keys}\n"
- f"└ 📋 По срокам:\n"
- f" • 🎁 Триал: {subs_all_time['trial']}\n"
- f" • 🗓️ 1 мес: {subs_all_time['1']}\n"
- f" • 🗓️ 3 мес: {subs_all_time['3']}\n"
- f" • 🗓️ 6 мес: {subs_all_time['6']}\n"
- f" • 🗓️ 12 мес: {subs_all_time['12']}\n\n"
+ f"├ 🎁 Только триал: {trial_only_count}\n"
+ f"└ 📋 По тарифам:\n{tariff_stats_text}\n\n"
"💰 Финансы:\n"
f"├ 📅 За день: {total_payments_today} ₽\n"
f"├ 📆 Вчера: {total_payments_yesterday} ₽\n"
f"├ 📆 За неделю: {total_payments_week} ₽\n"
f"├ 📆 За месяц: {total_payments_month} ₽\n"
- f"├ 📆 За прошлый месяц: {total_payments_last_month} ₽\n"
+ f"├ 📆 Прошлый месяц: {total_payments_last_month} ₽\n"
f"└ 🏦 Всего: {total_payments_all_time} ₽\n\n"
f"🔥 Горящие лиды: {hot_leads_count} (платили, но не продлили)\n\n"
f"⏱️ Последнее обновление: {update_time}"
)
await callback_query.message.edit_text(text=stats_message, reply_markup=build_stats_kb())
+
except TelegramBadRequest as e:
if "message is not modified" not in str(e):
logger.error(f"Error in user_stats_menu: {e}")
except Exception as e:
logger.error(f"Error in user_stats_menu: {e}")
- await callback_query.answer("Произошла ошибка при получении статистики", show_alert=True)
+ await callback_query.answer("\u041fроизошла ошибка при получении статистики", show_alert=True)
@router.callback_query(
diff --git a/handlers/keys/key_utils.py b/handlers/keys/key_utils.py
index 578ed73a..55a1522d 100644
--- a/handlers/keys/key_utils.py
+++ b/handlers/keys/key_utils.py
@@ -162,6 +162,7 @@ async def create_key_on_cluster(
server_id=server_id_to_store,
session=session,
remnawave_link=remnawave_key,
+ tariff_id=plan,
)
except Exception as e:
diff --git a/utils/csv_export.py b/utils/csv_export.py
index 565909fc..143724e7 100644
--- a/utils/csv_export.py
+++ b/utils/csv_export.py
@@ -173,26 +173,30 @@ async def export_hot_leads_csv(session: Any) -> BufferedInputFile:
async def export_keys_csv(session) -> BufferedInputFile:
"""
- Экспорт подписок в CSV с нормальными датами.
+ Экспорт подписок в CSV с нормальными датами и тарифом.
"""
keys = await session.fetch("""
- SELECT tg_id, client_id, email, created_at, expiry_time, key, server_id, is_frozen, alias
- FROM keys
- ORDER BY created_at ASC
+ SELECT k.tg_id, k.client_id, k.email, k.created_at, k.expiry_time,
+ k.key, k.server_id, k.is_frozen, k.alias,
+ t.name AS tariff_name
+ FROM keys k
+ LEFT JOIN tariffs t ON k.tariff_id = t.id
+ ORDER BY k.created_at ASC
""")
buffer = StringIO()
- buffer.write("tg_id,client_id,email,created_at,expiry_time,key,server_id,is_frozen,alias\n")
+ buffer.write("tg_id,client_id,email,created_at,expiry_time,key,server_id,is_frozen,alias,tariff\n")
for row in keys:
created_at = datetime.utcfromtimestamp(row["created_at"] / 1000).strftime("%Y-%m-%d %H:%M:%S")
expiry_time = datetime.utcfromtimestamp(row["expiry_time"] / 1000).strftime("%Y-%m-%d %H:%M:%S")
+ tariff = row["tariff_name"] or "—"
buffer.write(
f"{row['tg_id']},{row['client_id']},{row['email']},"
f"{created_at},{expiry_time},{row['key']},"
- f"{row['server_id']},{row['is_frozen']},{row['alias'] or ''}\n"
+ f"{row['server_id']},{row['is_frozen']},{row['alias'] or ''},{tariff}\n"
)
buffer.seek(0)
- return BufferedInputFile(file=buffer.getvalue().encode("utf-8-sig"), filename="keys_export.csv")
+ return BufferedInputFile(file=buffer.getvalue().encode("utf-8-sig"), filename="keys_export.csv")
\ No newline at end of file