diff --git a/.env.example b/.env.example
index 9d7cc33..50fec76 100644
--- a/.env.example
+++ b/.env.example
@@ -18,6 +18,7 @@ SERVER_STATUS_URL=https://status.yourdomain.tld/status/your_service
TERMS_OF_SERVICE_URL=https://example.com/tos
SUBSCRIPTION_MINI_APP_URL=
START_COMMAND_DESCRIPTION=
+DISABLE_WELCOME_MESSAGE=
# Webhook Base URL (used for Telegram and payment providers)
WEBHOOK_BASE_URL=https://webhooks.yourdomain.tld
diff --git a/bot/handlers/admin/broadcast.py b/bot/handlers/admin/broadcast.py
index 6d35275..bb2aa68 100644
--- a/bot/handlers/admin/broadcast.py
+++ b/bot/handlers/admin/broadcast.py
@@ -1,6 +1,7 @@
import logging
import asyncio
from aiogram import Router, F, types, Bot
+from aiogram.exceptions import TelegramRetryAfter
from aiogram.fsm.context import FSMContext
from typing import Optional
@@ -17,6 +18,7 @@ from bot.keyboards.inline.admin_keyboards import (
get_admin_panel_keyboard,
)
from bot.middlewares.i18n import JsonI18n
+from bot.utils.message_queue import get_queue_manager
router = Router(name="admin_broadcast_router")
@@ -171,22 +173,30 @@ async def confirm_broadcast_callback_handler(
f"Admin {admin_user.id} broadcasting '{text[:50]}...' to {len(user_ids)} users."
)
+ # Get message queue manager
+ queue_manager = get_queue_manager()
+ if not queue_manager:
+ await callback.message.edit_text("❌ Ошибка: система очередей не инициализирована", reply_markup=None)
+ return
+
+ # Queue all messages for sending
for uid in user_ids:
try:
- await bot.send_message(
+ await queue_manager.send_message(
chat_id=uid,
text=text,
entities=entities,
)
sent_count += 1
-
+
+ # Log successful queuing
await message_log_dal.create_message_log(
session,
{
"user_id": admin_user.id,
"telegram_username": admin_user.username,
"telegram_first_name": admin_user.first_name,
- "event_type": "admin_broadcast_sent",
+ "event_type": "admin_broadcast_queued",
"content": f"To user {uid}: {text[:70]}...",
"is_admin_event": True,
"target_user_id": uid,
@@ -195,7 +205,7 @@ async def confirm_broadcast_callback_handler(
except Exception as e:
failed_count += 1
logging.warning(
- f"Failed to send broadcast to {uid}: {type(e).__name__} – {e}"
+ f"Failed to queue broadcast to {uid}: {type(e).__name__} – {e}"
)
await message_log_dal.create_message_log(
session,
@@ -209,7 +219,6 @@ async def confirm_broadcast_callback_handler(
"target_user_id": uid,
},
)
- await asyncio.sleep(0.05)
try:
await session.commit()
@@ -217,7 +226,18 @@ async def confirm_broadcast_callback_handler(
await session.rollback()
logging.error(f"Error committing broadcast logs: {e_commit}")
- result_message = _("admin_broadcast_finished_stats", sent_count=sent_count, failed_count=failed_count)
+ # Get queue stats for detailed report
+ queue_stats = queue_manager.get_queue_stats()
+
+ result_message = f"""🚀 Рассылка поставлена в очередь!
+📤 В очередь добавлено: {sent_count}
+❌ Ошибок: {failed_count}
+
+📊 Статус очередей:
+👥 Очередь пользователей: {queue_stats['user_queue_size']} сообщений
+📢 Очередь групп: {queue_stats['group_queue_size']} сообщений
+
+ℹ️ Сообщения будут отправлены автоматически с соблюдением лимитов Telegram."""
await callback.message.answer(
result_message,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
diff --git a/bot/handlers/admin/common.py b/bot/handlers/admin/common.py
index e051dd6..a43c57d 100644
--- a/bot/handlers/admin/common.py
+++ b/bot/handlers/admin/common.py
@@ -14,6 +14,7 @@ from bot.keyboards.inline.admin_keyboards import (
from bot.middlewares.i18n import JsonI18n
from bot.services.panel_api_service import PanelApiService
from bot.services.subscription_service import SubscriptionService
+from bot.utils.message_queue import get_queue_manager
from . import broadcast as admin_broadcast_handlers
from .promo import create as admin_promo_create_handlers
@@ -120,6 +121,8 @@ async def admin_panel_actions_callback_handler(
panel_service=panel_service,
session=session)
await callback.answer(_("admin_sync_initiated_from_panel"))
+ elif action == "queue_status":
+ await show_queue_status_handler(callback, i18n_data)
elif action == "main":
try:
await callback.message.edit_text(
@@ -192,3 +195,52 @@ async def admin_section_handler(callback: types.CallbackQuery, state: FSMContext
reply_markup=get_admin_panel_keyboard(i18n, current_lang, settings)
)
await callback.answer()
+
+
+async def show_queue_status_handler(callback: types.CallbackQuery, i18n_data: dict):
+ """Show message queue status to admin"""
+ current_lang = i18n_data.get("current_language", "ru")
+ i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
+ if not i18n or not callback.message:
+ await callback.answer("Error processing request.", show_alert=True)
+ return
+ _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
+
+ queue_manager = get_queue_manager()
+ if not queue_manager:
+ from aiogram.utils.keyboard import InlineKeyboardBuilder
+ await callback.message.edit_text(
+ "❌ Система очередей не инициализирована",
+ reply_markup=InlineKeyboardBuilder().button(
+ text=_("back_to_admin_panel_button"),
+ callback_data="admin_action:main"
+ ).as_markup()
+ )
+ await callback.answer()
+ return
+
+ try:
+ stats = queue_manager.get_queue_stats()
+
+ message_text = _(
+ "admin_queue_status_info",
+ user_queue_size=stats['user_queue_size'],
+ user_processing="✅ Да" if stats['user_queue_processing'] else "❌ Нет",
+ user_recent=stats['user_recent_sends'],
+ group_queue_size=stats['group_queue_size'],
+ group_processing="✅ Да" if stats['group_queue_processing'] else "❌ Нет",
+ group_recent=stats['group_recent_sends']
+ )
+
+ from bot.keyboards.inline.admin_keyboards import get_back_to_admin_panel_keyboard
+
+ await callback.message.edit_text(
+ message_text,
+ reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
+ parse_mode="HTML"
+ )
+ await callback.answer()
+
+ except Exception as e:
+ logging.error(f"Error getting queue status: {e}")
+ await callback.answer("❌ Ошибка получения статуса очередей", show_alert=True)
diff --git a/bot/handlers/admin/promo/bulk.py b/bot/handlers/admin/promo/bulk.py
index 8a2f7c7..1b6b302 100644
--- a/bot/handlers/admin/promo/bulk.py
+++ b/bot/handlers/admin/promo/bulk.py
@@ -1,6 +1,8 @@
import logging
import random
import string
+import csv
+import io
from aiogram import Router, F, types
from aiogram.filters import StateFilter
from aiogram.fsm.context import FSMContext
@@ -421,15 +423,63 @@ async def create_bulk_promo_codes_final(callback_or_message,
)
)
+ # Create CSV file with promo codes if any were created
+ csv_file = None
if created_codes:
- success_lines.append("\n🎟 Созданные коды:")
- # Show first 20 codes, then indicate if there are more
- codes_to_show = created_codes[:20]
- for code in codes_to_show:
- success_lines.append(f"{code}")
+ success_lines.append(f"\n🎟 Создано {len(created_codes)} промокодов")
+ success_lines.append("📄 CSV файл с промокодами отправлен отдельным сообщением")
- if len(created_codes) > 20:
- success_lines.append(f"... и еще {len(created_codes) - 20} кодов")
+ # Create CSV file
+ output = io.StringIO()
+ writer = csv.writer(output)
+
+ # CSV headers
+ writer.writerow([
+ "Промокод", "Бонусные дни", "Макс. активации", "Действителен до",
+ "Команда для старта", "Ссылка для активации"
+ ])
+
+ # Get real bot username
+ bot_username = 'your_bot' # fallback
+ try:
+ if hasattr(callback_or_message, 'message'):
+ bot = callback_or_message.message.bot
+ else:
+ bot = callback_or_message.bot
+
+ bot_info = await bot.get_me()
+ bot_username = bot_info.username or 'your_bot'
+ except Exception as e:
+ logging.error(f"Failed to get bot username for CSV links: {e}")
+ bot_username = 'your_bot'
+
+ for code in created_codes:
+ # Determine validity info
+ if data.get("validity_days"):
+ valid_until = (datetime.now(timezone.utc) + timedelta(days=data["validity_days"])).strftime("%Y-%m-%d %H:%M:%S")
+ else:
+ valid_until = "Без ограничений"
+
+ start_command = f"/start promo_{code}"
+ telegram_link = f"https://t.me/{bot_username}?start=promo_{code}"
+
+ writer.writerow([
+ code,
+ data["bonus_days"],
+ data["max_activations"],
+ valid_until,
+ start_command,
+ telegram_link
+ ])
+
+ output.seek(0)
+
+ # Create file for sending
+ filename = f"bulk_promo_codes_{datetime.now().strftime('%Y%m%d_%H%M%S')}.csv"
+ csv_file = types.BufferedInputFile(
+ output.getvalue().encode('utf-8-sig'), # BOM for correct Excel display
+ filename=filename
+ )
if failed_codes:
success_lines.append(f"\n❌ Ошибки ({len(failed_codes)}):")
@@ -447,19 +497,26 @@ async def create_bulk_promo_codes_final(callback_or_message,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML"
)
+ message_obj = callback_or_message.message
except Exception:
- await callback_or_message.message.answer(
+ message_obj = await callback_or_message.message.answer(
success_text,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML"
)
+ await callback_or_message.answer()
else: # Message
- await callback_or_message.answer(
+ message_obj = await callback_or_message.answer(
success_text,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML"
)
+ # Send CSV file if created
+ if csv_file:
+ csv_caption = f"📄 Промокоды для массового создания\n💫 Всего: {len(created_codes)} промокодов\n🎁 Бонус: {data['bonus_days']} дней каждый"
+ await message_obj.answer_document(csv_file, caption=csv_caption)
+
await state.clear()
except Exception as e:
diff --git a/bot/handlers/admin/promo/create.py b/bot/handlers/admin/promo/create.py
index 7c7137f..6044913 100644
--- a/bot/handlers/admin/promo/create.py
+++ b/bot/handlers/admin/promo/create.py
@@ -345,18 +345,22 @@ async def create_promo_code_final(callback_or_message,
created_promo = await promo_code_dal.create_promo_code(session, promo_data)
await session.commit()
+ # Log successful creation
+ logging.info(f"Promo code '{data['promo_code']}' created with ID {created_promo.promo_code_id}")
+
# Success message
+ valid_until_str = _("admin_promo_unlimited", default="Без ограничений") if not data.get("validity_days") else f"{data['validity_days']} дней"
success_text = _(
"admin_promo_created_success",
default="✅ Промокод успешно создан!\n\n"
"🎟 Код: {code}\n"
"🎁 Бонусные дни: {bonus_days}\n"
"📊 Макс. активаций: {max_activations}\n"
- "⏰ Срок действия: {validity}",
+ "⏰ Срок действия: {valid_until_str}",
code=data["promo_code"],
bonus_days=data["bonus_days"],
max_activations=data["max_activations"],
- validity=_("admin_promo_unlimited", default="Без ограничений") if not data.get("validity_days") else f"{data['validity_days']} дней"
+ valid_until_str=valid_until_str
)
if hasattr(callback_or_message, 'message'): # CallbackQuery
diff --git a/bot/handlers/admin/promo/manage.py b/bot/handlers/admin/promo/manage.py
index 9d7c4f5..f8a5d97 100644
--- a/bot/handlers/admin/promo/manage.py
+++ b/bot/handlers/admin/promo/manage.py
@@ -8,7 +8,7 @@ from datetime import datetime, timedelta, timezone
from typing import Optional, List
from sqlalchemy.ext.asyncio import AsyncSession
-from config.settings import Settings
+from config.settings import Settings, get_settings
from db.dal import promo_code_dal
from db.models import PromoCode, PromoCodeActivation
from bot.states.admin_states import AdminStates
@@ -19,17 +19,27 @@ from bot.middlewares.i18n import JsonI18n
router = Router(name="promo_manage_router")
+def get_promo_status_emoji_and_text(promo: PromoCode, i18n: JsonI18n, current_lang: str):
+ """Determine promo code status and return emoji + text"""
+ _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
+
+ if promo.valid_until and promo.valid_until < datetime.now(timezone.utc):
+ return "⏰", _("admin_promo_status_expired")
+ elif promo.current_activations >= promo.max_activations:
+ return "🔄", _("admin_promo_status_used_up")
+ elif promo.is_active:
+ return "✅", _("admin_promo_status_active")
+ else:
+ return "🚫", _("admin_promo_status_inactive")
+
+
async def get_promo_detail_text_and_keyboard(promo_id: int, session: AsyncSession, i18n: JsonI18n, current_lang: str):
_ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
promo = await promo_code_dal.get_promo_code_by_id(session, promo_id)
if not promo:
return None, None
- status = _("admin_promo_status_active") if promo.is_active else _("admin_promo_status_inactive")
- if promo.valid_until and promo.valid_until < datetime.now(timezone.utc):
- status = _("admin_promo_status_expired")
- elif promo.current_activations >= promo.max_activations:
- status = _("admin_promo_status_used_up")
+ status_emoji, status = get_promo_status_emoji_and_text(promo, i18n, current_lang)
validity = _("admin_promo_valid_indefinitely")
if promo.valid_until:
@@ -68,7 +78,7 @@ async def view_promo_codes_handler(callback: types.CallbackQuery, i18n_data: dic
promo_models = await promo_code_dal.get_all_active_promo_codes(session, limit=20, offset=0)
text = f"{_('admin_active_promos_list_header')}\n\n{_('admin_no_active_promos')}" if not promo_models else "\n".join(
[_("admin_active_promos_list_header"), ""] + [
- f"🎟 {p.code} | 🎁 {p.bonus_days}д | 📊 {p.current_activations}/{p.max_activations} | ⏰ {p.valid_until.strftime('%d.%m.%Y') if p.valid_until else _('admin_promo_valid_indefinitely')}"
+ f"{get_promo_status_emoji_and_text(p, i18n, current_lang)[0]} {p.code} | 🎁 {p.bonus_days}д | 📊 {p.current_activations}/{p.max_activations} | ⏰ {p.valid_until.strftime('%d.%m.%Y') if p.valid_until else _('admin_promo_valid_indefinitely')}"
for p in promo_models
]
)
@@ -77,29 +87,66 @@ async def view_promo_codes_handler(callback: types.CallbackQuery, i18n_data: dic
await callback.answer()
-async def promo_management_handler(callback: types.CallbackQuery, i18n_data: dict, settings: Settings, session: AsyncSession):
- current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE)
+async def promo_management_handler(callback: types.CallbackQuery, i18n_data: dict, settings: Settings, session: AsyncSession, page: int = 0):
+ current_lang = i18n_data.get("current_language", "ru")
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
if not i18n or not callback.message:
await callback.answer("Error processing request.", show_alert=True)
return
_ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
- promo_models = await promo_code_dal.get_all_promo_codes_with_details(session, limit=50, offset=0)
- if not promo_models:
+ page_size = 10 # Количество промокодов на странице
+ offset = page * page_size
+
+ # Получаем общее количество промокодов
+ total_count = await promo_code_dal.get_promo_codes_count(session)
+ total_pages = (total_count + page_size - 1) // page_size if total_count > 0 else 1
+
+ promo_models = await promo_code_dal.get_all_promo_codes_with_details(session, limit=page_size, offset=offset)
+ if not promo_models and page == 0:
await callback.message.edit_text(_("admin_promo_management_empty"), reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n), parse_mode="HTML")
await callback.answer()
return
builder = InlineKeyboardBuilder()
for promo in promo_models:
- builder.row(InlineKeyboardButton(text=f"📝 {promo.code}", callback_data=f"promo_detail:{promo.promo_code_id}"))
+ status_emoji, status_text = get_promo_status_emoji_and_text(promo, i18n, current_lang)
+ button_text = f"{status_emoji} {promo.code} ({promo.current_activations}/{promo.max_activations})"
+ builder.row(InlineKeyboardButton(text=button_text, callback_data=f"promo_detail:{promo.promo_code_id}"))
+
+ # Добавляем кнопки пагинации если есть больше одной страницы
+ if total_pages > 1:
+ pagination_buttons = []
+ if page > 0:
+ pagination_buttons.append(InlineKeyboardButton(text=_("prev_page_button"), callback_data=f"promo_management:{page-1}"))
+ if page < total_pages - 1:
+ pagination_buttons.append(InlineKeyboardButton(text=_("next_page_button"), callback_data=f"promo_management:{page+1}"))
+
+ if pagination_buttons:
+ builder.row(*pagination_buttons)
+
+ # Добавляем кнопки экспорта и возврата
+ builder.row(InlineKeyboardButton(text="📄 Экспорт CSV", callback_data="promo_export_all"))
builder.row(InlineKeyboardButton(text=_("back_to_admin_panel_button"), callback_data="admin_action:main"))
- await callback.message.edit_text(_("admin_promo_management_title"), reply_markup=builder.as_markup(), parse_mode="HTML")
+ # Формируем заголовок с информацией о страницах
+ title = _("admin_promo_management_title")
+ if total_pages > 1:
+ title += f"\n{_('admin_promo_list_page_info', current=page+1, total=total_pages, count=total_count)}"
+
+ await callback.message.edit_text(title, reply_markup=builder.as_markup(), parse_mode="HTML")
await callback.answer()
+@router.callback_query(F.data.startswith("promo_management:"))
+async def promo_management_pagination_handler(callback: types.CallbackQuery, i18n_data: dict, settings: Settings, session: AsyncSession):
+ try:
+ page = int(callback.data.split(":")[1])
+ await promo_management_handler(callback, i18n_data, settings, session, page)
+ except (ValueError, IndexError):
+ await callback.answer("Error processing pagination.", show_alert=True)
+
+
@router.callback_query(F.data.startswith("promo_detail:"))
async def promo_detail_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession):
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
@@ -227,6 +274,63 @@ async def promo_export_activations_handler(callback: types.CallbackQuery, i18n_d
await callback.answer()
+@router.callback_query(F.data == "promo_export_all")
+async def promo_export_all_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession):
+ i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
+ current_lang = i18n_data.get("current_language")
+ if not i18n or not callback.message or not current_lang:
+ return await callback.answer("Error processing request.", show_alert=True)
+ _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
+
+ try:
+ await callback.answer("📄 Создаю CSV файл...", show_alert=True)
+
+ # Получаем все промокоды
+ all_promos = await promo_code_dal.get_all_promo_codes_with_details(session, limit=10000, offset=0)
+
+ output = io.StringIO()
+ writer = csv.writer(output)
+
+ # Заголовки CSV
+ writer.writerow([
+ "Код", "Бонусные дни", "Максимальные активации", "Текущие активации",
+ "Статус", "Активен", "Действителен до", "Создан", "Создал (Admin ID)"
+ ])
+
+ for promo in all_promos:
+ # Определяем статус
+ status_emoji, status_text = get_promo_status_emoji_and_text(promo, i18n, current_lang)
+
+ # Формируем данные для CSV
+ row = [
+ promo.code,
+ promo.bonus_days,
+ promo.max_activations,
+ promo.current_activations,
+ status_text,
+ "Да" if promo.is_active else "Нет",
+ promo.valid_until.strftime("%Y-%m-%d %H:%M:%S") if promo.valid_until else "Без ограничений",
+ promo.created_at.strftime("%Y-%m-%d %H:%M:%S") if promo.created_at else "N/A",
+ promo.created_by_admin_id or "N/A"
+ ]
+ writer.writerow(row)
+
+ output.seek(0)
+
+ # Создаем файл для отправки
+ filename = f"promo_codes_{datetime.now().strftime('%Y%m%d_%H%M%S')}.csv"
+ file = types.BufferedInputFile(
+ output.getvalue().encode('utf-8-sig'), # BOM для корректного отображения в Excel
+ filename=filename
+ )
+
+ caption = f"📄 Экспорт всех промокодов\n📊 Всего: {len(all_promos)} промокодов"
+ await callback.message.answer_document(file, caption=caption)
+
+ except Exception as e:
+ await callback.answer(f"❌ Ошибка экспорта: {str(e)}", show_alert=True)
+
+
@router.callback_query(F.data.startswith("promo_delete:"))
async def promo_delete_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession):
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
@@ -241,7 +345,7 @@ async def promo_delete_handler(callback: types.CallbackQuery, i18n_data: dict, s
if promo:
await session.commit()
await callback.answer(_("admin_promo_deleted_success", code=promo.code), show_alert=True)
- await promo_management_handler(callback, i18n_data, {}, session) # Settings not needed here
+ await promo_management_handler(callback, i18n_data, get_settings(), session, 0)
else:
await callback.answer(_("admin_promo_not_found"), show_alert=True)
except (ValueError, IndexError):
diff --git a/bot/handlers/admin/sync_admin.py b/bot/handlers/admin/sync_admin.py
index 4beca0f..603ecab 100644
--- a/bot/handlers/admin/sync_admin.py
+++ b/bot/handlers/admin/sync_admin.py
@@ -15,6 +15,222 @@ from bot.middlewares.i18n import JsonI18n
router = Router(name="admin_sync_router")
+async def perform_sync(panel_service: PanelApiService, session: AsyncSession,
+ settings: Settings, i18n_instance: JsonI18n) -> dict:
+ """
+ Perform panel synchronization and return results
+ Returns dict with status, details, and sync statistics
+ """
+ panel_records_checked = 0
+ users_found_in_db = 0
+ users_updated = 0
+ subscriptions_synced_count = 0
+ sync_errors = []
+
+ # Additional counters for detailed logging
+ users_without_telegram_id = 0
+ users_not_found_in_db = 0
+ users_uuid_updated = 0
+ subscriptions_created = 0
+ subscriptions_updated = 0
+
+ try:
+ panel_users_data = await panel_service.get_all_panel_users()
+
+ if panel_users_data is None:
+ error_msg = "Failed to fetch users from panel or panel API issue."
+ sync_errors.append(error_msg)
+ await panel_sync_dal.update_panel_sync_status(session, "failed", error_msg)
+ await session.commit()
+ return {"status": "failed", "details": error_msg, "errors": sync_errors}
+
+ if not panel_users_data:
+ status_msg = "No users found in the panel to sync."
+ await panel_sync_dal.update_panel_sync_status(
+ session, "success", status_msg, 0, 0
+ )
+ await session.commit()
+ return {"status": "success", "details": status_msg, "users_synced": 0, "subs_synced": 0}
+
+ total_panel_users = len(panel_users_data)
+ logging.info(f"Starting sync for {total_panel_users} panel users.")
+
+ for panel_user_dict in panel_users_data:
+ try:
+ panel_records_checked += 1
+ panel_uuid = panel_user_dict.get("uuid")
+ panel_subscription_uuid = panel_user_dict.get("subscriptionUuid") or panel_user_dict.get("shortUuid")
+ telegram_id_from_panel = panel_user_dict.get("telegramId")
+
+ if not panel_uuid:
+ sync_errors.append(f"Panel user missing UUID: {panel_user_dict}")
+ logging.warning(f"Skipping panel user without UUID: {panel_user_dict}")
+ continue
+
+ # Track users without telegram ID
+ if not telegram_id_from_panel:
+ users_without_telegram_id += 1
+
+ # Try to find existing user in local DB
+ existing_user = None
+
+ # First, try to find by telegram ID if available
+ if telegram_id_from_panel:
+ existing_user = await user_dal.get_user_by_id(session, telegram_id_from_panel)
+ if existing_user:
+ logging.debug(f"Found user by telegramId {telegram_id_from_panel}")
+
+ # If not found by telegram ID, try to find by panel UUID
+ if not existing_user:
+ existing_user = await user_dal.get_user_by_panel_uuid(session, panel_uuid)
+ if existing_user:
+ logging.info(f"Found user by panel UUID {panel_uuid}, telegramId: {existing_user.user_id}")
+ # Update telegram ID if it was missing in panel data but we have local user
+ if telegram_id_from_panel and existing_user.user_id != telegram_id_from_panel:
+ logging.warning(f"TelegramId mismatch: panel={telegram_id_from_panel}, local={existing_user.user_id}")
+
+ if not existing_user:
+ users_not_found_in_db += 1
+ if telegram_id_from_panel:
+ logging.debug(f"Panel user with telegramId {telegram_id_from_panel} and UUID {panel_uuid} not found in local DB")
+ else:
+ logging.debug(f"Panel user with UUID {panel_uuid} (no telegramId) not found in local DB")
+ continue
+
+ # User found in local DB
+ users_found_in_db += 1
+ user_was_updated = False
+
+ # Get the actual user_id for subscription operations
+ actual_user_id = existing_user.user_id
+
+ # Update panel UUID if different
+ if existing_user.panel_user_uuid != panel_uuid:
+ existing_user.panel_user_uuid = panel_uuid
+ user_was_updated = True
+ users_uuid_updated += 1
+ logging.info(f"Updated panel UUID for user {actual_user_id}: {panel_uuid}")
+
+ # Sync subscription data
+ panel_expire_at_iso = panel_user_dict.get("expireAt")
+ panel_status = panel_user_dict.get("status", "UNKNOWN")
+
+ if panel_expire_at_iso:
+ try:
+ panel_expire_at = datetime.fromisoformat(
+ panel_expire_at_iso.replace("Z", "+00:00")
+ )
+
+ # Update or create subscription
+ active_sub = await subscription_dal.get_active_subscription_by_user_id(
+ session, actual_user_id, panel_uuid
+ )
+
+ if active_sub:
+ # Check if subscription needs update
+ if (active_sub.end_date != panel_expire_at or
+ active_sub.status_from_panel != panel_status or
+ active_sub.is_active != (panel_status == "ACTIVE")):
+
+ await subscription_dal.update_subscription_end_date(
+ session, active_sub.subscription_id, panel_expire_at
+ )
+ # Update status fields
+ active_sub.status_from_panel = panel_status
+ active_sub.is_active = (panel_status == "ACTIVE")
+ subscriptions_synced_count += 1
+ subscriptions_updated += 1
+ user_was_updated = True
+ logging.info(f"Updated subscription for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}")
+ else:
+ # Create new subscription record
+ subscription_uuid_to_use = panel_subscription_uuid or panel_uuid
+
+ logging.info(f"Creating new subscription for user {actual_user_id} with UUID {subscription_uuid_to_use}")
+
+ sub_payload = {
+ "user_id": actual_user_id,
+ "panel_user_uuid": panel_uuid,
+ "panel_subscription_uuid": subscription_uuid_to_use,
+ "start_date": datetime.now(timezone.utc),
+ "end_date": panel_expire_at,
+ "duration_months": 1, # Default
+ "is_active": panel_status == "ACTIVE",
+ "status_from_panel": panel_status,
+ "traffic_limit_bytes": settings.user_traffic_limit_bytes,
+ }
+ await subscription_dal.upsert_subscription(session, sub_payload)
+ subscriptions_synced_count += 1
+ subscriptions_created += 1
+ user_was_updated = True
+
+ except Exception as e:
+ sync_errors.append(f"Error syncing subscription for user {actual_user_id}: {str(e)}")
+ logging.error(f"Error syncing subscription for user {actual_user_id}: {e}")
+
+ if user_was_updated:
+ users_updated += 1
+
+ except Exception as e_user:
+ sync_errors.append(f"Error processing panel user {panel_user_dict.get('uuid', 'unknown')}: {str(e_user)}")
+ logging.error(f"Error syncing user: {e_user}")
+
+ # Update sync status
+ status = "completed_with_errors" if sync_errors else "completed"
+ details = (f"📊 Статистика синхронизации:\n"
+ f"🔍 Проверено записей панели: {panel_records_checked}\n"
+ f"👥 Найдено пользователей в БД: {users_found_in_db}\n"
+ f"🔄 Пользователей обновлено: {users_updated}\n"
+ f"📋 Подписок синхронизировано: {subscriptions_synced_count}\n"
+ f" ├── Создано новых: {subscriptions_created}\n"
+ f" └── Обновлено существующих: {subscriptions_updated}")
+
+ if users_without_telegram_id > 0:
+ details += f"\n⚠️ Записей без telegramId: {users_without_telegram_id}"
+ if users_not_found_in_db > 0:
+ details += f"\n❌ Не найдено в БД: {users_not_found_in_db}"
+ if sync_errors:
+ details += f"\n🚫 Ошибок: {len(sync_errors)}"
+
+ await panel_sync_dal.update_panel_sync_status(
+ session, status, details, panel_records_checked, subscriptions_synced_count
+ )
+ await session.commit()
+
+ # Detailed logging summary
+ logging.info(f"Sync completed - Summary:")
+ logging.info(f" Panel records checked: {panel_records_checked}")
+ logging.info(f" Users without telegramId: {users_without_telegram_id}")
+ logging.info(f" Users not found in local DB: {users_not_found_in_db}")
+ logging.info(f" Users found in local DB: {users_found_in_db}")
+ logging.info(f" Users with UUID updated: {users_uuid_updated}")
+ logging.info(f" Users updated overall: {users_updated}")
+ logging.info(f" Subscriptions total synced: {subscriptions_synced_count}")
+ logging.info(f" Subscriptions created: {subscriptions_created}")
+ logging.info(f" Subscriptions updated: {subscriptions_updated}")
+ logging.info(f" Sync errors: {len(sync_errors)}")
+
+ return {
+ "status": status,
+ "details": details,
+ "users_processed": panel_records_checked,
+ "users_synced": users_found_in_db,
+ "subs_synced": subscriptions_synced_count,
+ "errors": sync_errors
+ }
+
+ except Exception as e_sync_global:
+ await session.rollback()
+ logging.error(f"Global error during sync: {e_sync_global}", exc_info=True)
+ error_detail = f"Unexpected error during sync: {str(e_sync_global)[:200]}"
+
+ await panel_sync_dal.update_panel_sync_status(
+ session, "failed", error_detail, panel_records_checked, subscriptions_synced_count
+ )
+
+ return {"status": "failed", "details": error_detail, "errors": [str(e_sync_global)]}
+
+
@router.message(Command("sync"))
async def sync_command_handler(
message_event: Union[types.Message, types.CallbackQuery],
@@ -52,265 +268,39 @@ async def sync_command_handler(
logging.info(f"Admin ({message_event.from_user.id}) triggered panel sync.")
- users_processed_count = 0
- users_synced_successfully = 0
- subscriptions_synced_count = 0
- sync_errors = []
-
+ # Use the extracted perform_sync function
try:
- panel_users_data = await panel_service.get_all_panel_users()
-
- if panel_users_data is None:
- error_msg = "Failed to fetch users from panel or panel API issue."
- sync_errors.append(error_msg)
- await panel_sync_dal.update_panel_sync_status(session, "failed", error_msg)
- await session.commit()
- await bot.send_message(target_chat_id, _("sync_failed", details=error_msg))
- return
-
- if not panel_users_data:
- status_msg = "No users found in the panel to sync."
- await panel_sync_dal.update_panel_sync_status(
- session, "success", status_msg, 0, 0
+ sync_result = await perform_sync(panel_service, session, settings, i18n)
+
+ status = sync_result.get("status")
+ details = sync_result.get("details", "No details available")
+ errors = sync_result.get("errors", [])
+
+ if status == "failed":
+ await bot.send_message(target_chat_id, _("sync_failed", details=details))
+ elif status == "completed_with_errors":
+ error_preview = "; ".join(errors[:3]) # Show first 3 errors
+ final_message = _(
+ "sync_completed_with_errors_details",
+ total_checked=sync_result.get("users_processed", 0),
+ users_synced=sync_result.get("users_synced", 0),
+ subs_synced=sync_result.get("subs_synced", 0),
+ errors_count=len(errors),
+ error_details_preview=error_preview[:200] + "..." if len(error_preview) > 200 else error_preview
)
- await session.commit()
- await bot.send_message(
- target_chat_id,
- _("sync_completed", status="Success", details=status_msg),
- )
- return
-
- total_panel_users = len(panel_users_data)
- logging.info(f"Starting sync for {total_panel_users} panel users.")
-
- for panel_user_dict in panel_users_data:
- users_processed_count += 1
- panel_uuid = panel_user_dict.get("uuid")
- telegram_id_from_panel_str = panel_user_dict.get("telegramId")
- panel_username = panel_user_dict.get("username")
-
- if not panel_uuid:
- logging.warning(
- f"Sync: Panel user data missing 'uuid'. Data: {str(panel_user_dict)[:200]}. Skipping."
- )
- sync_errors.append(
- f"Panel user data (username: {panel_username or 'N/A'}) missing UUID."
- )
- continue
-
- telegram_id_from_panel: Optional[int] = None
- if telegram_id_from_panel_str:
- try:
- telegram_id_from_panel = int(telegram_id_from_panel_str)
- except ValueError:
- logging.warning(
- f"Sync: Panel user {panel_uuid} (username: {panel_username}) has invalid 'telegramId': {telegram_id_from_panel_str}. Skipping TG ID based sync."
- )
-
- if not telegram_id_from_panel:
-
- logging.info(
- f"Sync: Panel user {panel_uuid} (username: {panel_username}) has no valid 'telegramId'. Skipping full sync for this user."
- )
-
- continue
-
- bot_user = await user_dal.get_user_by_id(session, telegram_id_from_panel)
- if not bot_user:
- user_data_to_create = {
- "user_id": telegram_id_from_panel,
- "username": panel_username,
- "panel_user_uuid": panel_uuid,
- "language_code": settings.DEFAULT_LANGUAGE,
- "registration_date": (
- datetime.fromisoformat(
- panel_user_dict["createdAt"].replace("Z", "+00:00")
- )
- if panel_user_dict.get("createdAt")
- else datetime.now(timezone.utc)
- ),
- }
- bot_user = await user_dal.create_user(session, user_data_to_create)
- logging.info(
- f"Sync: Created new local user {telegram_id_from_panel} from panel data {panel_uuid}."
- )
- else:
- if bot_user.panel_user_uuid != panel_uuid:
- if bot_user.panel_user_uuid is not None:
- logging.warning(
- f"Sync: Local user {telegram_id_from_panel} was linked to {bot_user.panel_user_uuid}, panel now gives {panel_uuid}. Updating."
- )
-
- conflicting_user = await user_dal.get_user_by_panel_uuid(
- session, panel_uuid
- )
- if (
- conflicting_user
- and conflicting_user.user_id != telegram_id_from_panel
- ):
- sync_errors.append(
- f"Panel UUID {panel_uuid} for TG {telegram_id_from_panel} already linked to another TG user {conflicting_user.user_id}."
- )
- logging.error(sync_errors[-1])
- continue
-
- await user_dal.update_user(
- session,
- telegram_id_from_panel,
- {"panel_user_uuid": panel_uuid, "username": panel_username},
- )
- logging.info(
- f"Sync: Updated panel_uuid for local user {telegram_id_from_panel} to {panel_uuid}."
- )
-
- panel_sub_link_id = panel_user_dict.get(
- "subscriptionUuid"
- ) or panel_user_dict.get("shortUuid")
- if panel_sub_link_id:
- end_date_str = panel_user_dict.get("expireAt")
- start_date_str = panel_user_dict.get("createdAt")
-
- if end_date_str:
- try:
- end_date_obj = datetime.fromisoformat(
- end_date_str.replace("Z", "+00:00")
- )
- start_date_obj = (
- datetime.fromisoformat(
- start_date_str.replace("Z", "+00:00")
- )
- if start_date_str
- else datetime.now(timezone.utc)
- )
-
- status_from_panel = panel_user_dict.get(
- "status", "UNKNOWN"
- ).upper()
- is_active_flag = (
- 1
- if status_from_panel == "ACTIVE"
- and end_date_obj > datetime.now(timezone.utc)
- else 0
- )
-
- sub_payload = {
- "user_id": telegram_id_from_panel,
- "panel_user_uuid": panel_uuid,
- "panel_subscription_uuid": panel_sub_link_id,
- "start_date": start_date_obj,
- "end_date": end_date_obj,
- "is_active": is_active_flag,
- "status_from_panel": status_from_panel,
- "traffic_limit_bytes": panel_user_dict.get(
- "trafficLimitBytes"
- ),
- "traffic_used_bytes": panel_user_dict.get(
- "usedTrafficBytes"
- ),
- }
-
- await subscription_dal.deactivate_other_active_subscriptions(
- session, panel_uuid, panel_sub_link_id
- )
- await subscription_dal.upsert_subscription(session, sub_payload)
- subscriptions_synced_count += 1
- users_synced_successfully += 1
- except ValueError as e_date:
- logging.warning(
- f"Sync: Bad date format for panel user {panel_uuid} (TG ID: {telegram_id_from_panel}). Sub data: {str(panel_user_dict)[:100]}. Error: {e_date}"
- )
- sync_errors.append(
- f"Bad date for panel user {panel_uuid} (TG ID: {telegram_id_from_panel})."
- )
- except Exception as e_sub_sync:
- logging.error(
- f"Sync: Error syncing subscription for panel user {panel_uuid} (TG ID: {telegram_id_from_panel}): {e_sub_sync}",
- exc_info=True,
- )
- sync_errors.append(
- f"Sub sync error for panel user {panel_uuid} (TG ID: {telegram_id_from_panel})."
- )
- else:
- logging.warning(
- f"Sync: Panel user {panel_uuid} (TG ID: {telegram_id_from_panel}) has sub link but no expireAt date. Skipping subscription sync."
- )
- else:
-
- await subscription_dal.deactivate_other_active_subscriptions(
- session, panel_uuid, None
- )
- logging.info(
- f"Sync: Panel user {panel_uuid} (TG ID: {telegram_id_from_panel}) has no subscription link on panel. Deactivated local subs if any."
- )
- users_synced_successfully += 1
-
- if users_processed_count % 20 == 0:
- logging.info(
- f"Sync progress: {users_processed_count}/{total_panel_users} users processed from panel."
- )
-
- panel_uuid_set = {u.get("uuid") for u in panel_users_data if u.get("uuid")}
- local_users_with_uuid = await user_dal.get_all_users_with_panel_uuid(session)
- for local_user in local_users_with_uuid:
- if local_user.panel_user_uuid not in panel_uuid_set:
- await subscription_dal.deactivate_other_active_subscriptions(
- session, local_user.panel_user_uuid, None
- )
- logging.info(
- f"Sync: Local user {local_user.user_id} with panel UUID {local_user.panel_user_uuid} not found on panel. Deactivated local subs."
- )
-
- status_msg_key = "sync_completed_details"
- final_status_type = "success"
-
- if sync_errors:
- final_status_type = "partial_success"
- status_msg_key = "sync_completed_with_errors_details"
- error_preview = "\n".join(sync_errors[:3])
- details_for_db = f"Users processed: {users_processed_count}. Subs synced: {subscriptions_synced_count}. Errors: {len(sync_errors)}. First few: {error_preview}"
+ await bot.send_message(target_chat_id, final_message)
else:
- details_for_db = f"Successfully processed {users_processed_count} users. Synced {subscriptions_synced_count} subscriptions."
-
- await panel_sync_dal.update_panel_sync_status(
- session,
- final_status_type,
- details_for_db,
- users_processed_count,
- subscriptions_synced_count,
- )
- await session.commit()
-
- final_user_message = _(
- status_msg_key,
- total_checked=total_panel_users,
- users_synced=users_synced_successfully,
- subs_synced=subscriptions_synced_count,
- errors_count=len(sync_errors),
- error_details_preview=(
- error_preview if sync_errors else _("no_errors_placeholder")
- ),
- )
- await bot.send_message(target_chat_id, final_user_message)
-
+ final_message = _(
+ "sync_completed_details",
+ total_checked=sync_result.get("users_processed", 0),
+ users_synced=sync_result.get("users_synced", 0),
+ subs_synced=sync_result.get("subs_synced", 0)
+ )
+ await bot.send_message(target_chat_id, _("sync_completed", status="Success", details=final_message))
+
except Exception as e_sync_global:
- await session.rollback()
- logging.error(
- f"Global error during /sync command: {e_sync_global}", exc_info=True
- )
- error_detail_for_db = (
- f"An unexpected error occurred during sync: {str(e_sync_global)[:200]}"
- )
- await panel_sync_dal.update_panel_sync_status(
- session,
- "failed",
- error_detail_for_db,
- users_processed_count,
- subscriptions_synced_count,
- )
-
- await bot.send_message(
- target_chat_id, _("sync_failed", details=error_detail_for_db)
- )
+ logging.error(f"Global error during /sync command: {e_sync_global}", exc_info=True)
+ await bot.send_message(target_chat_id, _("sync_failed", details=str(e_sync_global)))
@router.message(Command("syncstatus"))
@@ -350,4 +340,4 @@ async def sync_status_command_handler(
else:
response_text = _("admin_sync_status_never_run")
- await message.answer(response_text, parse_mode="HTML")
+ await message.answer(response_text, parse_mode="HTML")
\ No newline at end of file
diff --git a/bot/handlers/user/start.py b/bot/handlers/user/start.py
index ab9c19e..d0fa7f8 100644
--- a/bot/handlers/user/start.py
+++ b/bot/handlers/user/start.py
@@ -216,13 +216,15 @@ async def start_command_handler(message: types.Message,
f"Failed to update existing user {user_id} in session: {e_update}",
exc_info=True)
- await message.answer(_(key="welcome", user_name=hd.quote(user.full_name)))
+ # Send welcome message if not disabled
+ if not settings.DISABLE_WELCOME_MESSAGE:
+ await message.answer(_(key="welcome", user_name=hd.quote(user.full_name)))
# Auto-apply promo code if provided via start parameter
if promo_code_to_apply:
try:
from bot.services.promo_code_service import PromoCodeService
- promo_code_service = PromoCodeService()
+ promo_code_service = PromoCodeService(settings, subscription_service, message.bot, i18n)
success, result = await promo_code_service.apply_promo_code(
session, user_id, promo_code_to_apply, current_lang
diff --git a/bot/keyboards/inline/admin_keyboards.py b/bot/keyboards/inline/admin_keyboards.py
index 99eb072..bf0d4e7 100644
--- a/bot/keyboards/inline/admin_keyboards.py
+++ b/bot/keyboards/inline/admin_keyboards.py
@@ -105,10 +105,12 @@ def get_system_functions_keyboard(i18n_instance, lang: str) -> InlineKeyboardMar
callback_data="admin_action:broadcast")
builder.button(text=_(key="admin_sync_panel_button"),
callback_data="admin_action:sync_panel")
+ builder.button(text=_(key="admin_queue_status_button"),
+ callback_data="admin_action:queue_status")
builder.button(text=_(key="back_to_admin_panel_button"),
callback_data="admin_action:main")
- builder.adjust(2, 1)
+ builder.adjust(2, 1, 1)
return builder.as_markup()
diff --git a/bot/main_bot.py b/bot/main_bot.py
index 1408c22..21d2964 100644
--- a/bot/main_bot.py
+++ b/bot/main_bot.py
@@ -42,6 +42,8 @@ from bot.services.tribute_service import TributeService, tribute_webhook_route
from bot.services.crypto_pay_service import CryptoPayService, cryptopay_webhook_route
from bot.handlers.user import payment as user_payment_webhook_module
+from bot.handlers.admin.sync_admin import perform_sync
+from bot.utils.message_queue import init_queue_manager
class DBSessionMiddleware(BaseMiddleware):
@@ -191,6 +193,34 @@ async def on_startup_configured(dispatcher: Dispatcher):
except Exception as e:
logging.error(f"STARTUP: Failed to set bot commands: {e}", exc_info=True)
+ # Initialize message queue manager
+ try:
+ queue_manager = init_queue_manager(bot)
+ dispatcher["queue_manager"] = queue_manager
+ logging.info("STARTUP: Message queue manager initialized")
+ except Exception as e:
+ logging.error(f"STARTUP: Failed to initialize message queue manager: {e}", exc_info=True)
+
+ # Automatic sync on startup
+ try:
+ logging.info("STARTUP: Running automatic panel sync...")
+
+ async with async_session_factory() as session:
+ sync_result = await perform_sync(
+ panel_service=panel_service,
+ session=session,
+ settings=settings,
+ i18n_instance=i18n_instance
+ )
+
+ if sync_result.get("status") == "completed":
+ logging.info(f"STARTUP: Automatic sync completed successfully. Details: {sync_result.get('details', 'N/A')}")
+ else:
+ logging.warning(f"STARTUP: Automatic sync completed with issues. Status: {sync_result.get('status', 'unknown')}")
+
+ except Exception as e:
+ logging.error(f"STARTUP: Failed to run automatic sync: {e}", exc_info=True)
+
logging.info("STARTUP: Bot on_startup_configured completed.")
diff --git a/bot/services/notification_service.py b/bot/services/notification_service.py
index 15af8e3..8597597 100644
--- a/bot/services/notification_service.py
+++ b/bot/services/notification_service.py
@@ -2,12 +2,14 @@ import logging
import asyncio
from aiogram import Bot
from aiogram.utils.text_decorations import html_decoration as hd
+from aiogram.exceptions import TelegramRetryAfter
from datetime import datetime, timezone
from typing import Optional, Union, Dict, Any
from config.settings import Settings
from sqlalchemy.orm import sessionmaker
from bot.middlewares.i18n import JsonI18n
+from bot.utils.message_queue import get_queue_manager
class NotificationService:
@@ -19,16 +21,30 @@ class NotificationService:
self.i18n = i18n
async def _send_to_log_channel(self, message: str, thread_id: Optional[int] = None):
- """Send message to configured log channel/group"""
+ """Send message to configured log channel/group using message queue"""
if not self.settings.LOG_CHAT_ID:
return
+ queue_manager = get_queue_manager()
+ if not queue_manager:
+ logging.warning("Message queue manager not available, falling back to direct send")
+ try:
+ await self.bot.send_message(
+ chat_id=self.settings.LOG_CHAT_ID,
+ text=message,
+ parse_mode="HTML",
+ disable_web_page_preview=True,
+ message_thread_id=thread_id or self.settings.LOG_THREAD_ID
+ )
+ except Exception as e:
+ logging.error(f"Failed to send notification to log channel {self.settings.LOG_CHAT_ID}: {e}")
+ return
+
try:
# Use thread_id if provided, otherwise use from settings
final_thread_id = thread_id or self.settings.LOG_THREAD_ID
kwargs = {
- "chat_id": self.settings.LOG_CHAT_ID,
"text": message,
"parse_mode": "HTML",
"disable_web_page_preview": True
@@ -38,26 +54,42 @@ class NotificationService:
if final_thread_id:
kwargs["message_thread_id"] = final_thread_id
- await self.bot.send_message(**kwargs)
+ # Queue message for sending (groups are rate limited to 15/minute)
+ await queue_manager.send_message(self.settings.LOG_CHAT_ID, **kwargs)
except Exception as e:
- logging.error(f"Failed to send notification to log channel {self.settings.LOG_CHAT_ID}: {e}")
+ logging.error(f"Failed to queue notification to log channel {self.settings.LOG_CHAT_ID}: {e}")
async def _send_to_admins(self, message: str):
- """Send message to all admin users"""
+ """Send message to all admin users using message queue"""
if not self.settings.ADMIN_IDS:
return
+ queue_manager = get_queue_manager()
+ if not queue_manager:
+ logging.warning("Message queue manager not available, falling back to direct send")
+ for admin_id in self.settings.ADMIN_IDS:
+ try:
+ await self.bot.send_message(
+ chat_id=admin_id,
+ text=message,
+ parse_mode="HTML",
+ disable_web_page_preview=True
+ )
+ except Exception as e:
+ logging.error(f"Failed to send notification to admin {admin_id}: {e}")
+ return
+
for admin_id in self.settings.ADMIN_IDS:
try:
- await self.bot.send_message(
+ await queue_manager.send_message(
chat_id=admin_id,
text=message,
parse_mode="HTML",
disable_web_page_preview=True
)
except Exception as e:
- logging.error(f"Failed to send notification to admin {admin_id}: {e}")
+ logging.error(f"Failed to queue notification to admin {admin_id}: {e}")
async def notify_new_user_registration(self, user_id: int, username: Optional[str] = None,
first_name: Optional[str] = None,
diff --git a/bot/services/promo_code_service.py b/bot/services/promo_code_service.py
index 6aaaca8..77ec565 100644
--- a/bot/services/promo_code_service.py
+++ b/bot/services/promo_code_service.py
@@ -55,7 +55,6 @@ class PromoCodeService:
reason=f"promo code {code_input_upper}")
if new_end_date:
-
activation_recorded = await promo_code_dal.record_promo_activation(
session, promo_data.promo_code_id, user_id, payment_id=None)
promo_incremented = await promo_code_dal.increment_promo_code_usage(
@@ -83,5 +82,4 @@ class PromoCodeService:
)
return False, _("error_applying_promo_bonus")
else:
-
return False, _("error_applying_promo_bonus")
diff --git a/bot/services/subscription_service.py b/bot/services/subscription_service.py
index 0be251b..61795c3 100644
--- a/bot/services/subscription_service.py
+++ b/bot/services/subscription_service.py
@@ -563,6 +563,10 @@ class SubscriptionService:
)
start_date = datetime.now(timezone.utc)
new_end_date_obj = start_date + timedelta(days=bonus_days)
+
+ # For promo code activations, use the configured user traffic limit
+ traffic_limit = self.settings.user_traffic_limit_bytes if "promo code" in reason.lower() else self.settings.trial_traffic_limit_bytes
+
bonus_sub_payload = {
"user_id": user_id,
"panel_user_uuid": panel_uuid,
@@ -572,7 +576,7 @@ class SubscriptionService:
"duration_months": 0,
"is_active": True,
"status_from_panel": "ACTIVE_BONUS",
- "traffic_limit_bytes": self.settings.user_traffic_limit_bytes,
+ "traffic_limit_bytes": traffic_limit,
}
await subscription_dal.deactivate_other_active_subscriptions(
session, panel_uuid, panel_sub_uuid
@@ -593,14 +597,23 @@ class SubscriptionService:
)
if updated_sub_model:
+ # Prepare panel update payload
+ panel_update_payload = {
+ "expireAt": new_end_date_obj.isoformat(
+ timespec="milliseconds"
+ ).replace("+00:00", "Z")
+ }
+
+ # For promo code activations, remove traffic limit
+ if "promo code" in reason.lower():
+ panel_update_payload["trafficLimitBytes"] = self.settings.user_traffic_limit_bytes
+ panel_update_payload["trafficLimitStrategy"] = self.settings.USER_TRAFFIC_STRATEGY
+ logging.info(f"Updating traffic limit for user {user_id} to {self.settings.user_traffic_limit_bytes} bytes due to promo code activation")
+
panel_update_success = (
await self.panel_service.update_user_details_on_panel(
panel_uuid,
- {
- "expireAt": new_end_date_obj.isoformat(
- timespec="milliseconds"
- ).replace("+00:00", "Z")
- },
+ panel_update_payload,
)
)
if not panel_update_success:
diff --git a/bot/utils/__init__.py b/bot/utils/__init__.py
new file mode 100644
index 0000000..512ba75
--- /dev/null
+++ b/bot/utils/__init__.py
@@ -0,0 +1 @@
+# Bot utilities package
\ No newline at end of file
diff --git a/bot/utils/message_queue.py b/bot/utils/message_queue.py
new file mode 100644
index 0000000..a0704b8
--- /dev/null
+++ b/bot/utils/message_queue.py
@@ -0,0 +1,183 @@
+import asyncio
+import logging
+from typing import Dict, Any, Callable, Awaitable, Optional
+from dataclasses import dataclass
+from datetime import datetime, timedelta
+from collections import deque
+from aiogram import Bot
+
+
+@dataclass
+class QueuedMessage:
+ """Represents a queued message with all necessary parameters"""
+ chat_id: int
+ method_name: str # 'send_message', 'edit_message_text', etc.
+ kwargs: Dict[str, Any]
+ callback: Optional[Callable[[Any], Awaitable[None]]] = None # Optional callback for result
+
+
+class MessageQueue:
+ """Message queue with rate limiting for Telegram API"""
+
+ def __init__(self, messages_per_second: float, burst_size: int = 5):
+ self.messages_per_second = messages_per_second
+ self.burst_size = burst_size
+ self.queue: deque[QueuedMessage] = deque()
+ self.last_send_times: deque[datetime] = deque()
+ self.is_processing = False
+ self.delay_between_messages = 1.0 / messages_per_second
+
+ async def add_message(self, message: QueuedMessage) -> None:
+ """Add message to queue"""
+ self.queue.append(message)
+ if not self.is_processing:
+ asyncio.create_task(self._process_queue())
+
+ async def _process_queue(self) -> None:
+ """Process messages from queue with rate limiting"""
+ if self.is_processing:
+ return
+
+ self.is_processing = True
+
+ try:
+ while self.queue:
+ # Check if we need to wait
+ await self._wait_if_needed()
+
+ # Get and process next message
+ message = self.queue.popleft()
+ try:
+ await self._send_message(message)
+ self.last_send_times.append(datetime.now())
+
+ # Keep only recent send times (last minute)
+ cutoff_time = datetime.now() - timedelta(seconds=60)
+ while self.last_send_times and self.last_send_times[0] < cutoff_time:
+ self.last_send_times.popleft()
+
+ except Exception as e:
+ logging.error(f"Failed to send queued message to {message.chat_id}: {e}")
+
+ finally:
+ self.is_processing = False
+
+ async def _wait_if_needed(self) -> None:
+ """Wait if we need to respect rate limits"""
+ if not self.last_send_times:
+ return
+
+ # Calculate time since last message
+ time_since_last = (datetime.now() - self.last_send_times[-1]).total_seconds()
+
+ if time_since_last < self.delay_between_messages:
+ wait_time = self.delay_between_messages - time_since_last
+ await asyncio.sleep(wait_time)
+
+ async def _send_message(self, message: QueuedMessage) -> Any:
+ """Send a single message - to be implemented by subclass"""
+ raise NotImplementedError("Subclass must implement _send_message")
+
+
+class TelegramMessageQueue(MessageQueue):
+ """Telegram-specific message queue"""
+
+ def __init__(self, bot: Bot, messages_per_second: float, burst_size: int = 5):
+ super().__init__(messages_per_second, burst_size)
+ self.bot = bot
+
+ async def _send_message(self, message: QueuedMessage) -> Any:
+ """Send message using bot method"""
+ method = getattr(self.bot, message.method_name)
+ result = await method(chat_id=message.chat_id, **message.kwargs)
+
+ # Call callback if provided
+ if message.callback:
+ await message.callback(result)
+
+ return result
+
+
+class MessageQueueManager:
+ """Manager for different types of message queues"""
+
+ def __init__(self, bot: Bot):
+ self.bot = bot
+
+ # Different queues for different types of chats
+ self.group_queue = TelegramMessageQueue(
+ bot=bot,
+ messages_per_second=15/60, # 15 messages per minute for groups
+ burst_size=3
+ )
+
+ self.user_queue = TelegramMessageQueue(
+ bot=bot,
+ messages_per_second=25, # 25 messages per second for users
+ burst_size=10
+ )
+
+ def _is_group_chat(self, chat_id: int) -> bool:
+ """Check if chat_id belongs to a group or channel"""
+ return str(chat_id).startswith('-100')
+
+ async def send_message(self, chat_id: int, **kwargs) -> None:
+ """Queue a send_message call"""
+ queue = self.group_queue if self._is_group_chat(chat_id) else self.user_queue
+ message = QueuedMessage(
+ chat_id=chat_id,
+ method_name='send_message',
+ kwargs=kwargs
+ )
+ await queue.add_message(message)
+
+ async def edit_message_text(self, chat_id: int, **kwargs) -> None:
+ """Queue an edit_message_text call"""
+ queue = self.group_queue if self._is_group_chat(chat_id) else self.user_queue
+ message = QueuedMessage(
+ chat_id=chat_id,
+ method_name='edit_message_text',
+ kwargs=kwargs
+ )
+ await queue.add_message(message)
+
+ async def send_document(self, chat_id: int, **kwargs) -> None:
+ """Queue a send_document call"""
+ queue = self.group_queue if self._is_group_chat(chat_id) else self.user_queue
+ message = QueuedMessage(
+ chat_id=chat_id,
+ method_name='send_document',
+ kwargs=kwargs
+ )
+ await queue.add_message(message)
+
+ async def answer_callback_query(self, callback_query_id: str, **kwargs) -> None:
+ """Send callback query answer immediately (not rate limited)"""
+ await self.bot.answer_callback_query(callback_query_id, **kwargs)
+
+ def get_queue_stats(self) -> Dict[str, Any]:
+ """Get statistics about queues"""
+ return {
+ "group_queue_size": len(self.group_queue.queue),
+ "user_queue_size": len(self.user_queue.queue),
+ "group_queue_processing": self.group_queue.is_processing,
+ "user_queue_processing": self.user_queue.is_processing,
+ "group_recent_sends": len(self.group_queue.last_send_times),
+ "user_recent_sends": len(self.user_queue.last_send_times)
+ }
+
+
+# Global queue manager instance
+_queue_manager: Optional[MessageQueueManager] = None
+
+
+def init_queue_manager(bot: Bot) -> MessageQueueManager:
+ """Initialize global queue manager"""
+ global _queue_manager
+ _queue_manager = MessageQueueManager(bot)
+ return _queue_manager
+
+
+def get_queue_manager() -> Optional[MessageQueueManager]:
+ """Get global queue manager instance"""
+ return _queue_manager
\ No newline at end of file
diff --git a/config/settings.py b/config/settings.py
index aec0979..82e08ec 100644
--- a/config/settings.py
+++ b/config/settings.py
@@ -113,6 +113,7 @@ class Settings(BaseSettings):
SUBSCRIPTION_MINI_APP_URL: Optional[str] = Field(default=None)
START_COMMAND_DESCRIPTION: Optional[str] = Field(default=None)
+ DISABLE_WELCOME_MESSAGE: bool = Field(default=False, description="Disable welcome message on /start command")
# Inline mode thumbnail URLs
INLINE_REFERRAL_THUMBNAIL_URL: str = Field(default="https://cdn-icons-png.flaticon.com/512/1077/1077114.png")
diff --git a/db/dal/promo_code_dal.py b/db/dal/promo_code_dal.py
index 8d2f964..00adb21 100644
--- a/db/dal/promo_code_dal.py
+++ b/db/dal/promo_code_dal.py
@@ -64,6 +64,14 @@ async def get_all_promo_codes_with_details(session: AsyncSession, limit: int = 5
return result.scalars().all()
+async def get_promo_codes_count(session: AsyncSession) -> int:
+ """Get total count of all promo codes"""
+ from sqlalchemy import func
+ stmt = select(func.count(PromoCode.promo_code_id))
+ result = await session.execute(stmt)
+ return result.scalar_one()
+
+
async def get_promo_activations_by_code_id(session: AsyncSession, promo_code_id: int, limit: Optional[int] = None, offset: int = 0) -> List[PromoCodeActivation]:
"""Get activation history for a specific promo code with optional pagination."""
stmt = (select(PromoCodeActivation)
diff --git a/locales/ru.json b/locales/ru.json
index 9c75ebc..fd7177d 100644
--- a/locales/ru.json
+++ b/locales/ru.json
@@ -150,6 +150,16 @@
"admin_promo_invalid_values": "Неверные значения. {error}",
"admin_promo_invalid_format_general": "Ошибка парсинга деталей промокода. Проверьте формат.",
"admin_promo_created_success": "✅ Промокод {code} успешно создан!\nБонус: {bonus_days} дней\nМакс. активаций: {max_activations}\nДействителен: {valid_until_str}",
+ "admin_promo_set_validity_days": "⏰ Установить срок (дни)",
+ "admin_back_to_panel": "⬅️ В панель",
+ "admin_promo_unlimited": "♾️ Неограниченно",
+ "admin_bulk_promo_created_title": "📦 Массовое создание завершено",
+ "admin_bulk_promo_created_stats": "📊 Создано: {created} из {total}",
+ "admin_bulk_promo_settings": "🎁 Бонусные дни: {bonus_days}\n📊 Макс. активаций: {max_activations}\n⏰ Срок действия: {validity}",
+ "admin_promo_list_page_info": "Страница {current}/{total} ({count} промокодов)",
+ "admin_queue_status_button": "📊 Статус очередей",
+ "admin_queue_status_title": "📊 Статус очередей сообщений",
+ "admin_queue_status_info": "📤 Очереди сообщений:\n\n👥 Пользователи (25 сообщ/сек):\n 📋 В очереди: {user_queue_size}\n 🔄 Обрабатывается: {user_processing}\n 📈 Отправлено за минуту: {user_recent}\n\n📢 Группы/каналы (15 сообщ/мин):\n 📋 В очереди: {group_queue_size}\n 🔄 Обрабатывается: {group_processing}\n 📈 Отправлено за минуту: {group_recent}",
"admin_promo_creation_failed_duplicate": "❌ Ошибка: Промокод {code} уже существует.",
"admin_promo_creation_failed": "❌ Не удалось создать промокод. Пожалуйста, попробуйте позже.",
"admin_active_promos_list_header": "Активные промокоды:",