Compare commits

...
22 Commits
Author SHA1 Message Date
Machka PaslaandGitHub 6330d57c60 Merge pull request #68 from machka-pasla/dev
Tribute hotfix, payments logs
2025-08-07 23:51:55 +03:00
machka-pasla 149ff057a9 Refactor subscription synchronization logic for improved handling
- Enhanced the subscription syncing process to prioritize concrete subscription UUIDs for updates and creations, ensuring idempotency.
- Implemented atomic updates for existing subscriptions and streamlined the creation of new subscriptions when a UUID is available.
- Improved logging for subscription updates and creations to provide clearer feedback on sync actions.
- Added handling for cases where no subscription UUID is present, avoiding unnecessary record creation.
2025-08-07 23:47:19 +03:00
machka-pasla bb7641bb74 Improve amount conversion in TributeService to handle minor units
- Updated the amount handling logic to convert minor currency units (kopecks/cents) to major units before persisting.
- Added error handling for invalid amount inputs to ensure robustness in processing payment data.
- Ensured that the amount is rounded to two decimal places for accurate representation in the system.
2025-08-07 23:39:48 +03:00
machka-pasla 18f65ea493 Refactor notification handling and streamline router registration
- Replaced legacy notification functions with a unified NotificationService for better maintainability and clarity.
- Updated the main bot router registration to utilize a root router, simplifying the inclusion of user and admin routes.
- Removed unused middleware and helper functions to enhance code cleanliness and focus on essential components.
- Improved localization by adding new error messages for user interactions.
2025-08-07 23:30:32 +03:00
machka-pasla 5853b9da63 Enhance payment identifier normalization in TributeService
- Updated the TributeService to normalize provider payment identifiers, prioritizing true payment identifiers over subscription IDs for better uniqueness.
- Implemented fallback logic to append timestamps to subscription IDs, preventing deduplication of renewals.
- Expanded the list of successful charge events to include various payment-related events, improving event handling consistency.
2025-08-07 23:02:23 +03:00
machka-pasla 91cfe0baf3 Add payments feature to admin panel
- Integrated payments functionality into the admin panel by adding a new payments router and corresponding handlers.
- Updated the admin panel actions to include a view payments option, enhancing admin capabilities.
- Implemented new database functions to retrieve successful payment counts and details for export.
- Enhanced localization with new strings for payments management in both English and Russian.
2025-08-07 22:48:09 +03:00
machka-pasla 194f1b9e49 Enhance recent payment log retrieval to filter by succeeded status
- Updated the `get_recent_payment_logs_with_user` function to include a filter for payments with a 'succeeded' status, improving the relevance of retrieved payment logs.
- Adjusted the query structure for better readability and maintainability.
2025-08-07 19:00:27 +03:00
machka-pasla df15cfd25e Refactor promo handler functions to include session management
- Updated the promo_delete_handler, promo_edit_select_handler, and promo_edit_field_handler functions to accept an AsyncSession parameter, improving database interaction consistency.
- Enhanced the handling of expired subscriptions in the PanelWebhookService by modifying notification logic to only send messages if enabled, ensuring better control over user notifications.
- Added an import for the 'and_' function in payment_dal.py to support more complex query conditions.
2025-08-07 18:46:19 +03:00
machka-pasla 57e693fa37 Refactor synchronization messaging for clarity and simplicity
- Updated synchronization messages to provide simpler, more concise feedback to admins during the sync process.
- Replaced detailed sync status messages with straightforward notifications for success, failure, and errors.
- Enhanced localization for new message formats to improve user experience across languages.
2025-08-06 22:26:07 +03:00
Machka PaslaandGitHub 60ea6fff0d Merge pull request #67 from machka-pasla/dev
fix bug
2025-08-06 20:22:20 +03:00
machka-pasla 3cee4b243a Refactor details string handling in statistics display
- Simplified the logic for displaying synchronization details by removing the character limit, ensuring full visibility of the details or defaulting to "N/A" when not available.
- Improved code readability by streamlining the assignment of the details string.
2025-08-06 20:19:45 +03:00
machka-pasla a42f80160b Refactor broadcast confirmation prompt and enhance sync notification handling
- Updated the broadcast confirmation prompt to display the full message instead of a truncated preview.
- Improved error handling in the sync process by removing character limits on error details and ensuring comprehensive logging.
- Added a notification feature to inform admins about the panel synchronization status, including success and failure details.
- Enhanced localization for the broadcast confirmation prompt and added new log messages for sync notifications.
2025-08-06 19:04:03 +03:00
Machka PaslaandGitHub 74548527f1 Merge pull request #66 from machka-pasla/dev
autosync and bug fixes
2025-08-06 18:51:27 +03:00
Machka PaslaandGitHub f8e3bee52a Update .env.example 2025-08-06 18:49:21 +03:00
machka-pasla 649528e165 Refactor bot username retrieval in bulk promo code creation
- Updated the logic to retrieve the bot username dynamically, enhancing the accuracy of CSV links for promo codes.
- Improved error handling to log failures in fetching the bot username, ensuring better traceability.
- Adjusted the promo code service initialization in the start command handler to use the message context, improving consistency in bot interactions.
2025-08-06 18:03:14 +03:00
machka-pasla 859263dc2d Enhance bulk promo code creation with CSV export functionality
- Added the ability to generate and send a CSV file containing created promo codes, improving data accessibility for admins.
- Updated success messages to inform users about the CSV file and the number of created promo codes.
- Refactored the promo code creation logic to include detailed validity information and links for activation in the CSV output.
- Improved localization for bulk promo creation messages to enhance user experience.
2025-08-06 17:58:44 +03:00
machka-pasla a126c05365 Refactor variable assignment in promo detail retrieval
- Updated the variable assignment for the status emoji in the `get_promo_detail_text_and_keyboard` function to improve clarity and consistency in the code.
- Enhanced readability by using more descriptive variable names, aligning with recent refactoring efforts in the promo management handler.
2025-08-06 17:49:06 +03:00
machka-pasla d0c09f9e06 Enhance promo code creation and management logging
- Added logging for successful promo code creation, improving traceability of actions.
- Updated success message formatting to use a new variable for validity display, enhancing clarity for users.
- Refactored variable names in the promo management handler for better readability and consistency.
2025-08-06 17:41:58 +03:00
machka-pasla 9f7171a5c3 Refactor admin sync process for enhanced user and subscription tracking
- Updated the `perform_sync` function to improve tracking of users and subscriptions during synchronization.
- Introduced additional counters for detailed logging, including users without Telegram IDs and those not found in the database.
- Enhanced error handling and logging to provide clearer insights into the synchronization process and outcomes.
- Improved the summary details returned after synchronization, offering a comprehensive overview of the sync results.
2025-08-06 17:21:25 +03:00
machka-pasla 990b08cfdc Enhance subscription handling in admin sync process
- Added logic to retrieve the subscription UUID from the panel user data, improving the accuracy of subscription records.
- Updated the subscription creation process to use the actual subscription UUID when available, with fallback to the user UUID.
- Enhanced logging to provide clearer information on which UUID is being used for each user during synchronization.
2025-08-06 17:00:44 +03:00
machka-pasla 359a3c46a4 Implement panel synchronization functionality in admin handler
- Added a new `perform_sync` function to handle the synchronization of users and subscriptions from the admin panel.
- Enhanced error handling and logging during the sync process, providing detailed feedback on the synchronization status.
- Updated the `sync_command_handler` to utilize the new `perform_sync` function, improving code organization and readability.
- Improved messaging for sync results, including success and error details, to enhance admin user experience.
2025-08-06 16:46:52 +03:00
machka-pasla 00ffec13d8 Implement message queue management and automatic sync on bot startup
- Added initialization of the message queue manager during bot startup, enhancing message handling capabilities.
- Implemented automatic synchronization of the admin panel on startup, providing real-time updates and improved reliability.
- Updated admin handlers to utilize the message queue for broadcasting messages, improving efficiency and error handling.
- Introduced a new command for admins to check the status of message queues, enhancing monitoring and management capabilities.
- Enhanced localization for new features and messages related to queue management and synchronization.
2025-08-06 16:39:35 +03:00
30 changed files with 1483 additions and 579 deletions
+1
View File
@@ -18,6 +18,7 @@ SERVER_STATUS_URL=https://status.yourdomain.tld/status/your_service
TERMS_OF_SERVICE_URL=https://example.com/tos TERMS_OF_SERVICE_URL=https://example.com/tos
SUBSCRIPTION_MINI_APP_URL= SUBSCRIPTION_MINI_APP_URL=
START_COMMAND_DESCRIPTION= START_COMMAND_DESCRIPTION=
DISABLE_WELCOME_MESSAGE=
# Webhook Base URL (used for Telegram and payment providers) # Webhook Base URL (used for Telegram and payment providers)
WEBHOOK_BASE_URL=https://webhooks.yourdomain.tld WEBHOOK_BASE_URL=https://webhooks.yourdomain.tld
+2
View File
@@ -7,6 +7,7 @@ from . import user_management
from . import statistics from . import statistics
from . import sync_admin from . import sync_admin
from . import logs_admin from . import logs_admin
from . import payments
admin_router_aggregate = Router(name="admin_features_router") admin_router_aggregate = Router(name="admin_features_router")
@@ -17,5 +18,6 @@ admin_router_aggregate.include_router(user_management.router)
admin_router_aggregate.include_router(statistics.router) admin_router_aggregate.include_router(statistics.router)
admin_router_aggregate.include_router(sync_admin.router) admin_router_aggregate.include_router(sync_admin.router)
admin_router_aggregate.include_router(logs_admin.router) admin_router_aggregate.include_router(logs_admin.router)
admin_router_aggregate.include_router(payments.router)
__all__ = ("admin_router_aggregate", ) __all__ = ("admin_router_aggregate", )
+27 -8
View File
@@ -1,6 +1,7 @@
import logging import logging
import asyncio import asyncio
from aiogram import Router, F, types, Bot from aiogram import Router, F, types, Bot
from aiogram.exceptions import TelegramRetryAfter
from aiogram.fsm.context import FSMContext from aiogram.fsm.context import FSMContext
from typing import Optional from typing import Optional
@@ -17,6 +18,7 @@ from bot.keyboards.inline.admin_keyboards import (
get_admin_panel_keyboard, get_admin_panel_keyboard,
) )
from bot.middlewares.i18n import JsonI18n from bot.middlewares.i18n import JsonI18n
from bot.utils.message_queue import get_queue_manager
router = Router(name="admin_broadcast_router") router = Router(name="admin_broadcast_router")
@@ -82,8 +84,7 @@ async def process_broadcast_message_handler(
broadcast_entities=entities, broadcast_entities=entities,
) )
preview_snippet = (text[:200] + "...") if len(text) > 200 else text confirmation_prompt = _("admin_broadcast_confirm_prompt", message_preview=text)
confirmation_prompt = _("admin_broadcast_confirm_prompt", message_preview=preview_snippet)
await message.answer( await message.answer(
confirmation_prompt, confirmation_prompt,
@@ -171,22 +172,30 @@ async def confirm_broadcast_callback_handler(
f"Admin {admin_user.id} broadcasting '{text[:50]}...' to {len(user_ids)} users." 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: for uid in user_ids:
try: try:
await bot.send_message( await queue_manager.send_message(
chat_id=uid, chat_id=uid,
text=text, text=text,
entities=entities, entities=entities,
) )
sent_count += 1 sent_count += 1
# Log successful queuing
await message_log_dal.create_message_log( await message_log_dal.create_message_log(
session, session,
{ {
"user_id": admin_user.id, "user_id": admin_user.id,
"telegram_username": admin_user.username, "telegram_username": admin_user.username,
"telegram_first_name": admin_user.first_name, "telegram_first_name": admin_user.first_name,
"event_type": "admin_broadcast_sent", "event_type": "admin_broadcast_queued",
"content": f"To user {uid}: {text[:70]}...", "content": f"To user {uid}: {text[:70]}...",
"is_admin_event": True, "is_admin_event": True,
"target_user_id": uid, "target_user_id": uid,
@@ -195,7 +204,7 @@ async def confirm_broadcast_callback_handler(
except Exception as e: except Exception as e:
failed_count += 1 failed_count += 1
logging.warning( 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( await message_log_dal.create_message_log(
session, session,
@@ -209,7 +218,6 @@ async def confirm_broadcast_callback_handler(
"target_user_id": uid, "target_user_id": uid,
}, },
) )
await asyncio.sleep(0.05)
try: try:
await session.commit() await session.commit()
@@ -217,7 +225,18 @@ async def confirm_broadcast_callback_handler(
await session.rollback() await session.rollback()
logging.error(f"Error committing broadcast logs: {e_commit}") 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( await callback.message.answer(
result_message, result_message,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n), reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
+56
View File
@@ -14,6 +14,7 @@ from bot.keyboards.inline.admin_keyboards import (
from bot.middlewares.i18n import JsonI18n from bot.middlewares.i18n import JsonI18n
from bot.services.panel_api_service import PanelApiService from bot.services.panel_api_service import PanelApiService
from bot.services.subscription_service import SubscriptionService from bot.services.subscription_service import SubscriptionService
from bot.utils.message_queue import get_queue_manager
from . import broadcast as admin_broadcast_handlers from . import broadcast as admin_broadcast_handlers
from .promo import create as admin_promo_create_handlers from .promo import create as admin_promo_create_handlers
@@ -120,6 +121,12 @@ async def admin_panel_actions_callback_handler(
panel_service=panel_service, panel_service=panel_service,
session=session) session=session)
await callback.answer(_("admin_sync_initiated_from_panel")) await callback.answer(_("admin_sync_initiated_from_panel"))
elif action == "queue_status":
await show_queue_status_handler(callback, i18n_data)
elif action == "view_payments":
from . import payments as admin_payments_handlers
await admin_payments_handlers.view_payments_handler(
callback, i18n_data, settings, session)
elif action == "main": elif action == "main":
try: try:
await callback.message.edit_text( await callback.message.edit_text(
@@ -192,3 +199,52 @@ async def admin_section_handler(callback: types.CallbackQuery, state: FSMContext
reply_markup=get_admin_panel_keyboard(i18n, current_lang, settings) reply_markup=get_admin_panel_keyboard(i18n, current_lang, settings)
) )
await callback.answer() 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)
+245
View File
@@ -0,0 +1,245 @@
import logging
import csv
import io
from aiogram import Router, F, types
from aiogram.filters import StateFilter
from aiogram.fsm.context import FSMContext
from datetime import datetime, timedelta, timezone
from typing import Optional, List
from sqlalchemy.ext.asyncio import AsyncSession
from config.settings import Settings
from db.dal import payment_dal
from db.models import Payment
from bot.keyboards.inline.admin_keyboards import get_back_to_admin_panel_keyboard
from aiogram.utils.keyboard import InlineKeyboardBuilder, InlineKeyboardButton
from bot.middlewares.i18n import JsonI18n
router = Router(name="admin_payments_router")
async def get_payments_with_pagination(session: AsyncSession, page: int = 0,
page_size: int = 10) -> tuple[List[Payment], int]:
"""Get payments with pagination and total count."""
offset = page * page_size
# Get total count
total_count = await payment_dal.get_payments_count(session)
# Get payments for current page
payments = await payment_dal.get_recent_payment_logs_with_user(
session, limit=page_size, offset=offset
)
return payments, total_count
def format_payment_text(payment: Payment, i18n: JsonI18n, lang: str) -> str:
"""Format single payment info as text."""
_ = lambda key, **kwargs: i18n.gettext(lang, key, **kwargs)
status_emoji = "" if payment.status == 'succeeded' else (
"" if payment.status in ['pending', 'pending_yookassa'] else ""
)
user_info = f"User {payment.user_id}"
if payment.user and payment.user.username:
user_info += f" (@{payment.user.username})"
elif payment.user and payment.user.first_name:
user_info += f" ({payment.user.first_name})"
payment_date = payment.created_at.strftime('%Y-%m-%d %H:%M') if payment.created_at else "N/A"
provider_text = {
'yookassa': 'YooKassa',
'tribute': 'Tribute',
'telegram_stars': 'Telegram Stars',
'cryptopay': 'CryptoPay'
}.get(payment.provider, payment.provider or 'Unknown')
return (
f"{status_emoji} <b>{payment.amount} {payment.currency}</b>\n"
f"👤 {user_info}\n"
f"💳 {provider_text}\n"
f"📅 {payment_date}\n"
f"📋 {payment.status}\n"
f"📝 {payment.description or 'N/A'}"
)
async def view_payments_handler(callback: types.CallbackQuery, i18n_data: dict,
settings: Settings, session: AsyncSession, page: int = 0):
"""Display paginated list of all payments."""
current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE)
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)
page_size = 5 # Показываем по 5 платежей на странице
payments, total_count = await get_payments_with_pagination(session, page, page_size)
total_pages = (total_count + page_size - 1) // page_size if total_count > 0 else 1
if not payments and page == 0:
await callback.message.edit_text(
_("admin_no_payments_found", default="Платежи не найдены."),
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML"
)
await callback.answer()
return
# Format payments text
text_parts = [_("admin_payments_header", default="💰 <b>Все платежи</b>")]
text_parts.append(f"📊 Показано {len(payments)} из {total_count} платежей (стр. {page + 1}/{total_pages})\n")
for i, payment in enumerate(payments, 1):
text_parts.append(f"<b>{page * page_size + i}.</b> {format_payment_text(payment, i18n, current_lang)}")
text_parts.append("") # Empty line between payments
# Build keyboard with pagination and export
builder = InlineKeyboardBuilder()
# Pagination buttons
nav_buttons = []
if page > 0:
nav_buttons.append(InlineKeyboardButton(text="⬅️", callback_data=f"payments_page:{page-1}"))
nav_buttons.append(InlineKeyboardButton(text=f"{page + 1}/{total_pages}", callback_data="noop"))
if page < total_pages - 1:
nav_buttons.append(InlineKeyboardButton(text="➡️", callback_data=f"payments_page:{page+1}"))
if nav_buttons:
builder.row(*nav_buttons)
# Export and refresh buttons
builder.row(
InlineKeyboardButton(
text=_("admin_export_payments_csv", default="📊 Экспорт CSV"),
callback_data="payments_export_csv"
),
InlineKeyboardButton(
text=_("admin_refresh_payments", default="🔄 Обновить"),
callback_data=f"payments_page:{page}"
)
)
# Back button
builder.row(InlineKeyboardButton(
text=_("back_to_admin_panel_button"),
callback_data="admin_section:stats_monitoring"
))
await callback.message.edit_text(
"\n".join(text_parts),
reply_markup=builder.as_markup(),
parse_mode="HTML"
)
await callback.answer()
@router.callback_query(F.data.startswith("payments_page:"))
async def payments_pagination_handler(callback: types.CallbackQuery, i18n_data: dict,
settings: Settings, session: AsyncSession):
"""Handle pagination for payments list."""
try:
page = int(callback.data.split(":")[1])
await view_payments_handler(callback, i18n_data, settings, session, page)
except (ValueError, IndexError):
await callback.answer("Error processing pagination.", show_alert=True)
@router.callback_query(F.data == "payments_export_csv")
async def export_payments_csv_handler(callback: types.CallbackQuery, i18n_data: dict,
settings: Settings, session: AsyncSession):
"""Export all successful payments to CSV file."""
current_lang = i18n_data.get("current_language", settings.DEFAULT_LANGUAGE)
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
if not i18n:
await callback.answer("Language service error.", show_alert=True)
return
_ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
try:
# Get all successful payments
all_payments = await payment_dal.get_all_succeeded_payments_with_user(session)
if not all_payments:
await callback.answer(
_("admin_no_payments_to_export", default="Нет платежей для экспорта."),
show_alert=True
)
return
# Create CSV in memory
output = io.StringIO()
writer = csv.writer(output)
# Write header
writer.writerow([
_("admin_csv_payment_id", default="ID"),
_("admin_csv_user_id", default="User ID"),
_("admin_csv_username", default="Username"),
_("admin_csv_first_name", default="First Name"),
_("admin_csv_amount", default="Amount"),
_("admin_csv_currency", default="Currency"),
_("admin_csv_provider", default="Provider"),
_("admin_csv_status", default="Status"),
_("admin_csv_description", default="Description"),
_("admin_csv_months", default="Months"),
_("admin_csv_created_at", default="Created At"),
_("admin_csv_provider_payment_id", default="Provider Payment ID")
])
# Write payment data
for payment in all_payments:
writer.writerow([
payment.payment_id,
payment.user_id,
payment.user.username if payment.user and payment.user.username else "",
payment.user.first_name if payment.user and payment.user.first_name else "",
payment.amount,
payment.currency,
payment.provider or "",
payment.status,
payment.description or "",
payment.subscription_duration_months or "",
payment.created_at.strftime('%Y-%m-%d %H:%M:%S') if payment.created_at else "",
payment.provider_payment_id or ""
])
# Prepare file
csv_content = output.getvalue().encode('utf-8-sig') # UTF-8 with BOM for Excel
output.close()
# Generate filename with current date
current_time = datetime.now().strftime('%Y-%m-%d_%H-%M')
filename = f"payments_export_{current_time}.csv"
# Send file
from aiogram.types import BufferedInputFile
file = BufferedInputFile(csv_content, filename=filename)
await callback.message.reply_document(
document=file,
caption=_("admin_payments_export_success",
default="📊 Экспорт платежей завершен!\nВсего записей: {count}",
count=len(all_payments))
)
await callback.answer(
_("admin_export_sent", default="Файл отправлен!"),
show_alert=False
)
except Exception as e:
logging.error(f"Failed to export payments CSV: {e}", exc_info=True)
await callback.answer(f"❌ Ошибка экспорта: {str(e)}", show_alert=True)
@router.callback_query(F.data == "noop")
async def noop_handler(callback: types.CallbackQuery):
"""Handle no-op callback (for pagination display)."""
await callback.answer()
+66 -9
View File
@@ -1,6 +1,8 @@
import logging import logging
import random import random
import string import string
import csv
import io
from aiogram import Router, F, types from aiogram import Router, F, types
from aiogram.filters import StateFilter from aiogram.filters import StateFilter
from aiogram.fsm.context import FSMContext 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: if created_codes:
success_lines.append("\n🎟 <b>Созданные коды:</b>") success_lines.append(f"\n🎟 <b>Создано {len(created_codes)} промокодов</b>")
# Show first 20 codes, then indicate if there are more success_lines.append("📄 CSV файл с промокодами отправлен отдельным сообщением")
codes_to_show = created_codes[:20]
for code in codes_to_show:
success_lines.append(f"<code>{code}</code>")
if len(created_codes) > 20: # Create CSV file
success_lines.append(f"... и еще {len(created_codes) - 20} кодов") 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: if failed_codes:
success_lines.append(f"\n❌ <b>Ошибки ({len(failed_codes)}):</b>") success_lines.append(f"\n❌ <b>Ошибки ({len(failed_codes)}):</b>")
@@ -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), reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML" parse_mode="HTML"
) )
message_obj = callback_or_message.message
except Exception: except Exception:
await callback_or_message.message.answer( message_obj = await callback_or_message.message.answer(
success_text, success_text,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n), reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML" parse_mode="HTML"
) )
await callback_or_message.answer()
else: # Message else: # Message
await callback_or_message.answer( message_obj = await callback_or_message.answer(
success_text, success_text,
reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n), reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n),
parse_mode="HTML" 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() await state.clear()
except Exception as e: except Exception as e:
+6 -2
View File
@@ -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) created_promo = await promo_code_dal.create_promo_code(session, promo_data)
await session.commit() await session.commit()
# Log successful creation
logging.info(f"Promo code '{data['promo_code']}' created with ID {created_promo.promo_code_id}")
# Success message # Success message
valid_until_str = _("admin_promo_unlimited", default="Без ограничений") if not data.get("validity_days") else f"{data['validity_days']} дней"
success_text = _( success_text = _(
"admin_promo_created_success", "admin_promo_created_success",
default="✅ <b>Промокод успешно создан!</b>\n\n" default="✅ <b>Промокод успешно создан!</b>\n\n"
"🎟 Код: <code>{code}</code>\n" "🎟 Код: <code>{code}</code>\n"
"🎁 Бонусные дни: <b>{bonus_days}</b>\n" "🎁 Бонусные дни: <b>{bonus_days}</b>\n"
"📊 Макс. активаций: <b>{max_activations}</b>\n" "📊 Макс. активаций: <b>{max_activations}</b>\n"
"⏰ Срок действия: <b>{validity}</b>", "⏰ Срок действия: <b>{valid_until_str}</b>",
code=data["promo_code"], code=data["promo_code"],
bonus_days=data["bonus_days"], bonus_days=data["bonus_days"],
max_activations=data["max_activations"], 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 if hasattr(callback_or_message, 'message'): # CallbackQuery
+121 -17
View File
@@ -8,7 +8,7 @@ from datetime import datetime, timedelta, timezone
from typing import Optional, List from typing import Optional, List
from sqlalchemy.ext.asyncio import AsyncSession 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.dal import promo_code_dal
from db.models import PromoCode, PromoCodeActivation from db.models import PromoCode, PromoCodeActivation
from bot.states.admin_states import AdminStates from bot.states.admin_states import AdminStates
@@ -19,17 +19,27 @@ from bot.middlewares.i18n import JsonI18n
router = Router(name="promo_manage_router") 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): 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) _ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs)
promo = await promo_code_dal.get_promo_code_by_id(session, promo_id) promo = await promo_code_dal.get_promo_code_by_id(session, promo_id)
if not promo: if not promo:
return None, None return None, None
status = _("admin_promo_status_active") if promo.is_active else _("admin_promo_status_inactive") status_emoji, status = get_promo_status_emoji_and_text(promo, i18n, current_lang)
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")
validity = _("admin_promo_valid_indefinitely") validity = _("admin_promo_valid_indefinitely")
if promo.valid_until: 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) 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( 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"), ""] + [ [_("admin_active_promos_list_header"), ""] + [
f"🎟 <code>{p.code}</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]} <code>{p.code}</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 for p in promo_models
] ]
) )
@@ -77,29 +87,66 @@ async def view_promo_codes_handler(callback: types.CallbackQuery, i18n_data: dic
await callback.answer() await callback.answer()
async def promo_management_handler(callback: types.CallbackQuery, i18n_data: dict, settings: Settings, session: AsyncSession): 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", settings.DEFAULT_LANGUAGE) current_lang = i18n_data.get("current_language", "ru")
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
if not i18n or not callback.message: if not i18n or not callback.message:
await callback.answer("Error processing request.", show_alert=True) await callback.answer("Error processing request.", show_alert=True)
return return
_ = lambda key, **kwargs: i18n.gettext(current_lang, key, **kwargs) _ = 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) page_size = 10 # Количество промокодов на странице
if not promo_models: 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.message.edit_text(_("admin_promo_management_empty"), reply_markup=get_back_to_admin_panel_keyboard(current_lang, i18n), parse_mode="HTML")
await callback.answer() await callback.answer()
return return
builder = InlineKeyboardBuilder() builder = InlineKeyboardBuilder()
for promo in promo_models: 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")) 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() 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:")) @router.callback_query(F.data.startswith("promo_detail:"))
async def promo_detail_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession): async def promo_detail_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession):
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
@@ -227,8 +274,65 @@ async def promo_export_activations_handler(callback: types.CallbackQuery, i18n_d
await callback.answer() 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:")) @router.callback_query(F.data.startswith("promo_delete:"))
async def promo_delete_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession): async def promo_delete_handler(callback: types.CallbackQuery, i18n_data: dict, settings: Settings, session: AsyncSession):
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
current_lang = i18n_data.get("current_language") current_lang = i18n_data.get("current_language")
if not i18n or not callback.message or not current_lang: if not i18n or not callback.message or not current_lang:
@@ -241,7 +345,7 @@ async def promo_delete_handler(callback: types.CallbackQuery, i18n_data: dict, s
if promo: if promo:
await session.commit() await session.commit()
await callback.answer(_("admin_promo_deleted_success", code=promo.code), show_alert=True) 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, settings, session, 0)
else: else:
await callback.answer(_("admin_promo_not_found"), show_alert=True) await callback.answer(_("admin_promo_not_found"), show_alert=True)
except (ValueError, IndexError): except (ValueError, IndexError):
@@ -250,7 +354,7 @@ async def promo_delete_handler(callback: types.CallbackQuery, i18n_data: dict, s
# --- Promo Edit Handlers --- # --- Promo Edit Handlers ---
@router.callback_query(F.data.startswith("promo_edit_select:")) @router.callback_query(F.data.startswith("promo_edit_select:"))
async def promo_edit_select_handler(callback: types.CallbackQuery, i18n_data: dict): async def promo_edit_select_handler(callback: types.CallbackQuery, i18n_data: dict, session: AsyncSession):
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
current_lang = i18n_data.get("current_language") current_lang = i18n_data.get("current_language")
if not i18n or not callback.message or not current_lang: if not i18n or not callback.message or not current_lang:
@@ -269,7 +373,7 @@ async def promo_edit_select_handler(callback: types.CallbackQuery, i18n_data: di
@router.callback_query(F.data.startswith("promo_edit_field:")) @router.callback_query(F.data.startswith("promo_edit_field:"))
async def promo_edit_field_handler(callback: types.CallbackQuery, state: FSMContext, i18n_data: dict): async def promo_edit_field_handler(callback: types.CallbackQuery, state: FSMContext, i18n_data: dict, session: AsyncSession):
i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance") i18n: Optional[JsonI18n] = i18n_data.get("i18n_instance")
current_lang = i18n_data.get("current_language") current_lang = i18n_data.get("current_language")
if not i18n or not callback.message or not current_lang: return if not i18n or not callback.message or not current_lang: return
+1 -3
View File
@@ -197,9 +197,7 @@ async def show_statistics_handler(callback: types.CallbackQuery,
'%Y-%m-%d %H:%M:%S UTC') if sync_time_val else "N/A" '%Y-%m-%d %H:%M:%S UTC') if sync_time_val else "N/A"
details_val = sync_status_model.details details_val = sync_status_model.details
details_str = (details_val[:100] + details_str = details_val or "N/A"
"...") if details_val and len(details_val) > 100 else (
details_val or "N/A")
stats_text_parts.append( stats_text_parts.append(
f" {_('admin_stats_sync_time')}: {sync_time_str}") f" {_('admin_stats_sync_time')}: {sync_time_str}")
+296 -263
View File
@@ -7,6 +7,7 @@ from datetime import datetime, timezone
from config.settings import Settings from config.settings import Settings
from bot.services.panel_api_service import PanelApiService from bot.services.panel_api_service import PanelApiService
from bot.services.notification_service import NotificationService
from db.dal import user_dal, subscription_dal, panel_sync_dal from db.dal import user_dal, subscription_dal, panel_sync_dal
@@ -15,6 +16,262 @@ from bot.middlewares.i18n import JsonI18n
router = Router(name="admin_sync_router") 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")
)
# Prefer syncing by concrete subscription UUID (shortUuid/subscriptionUuid)
subscription_uuid_from_panel = (
panel_user_dict.get("subscriptionUuid")
or panel_user_dict.get("shortUuid")
)
if subscription_uuid_from_panel:
# Try to find subscription by its panel_subscription_uuid first (idempotent)
existing_sub_by_uuid = (
await subscription_dal.get_subscription_by_panel_subscription_uuid(
session, subscription_uuid_from_panel
)
)
if existing_sub_by_uuid:
# Atomic update of all relevant fields
await subscription_dal.update_subscription(
session,
existing_sub_by_uuid.subscription_id,
{
"user_id": actual_user_id,
"panel_user_uuid": panel_uuid,
"end_date": panel_expire_at,
"is_active": panel_status == "ACTIVE",
"status_from_panel": panel_status,
},
)
subscriptions_synced_count += 1
subscriptions_updated += 1
user_was_updated = True
logging.info(
f"Synced existing subscription {existing_sub_by_uuid.subscription_id} for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}"
)
else:
# Create a new subscription only when we have a concrete subscription UUID
sub_payload = {
"user_id": actual_user_id,
"panel_user_uuid": panel_uuid,
"panel_subscription_uuid": subscription_uuid_from_panel,
# Do not guess precise start_date from panel; keep nullable
"start_date": None,
"end_date": panel_expire_at,
"duration_months": None,
"is_active": panel_status == "ACTIVE",
"status_from_panel": panel_status,
"traffic_limit_bytes": settings.user_traffic_limit_bytes,
}
created_sub = await subscription_dal.upsert_subscription(
session, sub_payload
)
subscriptions_synced_count += 1
subscriptions_created += 1
user_was_updated = True
logging.info(
f"Created subscription {created_sub.subscription_id} for user {actual_user_id} by panel_sub_uuid {subscription_uuid_from_panel}"
)
else:
# No subscription UUID from panel: only update an already active subscription for this user/panel UUID
active_sub = await subscription_dal.get_active_subscription_by_user_id(
session, actual_user_id, panel_uuid
)
if active_sub:
await subscription_dal.update_subscription(
session,
active_sub.subscription_id,
{
"end_date": panel_expire_at,
"is_active": panel_status == "ACTIVE",
"status_from_panel": panel_status,
},
)
subscriptions_synced_count += 1
subscriptions_updated += 1
user_was_updated = True
logging.info(
f"Updated active subscription {active_sub.subscription_id} for user {actual_user_id}: expires {panel_expire_at}, status {panel_status}"
)
else:
# Without a concrete subscription UUID we avoid creating new records to keep sync idempotent
logging.debug(
f"No subscriptionUuid for panel user {panel_uuid}; skipped creation for user {actual_user_id}"
)
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)}"
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")) @router.message(Command("sync"))
async def sync_command_handler( async def sync_command_handler(
message_event: Union[types.Message, types.CallbackQuery], message_event: Union[types.Message, types.CallbackQuery],
@@ -48,269 +305,49 @@ async def sync_command_handler(
return return
if isinstance(message_event, types.Message): if isinstance(message_event, types.Message):
await message_event.answer(_("sync_started")) await message_event.answer(_("sync_started_simple"))
logging.info(f"Admin ({message_event.from_user.id}) triggered panel sync.") logging.info(f"Admin ({message_event.from_user.id}) triggered panel sync.")
users_processed_count = 0 # Use the extracted perform_sync function
users_synced_successfully = 0
subscriptions_synced_count = 0
sync_errors = []
try: try:
panel_users_data = await panel_service.get_all_panel_users() sync_result = await perform_sync(panel_service, session, settings, i18n)
if panel_users_data is None: status = sync_result.get("status")
error_msg = "Failed to fetch users from panel or panel API issue." details = sync_result.get("details", "No details available")
sync_errors.append(error_msg) errors = sync_result.get("errors", [])
await panel_sync_dal.update_panel_sync_status(session, "failed", error_msg)
await session.commit() # Simple confirmation message to admin
await bot.send_message(target_chat_id, _("sync_failed", details=error_msg)) if status == "failed":
return await bot.send_message(target_chat_id, _("sync_failed_simple"))
elif status == "completed_with_errors":
if not panel_users_data: await bot.send_message(target_chat_id, _("sync_errors_simple", errors_count=len(errors)))
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()
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}"
else: else:
details_for_db = f"Successfully processed {users_processed_count} users. Synced {subscriptions_synced_count} subscriptions." await bot.send_message(target_chat_id, _("sync_success_simple"))
await panel_sync_dal.update_panel_sync_status( # Send notification to log channel with proper thread handling
session, try:
final_status_type, notification_service = NotificationService(bot, settings, i18n)
details_for_db, await notification_service.notify_panel_sync(
users_processed_count, status, details,
subscriptions_synced_count, sync_result.get("users_processed", 0),
) sync_result.get("subs_synced", 0)
await session.commit() )
except Exception as e_notification:
final_user_message = _( logging.error(f"Failed to send sync notification: {e_notification}")
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)
except Exception as e_sync_global: except Exception as e_sync_global:
await session.rollback() logging.error(f"Global error during /sync command: {e_sync_global}", exc_info=True)
logging.error( await bot.send_message(target_chat_id, _("sync_critical_error"))
f"Global error during /sync command: {e_sync_global}", exc_info=True
) # Send notification to log channel about failure
error_detail_for_db = ( try:
f"An unexpected error occurred during sync: {str(e_sync_global)[:200]}" notification_service = NotificationService(bot, settings, i18n)
) await notification_service.notify_panel_sync(
await panel_sync_dal.update_panel_sync_status( "failed", str(e_sync_global), 0, 0
session, )
"failed", except Exception as e_notification:
error_detail_for_db, logging.error(f"Failed to send sync failure notification: {e_notification}")
users_processed_count,
subscriptions_synced_count,
)
await bot.send_message(
target_chat_id, _("sync_failed", details=error_detail_for_db)
)
@router.message(Command("syncstatus")) @router.message(Command("syncstatus"))
@@ -333,11 +370,7 @@ async def sync_status_command_handler(
) )
details_val = status_record_model.details details_val = status_record_model.details
details_str = ( details_str = details_val or "N/A"
(details_val[:200] + "...")
if details_val and len(details_val) > 200
else (details_val or "N/A")
)
response_text = ( response_text = (
f"<b>{_('admin_stats_last_sync_header')}</b>\n" f"<b>{_('admin_stats_last_sync_header')}</b>\n"
@@ -350,4 +383,4 @@ async def sync_status_command_handler(
else: else:
response_text = _("admin_sync_status_never_run") response_text = _("admin_sync_status_never_run")
await message.answer(response_text, parse_mode="HTML") await message.answer(response_text, parse_mode="HTML")
+4 -2
View File
@@ -216,13 +216,15 @@ async def start_command_handler(message: types.Message,
f"Failed to update existing user {user_id} in session: {e_update}", f"Failed to update existing user {user_id} in session: {e_update}",
exc_info=True) 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 # Auto-apply promo code if provided via start parameter
if promo_code_to_apply: if promo_code_to_apply:
try: try:
from bot.services.promo_code_service import PromoCodeService 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( success, result = await promo_code_service.apply_promo_code(
session, user_id, promo_code_to_apply, current_lang session, user_id, promo_code_to_apply, current_lang
+5 -15
View File
@@ -7,7 +7,7 @@ from datetime import datetime
from config.settings import Settings from config.settings import Settings
from bot.services.subscription_service import SubscriptionService from bot.services.subscription_service import SubscriptionService
from bot.services.panel_api_service import PanelApiService from bot.services.panel_api_service import PanelApiService
from bot.services.notification_service import notify_admin_new_trial from bot.services.notification_service import NotificationService
from bot.keyboards.inline.user_keyboards import ( from bot.keyboards.inline.user_keyboards import (
get_trial_confirmation_keyboard, get_trial_confirmation_keyboard,
get_main_menu_inline_keyboard, get_main_menu_inline_keyboard,
@@ -97,13 +97,8 @@ async def request_trial_confirmation_handler(
) )
# Send notification to admin about new trial # Send notification to admin about new trial
await notify_admin_new_trial( notification_service = NotificationService(callback.bot, settings, i18n)
callback.bot, await notification_service.notify_trial_activation(user_id, end_date_obj)
settings,
i18n,
user_id,
end_date_obj,
)
else: else:
message_key_from_service = ( message_key_from_service = (
activation_result.get("message_key", "trial_activation_failed") activation_result.get("message_key", "trial_activation_failed")
@@ -264,13 +259,8 @@ async def confirm_activate_trial_handler(
) )
if activation_result and activation_result.get("activated") and end_date_obj: if activation_result and activation_result.get("activated") and end_date_obj:
await notify_admin_new_trial( notification_service = NotificationService(callback.bot, settings, i18n)
callback.bot, await notification_service.notify_trial_activation(user_id, end_date_obj)
settings,
i18n,
user_id,
end_date_obj,
)
@router.callback_query(F.data == "main_action:cancel_trial") @router.callback_query(F.data == "main_action:cancel_trial")
+6 -2
View File
@@ -39,12 +39,14 @@ def get_stats_monitoring_keyboard(i18n_instance, lang: str) -> InlineKeyboardMar
builder.button(text=_(key="admin_stats_button"), builder.button(text=_(key="admin_stats_button"),
callback_data="admin_action:stats") callback_data="admin_action:stats")
builder.button(text=_(key="admin_view_payments_button", default="💰 Платежи"),
callback_data="admin_action:view_payments")
builder.button(text=_(key="admin_view_logs_menu_button"), builder.button(text=_(key="admin_view_logs_menu_button"),
callback_data="admin_action:view_logs_menu") callback_data="admin_action:view_logs_menu")
builder.button(text=_(key="back_to_admin_panel_button"), builder.button(text=_(key="back_to_admin_panel_button"),
callback_data="admin_action:main") callback_data="admin_action:main")
builder.adjust(2, 1) builder.adjust(2, 1, 1)
return builder.as_markup() return builder.as_markup()
@@ -105,10 +107,12 @@ def get_system_functions_keyboard(i18n_instance, lang: str) -> InlineKeyboardMar
callback_data="admin_action:broadcast") callback_data="admin_action:broadcast")
builder.button(text=_(key="admin_sync_panel_button"), builder.button(text=_(key="admin_sync_panel_button"),
callback_data="admin_action:sync_panel") 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"), builder.button(text=_(key="back_to_admin_panel_button"),
callback_data="admin_action:main") callback_data="admin_action:main")
builder.adjust(2, 1) builder.adjust(2, 1, 1)
return builder.as_markup() return builder.as_markup()
+36 -61
View File
@@ -1,17 +1,10 @@
import logging import logging
import asyncio import asyncio
from typing import Callable, Dict, Any, Awaitable, Optional from typing import Dict, Any, Optional
from aiogram import Bot, Dispatcher, BaseMiddleware, Router, F from aiogram import Bot, Dispatcher
from aiogram.types import ( from aiogram.types import (MenuButtonDefault, MenuButtonWebApp, WebAppInfo, BotCommand)
Update,
MenuButtonDefault,
MenuButtonWebApp,
WebAppInfo,
BotCommand,
)
from aiogram.enums import ParseMode from aiogram.enums import ParseMode
from aiogram.filters import CommandStart, Command
from aiogram.client.default import DefaultBotProperties from aiogram.client.default import DefaultBotProperties
from aiogram.webhook.aiohttp_server import SimpleRequestHandler, setup_application from aiogram.webhook.aiohttp_server import SimpleRequestHandler, setup_application
from aiogram.fsm.storage.memory import MemoryStorage from aiogram.fsm.storage.memory import MemoryStorage
@@ -24,13 +17,11 @@ from config.settings import Settings
from db.database_setup import init_db_connection from db.database_setup import init_db_connection
from bot.middlewares.i18n import I18nMiddleware, get_i18n_instance, JsonI18n from bot.middlewares.i18n import I18nMiddleware, get_i18n_instance, JsonI18n
from bot.middlewares.db_session import DBSessionMiddleware
from bot.middlewares.ban_check_middleware import BanCheckMiddleware from bot.middlewares.ban_check_middleware import BanCheckMiddleware
from bot.middlewares.action_logger_middleware import ActionLoggerMiddleware from bot.middlewares.action_logger_middleware import ActionLoggerMiddleware
from bot.handlers.user import user_router_aggregate from bot.routers import build_root_router
from bot.handlers.admin import admin_router_aggregate
from bot.handlers import inline_mode
from bot.filters.admin_filter import AdminFilter
from bot.services.yookassa_service import YooKassaService from bot.services.yookassa_service import YooKassaService
from bot.services.panel_api_service import PanelApiService from bot.services.panel_api_service import PanelApiService
@@ -42,56 +33,12 @@ from bot.services.tribute_service import TributeService, tribute_webhook_route
from bot.services.crypto_pay_service import CryptoPayService, cryptopay_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.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):
def __init__(self, async_session_factory: sessionmaker):
super().__init__()
self.async_session_factory = async_session_factory
async def __call__(
self,
handler: Callable[[Update, Dict[str, Any]], Awaitable[Any]],
event: Update,
data: Dict[str, Any],
) -> Any:
if self.async_session_factory is None:
logging.critical("DBSessionMiddleware: async_session_factory is None!")
raise RuntimeError(
"async_session_factory not provided to DBSessionMiddleware"
)
async with self.async_session_factory() as session:
data["session"] = session
try:
result = await handler(event, data)
await session.commit()
return result
except Exception:
await session.rollback()
logging.error(
"DBSessionMiddleware: Exception caused rollback.", exc_info=True
)
raise
async def register_all_routers(dp: Dispatcher, settings: Settings): async def register_all_routers(dp: Dispatcher, settings: Settings):
dp.include_router(user_router_aggregate) dp.include_router(build_root_router(settings))
# Add inline mode router (available for all users)
dp.include_router(inline_mode.router)
admin_main_router = Router(name="admin_main_filtered_router")
admin_filter_instance = AdminFilter(admin_ids=settings.ADMIN_IDS)
admin_main_router.message.filter(admin_filter_instance)
admin_main_router.callback_query.filter(admin_filter_instance)
admin_main_router.include_router(admin_router_aggregate)
dp.include_router(admin_main_router)
logging.info("All application routers registered.") logging.info("All application routers registered.")
@@ -191,6 +138,34 @@ async def on_startup_configured(dispatcher: Dispatcher):
except Exception as e: except Exception as e:
logging.error(f"STARTUP: Failed to set bot commands: {e}", exc_info=True) 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.") logging.info("STARTUP: Bot on_startup_configured completed.")
+40
View File
@@ -0,0 +1,40 @@
import logging
from typing import Callable, Dict, Any, Awaitable
from aiogram import BaseMiddleware
from aiogram.types import Update
from sqlalchemy.orm import sessionmaker
class DBSessionMiddleware(BaseMiddleware):
def __init__(self, async_session_factory: sessionmaker):
super().__init__()
self.async_session_factory = async_session_factory
async def __call__(
self,
handler: Callable[[Update, Dict[str, Any]], Awaitable[Any]],
event: Update,
data: Dict[str, Any],
) -> Any:
if self.async_session_factory is None:
logging.critical("DBSessionMiddleware: async_session_factory is None!")
raise RuntimeError(
"async_session_factory not provided to DBSessionMiddleware"
)
async with self.async_session_factory() as session:
data["session"] = session
try:
result = await handler(event, data)
await session.commit()
return result
except Exception:
await session.rollback()
logging.error(
"DBSessionMiddleware: Exception caused rollback.", exc_info=True
)
raise
+26
View File
@@ -0,0 +1,26 @@
from aiogram import Router
from bot.handlers.user import user_router_aggregate
from bot.handlers import inline_mode
from bot.handlers.admin import admin_router_aggregate
from bot.filters.admin_filter import AdminFilter
from config.settings import Settings
def build_root_router(settings: Settings) -> Router:
root = Router(name="root")
# Public routers
root.include_router(user_router_aggregate)
root.include_router(inline_mode.router)
# Admin routers behind filter
admin_main_router = Router(name="admin_main_filtered_router")
admin_filter_instance = AdminFilter(admin_ids=settings.ADMIN_IDS)
admin_main_router.message.filter(admin_filter_instance)
admin_main_router.callback_query.filter(admin_filter_instance)
admin_main_router.include_router(admin_router_aggregate)
root.include_router(admin_main_router)
return root
+76 -46
View File
@@ -2,12 +2,14 @@ import logging
import asyncio import asyncio
from aiogram import Bot from aiogram import Bot
from aiogram.utils.text_decorations import html_decoration as hd from aiogram.utils.text_decorations import html_decoration as hd
from aiogram.exceptions import TelegramRetryAfter
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Optional, Union, Dict, Any from typing import Optional, Union, Dict, Any
from config.settings import Settings from config.settings import Settings
from sqlalchemy.orm import sessionmaker from sqlalchemy.orm import sessionmaker
from bot.middlewares.i18n import JsonI18n from bot.middlewares.i18n import JsonI18n
from bot.utils.message_queue import get_queue_manager
class NotificationService: class NotificationService:
@@ -19,16 +21,30 @@ class NotificationService:
self.i18n = i18n self.i18n = i18n
async def _send_to_log_channel(self, message: str, thread_id: Optional[int] = None): 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: if not self.settings.LOG_CHAT_ID:
return 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: try:
# Use thread_id if provided, otherwise use from settings # Use thread_id if provided, otherwise use from settings
final_thread_id = thread_id or self.settings.LOG_THREAD_ID final_thread_id = thread_id or self.settings.LOG_THREAD_ID
kwargs = { kwargs = {
"chat_id": self.settings.LOG_CHAT_ID,
"text": message, "text": message,
"parse_mode": "HTML", "parse_mode": "HTML",
"disable_web_page_preview": True "disable_web_page_preview": True
@@ -38,26 +54,42 @@ class NotificationService:
if final_thread_id: if final_thread_id:
kwargs["message_thread_id"] = 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: 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): 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: if not self.settings.ADMIN_IDS:
return 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: for admin_id in self.settings.ADMIN_IDS:
try: try:
await self.bot.send_message( await queue_manager.send_message(
chat_id=admin_id, chat_id=admin_id,
text=message, text=message,
parse_mode="HTML", parse_mode="HTML",
disable_web_page_preview=True disable_web_page_preview=True
) )
except Exception as e: 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, async def notify_new_user_registration(self, user_id: int, username: Optional[str] = None,
first_name: Optional[str] = None, first_name: Optional[str] = None,
@@ -189,6 +221,42 @@ class NotificationService:
# Send to log channel # Send to log channel
await self._send_to_log_channel(message) await self._send_to_log_channel(message)
async def notify_panel_sync(self, status: str, details: str,
users_processed: int, subs_synced: int,
username: Optional[str] = None):
"""Send notification about panel synchronization"""
if not getattr(self.settings, 'LOG_PANEL_SYNC', True):
return
admin_lang = self.settings.DEFAULT_LANGUAGE
_ = lambda k, **kw: self.i18n.gettext(admin_lang, k, **kw) if self.i18n else k
# Status emoji based on sync result
status_emoji = {
"completed": "",
"completed_with_errors": "⚠️",
"failed": ""
}.get(status, "🔄")
message = _(
"log_panel_sync",
default="{status_emoji} <b>Синхронизация с панелью</b>\n\n"
"📊 Статус: <b>{status}</b>\n"
"👥 Обработано пользователей: <b>{users_processed}</b>\n"
"📋 Синхронизировано подписок: <b>{subs_synced}</b>\n"
"🕐 Время: {timestamp}\n\n"
"📝 Детали:\n{details}",
status_emoji=status_emoji,
status=status,
users_processed=users_processed,
subs_synced=subs_synced,
timestamp=datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S %Z"),
details=details
)
# Send to log channel
await self._send_to_log_channel(message)
async def notify_suspicious_promo_attempt( async def notify_suspicious_promo_attempt(
self, user_id: int, suspicious_input: str, self, user_id: int, suspicious_input: str,
username: Optional[str] = None, first_name: Optional[str] = None): username: Optional[str] = None, first_name: Optional[str] = None):
@@ -227,42 +295,4 @@ class NotificationService:
if to_admins: if to_admins:
await self._send_to_admins(message) await self._send_to_admins(message)
# Removed legacy helper functions that duplicated NotificationService API
# Legacy functions for backward compatibility
async def notify_admins(bot: Bot, settings: Settings, i18n: JsonI18n,
message_key: str, parse_mode: str | None = None,
**kwargs) -> None:
if not settings.ADMIN_IDS:
return
admin_lang = settings.DEFAULT_LANGUAGE
msg = i18n.gettext(admin_lang, message_key, **kwargs)
for admin_id in settings.ADMIN_IDS:
try:
await bot.send_message(admin_id, msg, parse_mode=parse_mode)
except Exception as e:
logging.error(f"Failed to send admin notification to {admin_id}: {e}")
async def notify_admin_new_trial(bot: Bot, settings: Settings, i18n: JsonI18n,
user_id: int, end_date: datetime) -> None:
"""Send notification to admins about new trial activation (legacy)"""
notification_service = NotificationService(bot, settings, i18n)
await notification_service.notify_trial_activation(user_id, end_date)
async def notify_admin_promo_activation(bot: Bot, settings: Settings,
i18n: JsonI18n, user_id: int,
code: str,
bonus_days: int) -> None:
await notify_admins(
bot,
settings,
i18n,
"admin_promo_activation_notification",
user_id=user_id,
code=code,
bonus_days=bonus_days,
)
+12 -10
View File
@@ -133,18 +133,20 @@ class PanelWebhookService:
user_name=first_name, user_name=first_name,
end_date=user_payload.get("expireAt", "")[:10], end_date=user_payload.get("expireAt", "")[:10],
) )
elif event_name == "user.expired" and self.settings.SUBSCRIPTION_NOTIFY_ON_EXPIRE: elif event_name == "user.expired":
# Check if this is a tribute user that should be auto-renewed # Check if this is a tribute user that should be auto-renewed (regardless of notification settings)
await self._handle_expired_subscription(session, user_id, user_payload, lang, markup, first_name) await self._handle_expired_subscription(session, user_id, user_payload, lang, markup, first_name)
await self._send_message( # Send notification only if enabled
user_id, if self.settings.SUBSCRIPTION_NOTIFY_ON_EXPIRE:
lang, await self._send_message(
"subscription_expired_notification", user_id,
reply_markup=markup, lang,
user_name=first_name, "subscription_expired_notification",
end_date=user_payload.get("expireAt", "")[:10], reply_markup=markup,
) user_name=first_name,
end_date=user_payload.get("expireAt", "")[:10],
)
elif event_name == "user.expired_24_hours_ago" and self.settings.SUBSCRIPTION_NOTIFY_AFTER_EXPIRE: elif event_name == "user.expired_24_hours_ago" and self.settings.SUBSCRIPTION_NOTIFY_AFTER_EXPIRE:
await self._send_message( await self._send_message(
user_id, user_id,
-2
View File
@@ -55,7 +55,6 @@ class PromoCodeService:
reason=f"promo code {code_input_upper}") reason=f"promo code {code_input_upper}")
if new_end_date: if new_end_date:
activation_recorded = await promo_code_dal.record_promo_activation( activation_recorded = await promo_code_dal.record_promo_activation(
session, promo_data.promo_code_id, user_id, payment_id=None) session, promo_data.promo_code_id, user_id, payment_id=None)
promo_incremented = await promo_code_dal.increment_promo_code_usage( promo_incremented = await promo_code_dal.increment_promo_code_usage(
@@ -83,5 +82,4 @@ class PromoCodeService:
) )
return False, _("error_applying_promo_bonus") return False, _("error_applying_promo_bonus")
else: else:
return False, _("error_applying_promo_bonus") return False, _("error_applying_promo_bonus")
+52 -36
View File
@@ -34,10 +34,7 @@ class SubscriptionService:
else self.settings.DEFAULT_LANGUAGE else self.settings.DEFAULT_LANGUAGE
) )
async def has_had_any_subscription( async def has_had_any_subscription(self, session: AsyncSession, user_id: int) -> bool:
self, session: AsyncSession, user_id: int
) -> bool:
return await subscription_dal.has_any_subscription_for_user(session, user_id) return await subscription_dal.has_any_subscription_for_user(session, user_id)
async def _notify_admin_panel_user_creation_failed(self, user_id: int): async def _notify_admin_panel_user_creation_failed(self, user_id: int):
@@ -348,19 +345,12 @@ class SubscriptionService:
"message_key": "trial_activation_failed_db", "message_key": "trial_activation_failed_db",
} }
panel_update_payload: Dict[str, Any] = { panel_update_payload = self._build_panel_update_payload(
"uuid": panel_user_uuid, panel_user_uuid=panel_user_uuid,
"expireAt": end_date.isoformat(timespec="milliseconds").replace( expire_at=end_date,
"+00:00", "Z" status="ACTIVE",
), traffic_limit_bytes=self.settings.trial_traffic_limit_bytes,
"status": "ACTIVE", )
"trafficLimitBytes": self.settings.trial_traffic_limit_bytes,
"trafficLimitStrategy": self.settings.USER_TRAFFIC_STRATEGY,
}
if self.settings.parsed_user_squad_uuids:
panel_update_payload["activeInternalSquads"] = (
self.settings.parsed_user_squad_uuids
)
updated_panel_user = await self.panel_service.update_user_details_on_panel( updated_panel_user = await self.panel_service.update_user_details_on_panel(
panel_user_uuid, panel_update_payload panel_user_uuid, panel_update_payload
@@ -495,19 +485,12 @@ class SubscriptionService:
) )
return None return None
panel_update_payload = { panel_update_payload = self._build_panel_update_payload(
"uuid": panel_user_uuid, panel_user_uuid=panel_user_uuid,
"expireAt": final_end_date.isoformat(timespec="milliseconds").replace( expire_at=final_end_date,
"+00:00", "Z" status="ACTIVE",
), traffic_limit_bytes=self.settings.user_traffic_limit_bytes,
"status": "ACTIVE", )
"trafficLimitBytes": self.settings.user_traffic_limit_bytes,
"trafficLimitStrategy": self.settings.USER_TRAFFIC_STRATEGY,
}
if self.settings.parsed_user_squad_uuids:
panel_update_payload["activeInternalSquads"] = (
self.settings.parsed_user_squad_uuids
)
updated_panel_user = await self.panel_service.update_user_details_on_panel( updated_panel_user = await self.panel_service.update_user_details_on_panel(
panel_user_uuid, panel_update_payload panel_user_uuid, panel_update_payload
@@ -563,6 +546,10 @@ class SubscriptionService:
) )
start_date = datetime.now(timezone.utc) start_date = datetime.now(timezone.utc)
new_end_date_obj = start_date + timedelta(days=bonus_days) 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 = { bonus_sub_payload = {
"user_id": user_id, "user_id": user_id,
"panel_user_uuid": panel_uuid, "panel_user_uuid": panel_uuid,
@@ -572,7 +559,7 @@ class SubscriptionService:
"duration_months": 0, "duration_months": 0,
"is_active": True, "is_active": True,
"status_from_panel": "ACTIVE_BONUS", "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( await subscription_dal.deactivate_other_active_subscriptions(
session, panel_uuid, panel_sub_uuid session, panel_uuid, panel_sub_uuid
@@ -593,14 +580,19 @@ class SubscriptionService:
) )
if updated_sub_model: if updated_sub_model:
# Prepare panel update payload
panel_update_payload = self._build_panel_update_payload(
expire_at=new_end_date_obj,
traffic_limit_bytes=(
self.settings.user_traffic_limit_bytes if "promo code" in reason.lower() else None
),
include_uuid=False,
)
panel_update_success = ( panel_update_success = (
await self.panel_service.update_user_details_on_panel( await self.panel_service.update_user_details_on_panel(
panel_uuid, panel_uuid,
{ panel_update_payload,
"expireAt": new_end_date_obj.isoformat(
timespec="milliseconds"
).replace("+00:00", "Z")
},
) )
) )
if not panel_update_success: if not panel_update_success:
@@ -762,3 +754,27 @@ class SubscriptionService:
logging.warning( logging.warning(
f"Could not find subscription for user {user_id} ending at {subscription_end_date.isoformat()} to update notification time." f"Could not find subscription for user {user_id} ending at {subscription_end_date.isoformat()} to update notification time."
) )
# Helpers
def _build_panel_update_payload(
self,
*,
panel_user_uuid: Optional[str] = None,
expire_at: Optional[datetime] = None,
status: Optional[str] = None,
traffic_limit_bytes: Optional[int] = None,
include_uuid: bool = True,
) -> Dict[str, Any]:
payload: Dict[str, Any] = {}
if include_uuid and panel_user_uuid:
payload["uuid"] = panel_user_uuid
if expire_at is not None:
payload["expireAt"] = expire_at.isoformat(timespec="milliseconds").replace("+00:00", "Z")
if status is not None:
payload["status"] = status
if traffic_limit_bytes is not None:
payload["trafficLimitBytes"] = traffic_limit_bytes
payload["trafficLimitStrategy"] = self.settings.USER_TRAFFIC_STRATEGY
if self.settings.parsed_user_squad_uuids:
payload["activeInternalSquads"] = self.settings.parsed_user_squad_uuids
return payload
+61 -74
View File
@@ -39,11 +39,16 @@ def convert_period_to_months(period: Optional[str]) -> int:
class TributeService: class TributeService:
def __init__(self, bot: Bot, settings: Settings, i18n: JsonI18n, def __init__(
async_session_factory: sessionmaker, self,
panel_service: PanelApiService, bot: Bot,
subscription_service: SubscriptionService, settings: Settings,
referral_service: ReferralService): i18n: JsonI18n,
async_session_factory: sessionmaker,
panel_service: PanelApiService,
subscription_service: SubscriptionService,
referral_service: ReferralService,
):
self.bot = bot self.bot = bot
self.settings = settings self.settings = settings
self.i18n = i18n self.i18n = i18n
@@ -52,8 +57,7 @@ class TributeService:
self.subscription_service = subscription_service self.subscription_service = subscription_service
self.referral_service = referral_service self.referral_service = referral_service
async def handle_webhook(self, raw_body: bytes, async def handle_webhook(self, raw_body: bytes, signature_header: Optional[str]) -> web.Response:
signature_header: Optional[str]) -> web.Response:
settings = self.settings settings = self.settings
bot = self.bot bot = self.bot
i18n = self.i18n i18n = self.i18n
@@ -79,60 +83,60 @@ class TributeService:
json.dumps(payload, ensure_ascii=False), json.dumps(payload, ensure_ascii=False),
) )
event_name = payload.get('name') # Tribute webhook spec: only two events are sent
data = payload.get('payload', {}) # name: new_subscription | cancelled_subscription
user_id = data.get('telegram_user_id') event_name = payload.get("name")
price_val = ( data = payload.get("payload", {})
data.get('amount')
or data.get('amount_paid')
or data.get('price')
)
if not user_id or price_val is None: # Mandatory routing fields
return web.Response(status=200, text="ok_missing_fields") user_id = data.get("telegram_user_id")
if not user_id:
return web.Response(status=400, text="missing_telegram_user_id")
period_val = data.get('period') period_val = data.get("period")
months = convert_period_to_months(period_val) months = convert_period_to_months(period_val)
price_rub = price_val / 100
# Tribute sends amount in minor units (kopecks/cents). Convert to major units before persisting.
amount_value = data.get("amount") or data.get("price")
currency = (data.get("currency") or settings.DEFAULT_CURRENCY_SYMBOL or "RUB").upper()
if amount_value is not None:
try:
amount_minor_units = float(amount_value)
except (TypeError, ValueError):
amount_minor_units = 0.0
amount_float = round(amount_minor_units / 100.0, 2)
else:
amount_float = 0.0
async with async_session_factory() as session: async with async_session_factory() as session:
if event_name == 'new_subscription': if event_name == "new_subscription":
provider_payment_id = str(data.get('subscription_id')) # Build a stable provider payment id from subscription and timestamps
existing_payment = await payment_dal.get_payment_by_provider_payment_id( provider_payment_id = str(data.get("subscription_id"))
session, provider_payment_id) # Idempotent ensure payment
if existing_payment: payment_record = await payment_dal.ensure_payment_with_provider_id(
logging.info( session,
"Duplicate Tribute payment webhook ignored for provider_payment_id %s", user_id=int(user_id),
provider_payment_id, amount=amount_float,
) currency=currency,
payment_record = existing_payment months=months,
else: description="Tribute subscription",
payment_record = await payment_dal.create_payment_record( provider="tribute",
session, provider_payment_id=provider_payment_id,
{ )
'user_id': user_id,
'amount': float(price_rub),
'currency': 'RUB',
'status': 'succeeded',
'description': 'Tribute subscription',
'subscription_duration_months': months,
'provider_payment_id': provider_payment_id,
'provider': 'tribute',
},
)
activation_details = await subscription_service.activate_subscription( activation_details = await subscription_service.activate_subscription(
session, session,
user_id, int(user_id),
months, months,
float(price_rub), float(amount_float),
payment_record.payment_id, payment_record.payment_id,
provider='tribute', provider="tribute",
) )
referral_bonus = await referral_service.apply_referral_bonuses_for_payment( referral_bonus = await referral_service.apply_referral_bonuses_for_payment(
session, user_id, months) session, int(user_id), months)
await session.commit() await session.commit()
db_user = await user_dal.get_user_by_id(session, user_id) db_user = await user_dal.get_user_by_id(session, int(user_id))
lang = db_user.language_code if db_user and db_user.language_code else settings.DEFAULT_LANGUAGE lang = db_user.language_code if db_user and db_user.language_code else settings.DEFAULT_LANGUAGE
_ = lambda k, **kw: i18n.gettext(lang, k, **kw) _ = lambda k, **kw: i18n.gettext(lang, k, **kw)
@@ -177,7 +181,7 @@ class TributeService:
try: try:
await bot.send_message( await bot.send_message(
user_id, int(user_id),
success_msg, success_msg,
reply_markup=markup, reply_markup=markup,
parse_mode="HTML", parse_mode="HTML",
@@ -190,21 +194,19 @@ class TributeService:
# Send notification about payment # Send notification about payment
try: try:
notification_service = NotificationService(bot, settings, i18n) notification_service = NotificationService(bot, settings, i18n)
user = await user_dal.get_user_by_id(session, user_id) user = await user_dal.get_user_by_id(session, int(user_id))
await notification_service.notify_payment_received( await notification_service.notify_payment_received(
user_id=user_id, user_id=int(user_id),
amount=float(price_rub), amount=float(amount_float),
currency="RUB", currency=currency,
months=months, months=months,
payment_provider="tribute", payment_provider="tribute",
username=user.username if user else None username=user.username if user else None
) )
except Exception as e: except Exception as e:
logging.error(f"Failed to send tribute payment notification: {e}") logging.error(f"Failed to send tribute payment notification: {e}")
elif event_name == "cancelled_subscription":
elif event_name == 'subscription_cancelled': await self._handle_tribute_cancellation(session, int(user_id), bot, i18n)
# Handle tribute subscription cancellation
await self._handle_tribute_cancellation(session, user_id, bot, i18n)
else: else:
await session.commit() await session.commit()
@@ -218,22 +220,7 @@ class TributeService:
try: try:
# Set all user's subscriptions to expire in 1 day (grace period) # Set all user's subscriptions to expire in 1 day (grace period)
grace_end_date = datetime.now(timezone.utc) + timedelta(days=1) await subscription_dal.set_user_subscriptions_cancelled_with_grace(session, user_id, grace_days=1)
# Get all active subscriptions for the user
user_subs = await subscription_dal.get_active_subscriptions_for_user(session, user_id)
for sub in user_subs:
await subscription_dal.update_subscription(
session,
sub.subscription_id,
{
'end_date': grace_end_date,
'status_from_panel': 'CANCELLED',
'skip_notifications': True # Skip future notifications for cancelled subs
}
)
await session.commit() await session.commit()
# Send notification about cancellation if enabled # Send notification about cancellation if enabled
@@ -256,7 +243,7 @@ class TributeService:
try: try:
await bot.send_message( await bot.send_message(
user_id, int(user_id),
cancellation_msg, cancellation_msg,
reply_markup=markup, reply_markup=markup,
parse_mode="HTML" parse_mode="HTML"
+1
View File
@@ -0,0 +1 @@
# Bot utilities package
+183
View File
@@ -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
+1
View File
@@ -113,6 +113,7 @@ class Settings(BaseSettings):
SUBSCRIPTION_MINI_APP_URL: Optional[str] = Field(default=None) SUBSCRIPTION_MINI_APP_URL: Optional[str] = Field(default=None)
START_COMMAND_DESCRIPTION: 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 mode thumbnail URLs
INLINE_REFERRAL_THUMBNAIL_URL: str = Field(default="https://cdn-icons-png.flaticon.com/512/1077/1077114.png") INLINE_REFERRAL_THUMBNAIL_URL: str = Field(default="https://cdn-icons-png.flaticon.com/512/1077/1077114.png")
+53 -3
View File
@@ -2,7 +2,7 @@ import logging
from typing import Optional, List, Dict, Any from typing import Optional, List, Dict, Any
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.future import select from sqlalchemy.future import select
from sqlalchemy import update, func from sqlalchemy import update, func, and_
from sqlalchemy.orm import selectinload from sqlalchemy.orm import selectinload
from db.models import Payment, User from db.models import Payment, User
@@ -47,6 +47,38 @@ async def get_payment_by_provider_payment_id(
return result.scalar_one_or_none() return result.scalar_one_or_none()
async def ensure_payment_with_provider_id(
session: AsyncSession,
*,
user_id: int,
amount: float,
currency: str,
months: int,
description: str,
provider: str,
provider_payment_id: str) -> Payment:
"""Idempotently create a payment record for a provider event.
If a payment with the same provider_payment_id already exists, returns it.
Otherwise creates a new succeeded payment with provided data.
"""
existing = await get_payment_by_provider_payment_id(session, provider_payment_id)
if existing:
return existing
payment_payload: Dict[str, Any] = {
"user_id": user_id,
"amount": float(amount),
"currency": currency,
"status": "succeeded",
"description": description,
"subscription_duration_months": months,
"provider_payment_id": provider_payment_id,
"provider": provider,
}
return await create_payment_record(session, payment_payload)
async def get_payment_by_db_id(session: AsyncSession, async def get_payment_by_db_id(session: AsyncSession,
payment_db_id: int) -> Optional[Payment]: payment_db_id: int) -> Optional[Payment]:
@@ -82,8 +114,26 @@ async def update_payment_status_by_db_id(
async def get_recent_payment_logs_with_user(session: AsyncSession, async def get_recent_payment_logs_with_user(session: AsyncSession,
limit: int = 20, limit: int = 20,
offset: int = 0) -> List[Payment]: offset: int = 0) -> List[Payment]:
stmt = (select(Payment).options(selectinload(Payment.user)).order_by( stmt = (select(Payment).options(selectinload(Payment.user))
Payment.created_at.desc()).limit(limit).offset(offset)) .where(Payment.status == 'succeeded')
.order_by(Payment.created_at.desc())
.limit(limit).offset(offset))
result = await session.execute(stmt)
return result.scalars().all()
async def get_payments_count(session: AsyncSession) -> int:
"""Get total count of successful payments."""
stmt = select(func.count(Payment.payment_id)).where(Payment.status == 'succeeded')
result = await session.execute(stmt)
return result.scalar() or 0
async def get_all_succeeded_payments_with_user(session: AsyncSession) -> List[Payment]:
"""Get all successful payments with user data for export."""
stmt = (select(Payment).options(selectinload(Payment.user))
.where(Payment.status == 'succeeded')
.order_by(Payment.created_at.desc()))
result = await session.execute(stmt) result = await session.execute(stmt)
return result.scalars().all() return result.scalars().all()
+8
View File
@@ -64,6 +64,14 @@ async def get_all_promo_codes_with_details(session: AsyncSession, limit: int = 5
return result.scalars().all() 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]: 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.""" """Get activation history for a specific promo code with optional pagination."""
stmt = (select(PromoCodeActivation) stmt = (select(PromoCodeActivation)
+23
View File
@@ -53,6 +53,29 @@ async def update_subscription(
return sub return sub
async def set_user_subscriptions_cancelled_with_grace(
session: AsyncSession, user_id: int, grace_days: int = 1) -> int:
"""Mark all active user subscriptions as cancelled with a short grace period.
Sets end_date to now + grace_days, status_from_panel to 'CANCELLED', and
skip future notifications to reduce noise after cancellation.
Returns number of updated rows.
"""
from datetime import datetime, timezone, timedelta
grace_end = datetime.now(timezone.utc) + timedelta(days=grace_days)
stmt = (
update(Subscription)
.where(Subscription.user_id == user_id, Subscription.is_active == True)
.values(
end_date=grace_end,
status_from_panel="CANCELLED",
skip_notifications=True,
)
)
result = await session.execute(stmt)
return result.rowcount or 0
async def upsert_subscription(session: AsyncSession, async def upsert_subscription(session: AsyncSession,
sub_payload: Dict[str, Any]) -> Subscription: sub_payload: Dict[str, Any]) -> Subscription:
panel_sub_uuid = sub_payload.get("panel_subscription_uuid") panel_sub_uuid = sub_payload.get("panel_subscription_uuid")
+1 -14
View File
@@ -30,20 +30,7 @@ async def get_user_by_panel_uuid(
return result.scalar_one_or_none() return result.scalar_one_or_none()
async def get_user( ## Removed unused generic get_user helper to keep DAL explicit and simple
session: AsyncSession,
*,
user_id: Optional[int] = None,
username: Optional[str] = None,
panel_uuid: Optional[str] = None,
) -> Optional[User]:
if user_id is not None:
return await get_user_by_id(session, user_id)
if username is not None:
return await get_user_by_username(session, username)
if panel_uuid is not None:
return await get_user_by_panel_uuid(session, panel_uuid)
return None
async def create_user(session: AsyncSession, user_data: Dict[str, Any]) -> User: async def create_user(session: AsyncSession, user_data: Dict[str, Any]) -> User:
+32 -6
View File
@@ -16,6 +16,7 @@
"choose_language": "Choose language / Выберите язык:", "choose_language": "Choose language / Выберите язык:",
"language_set_alert": "Language changed!", "language_set_alert": "Language changed!",
"error_occurred_try_again": "An error occurred, please try again.", "error_occurred_try_again": "An error occurred, please try again.",
"error_try_again": "Please try again.",
"error_displaying_menu": "Error displaying menu.", "error_displaying_menu": "Error displaying menu.",
"main_menu_unknown_action": "Unknown action.", "main_menu_unknown_action": "Unknown action.",
@@ -122,6 +123,28 @@
"admin_stats_recent_payments_header": "Recent Payments:", "admin_stats_recent_payments_header": "Recent Payments:",
"admin_stats_payment_item": "{status_emoji} {amount} {currency} from {user_info} ({p_status}) [{p_date}]", "admin_stats_payment_item": "{status_emoji} {amount} {currency} from {user_info} ({p_status}) [{p_date}]",
"admin_stats_no_payments_found": "No payments found yet.", "admin_stats_no_payments_found": "No payments found yet.",
"admin_view_payments_button": "💰 Payments",
"admin_payments_header": "💰 <b>All Payments</b>",
"admin_no_payments_found": "No payments found.",
"admin_export_payments_csv": "📊 Export CSV",
"admin_refresh_payments": "🔄 Refresh",
"admin_no_payments_to_export": "No payments to export.",
"admin_payments_export_success": "📊 Payments export completed!\nTotal records: {count}",
"admin_export_sent": "File sent!",
"admin_csv_payment_id": "ID",
"admin_csv_user_id": "User ID",
"admin_csv_username": "Username",
"admin_csv_first_name": "First Name",
"admin_csv_amount": "Amount",
"admin_csv_currency": "Currency",
"admin_csv_provider": "Provider",
"admin_csv_status": "Status",
"admin_csv_description": "Description",
"admin_csv_months": "Months",
"admin_csv_created_at": "Created At",
"admin_csv_provider_payment_id": "Provider Payment ID",
"admin_stats_last_sync_header": "Last Panel Sync:", "admin_stats_last_sync_header": "Last Panel Sync:",
"admin_stats_sync_time": "Time", "admin_stats_sync_time": "Time",
"admin_stats_sync_status": "Status", "admin_stats_sync_status": "Status",
@@ -131,7 +154,7 @@
"admin_sync_status_never_run": "Panel sync never run.", "admin_sync_status_never_run": "Panel sync never run.",
"admin_broadcast_enter_message": "Enter the broadcast message (HTML supported):", "admin_broadcast_enter_message": "Enter the broadcast message (HTML supported):",
"admin_broadcast_confirm_prompt": "You are about to send the following message (first 200 characters):\n\n{message_preview}\n\nConfirm sending?", "admin_broadcast_confirm_prompt": "You are about to send the following message:\n\n{message_preview}\n\nConfirm sending?",
"confirm_broadcast_send_button": "✅ Send", "confirm_broadcast_send_button": "✅ Send",
"cancel_broadcast_button": "❌ Cancel", "cancel_broadcast_button": "❌ Cancel",
"admin_broadcast_sending_started": "Starting broadcast...", "admin_broadcast_sending_started": "Starting broadcast...",
@@ -145,6 +168,8 @@
"admin_promo_create_prompt": "Enter promo details in the format: CODE BONUS_DAYS MAX_USES [VALIDITY_DAYS]\nExample: <code>{example_format}</code>\n(Validity is optional; default is indefinite)", "admin_promo_create_prompt": "Enter promo details in the format: CODE BONUS_DAYS MAX_USES [VALIDITY_DAYS]\nExample: <code>{example_format}</code>\n(Validity is optional; default is indefinite)",
"admin_promo_invalid_format": "Invalid format. Please use: CODE BONUS_DAYS MAX_USES [VALIDITY_DAYS]", "admin_promo_invalid_format": "Invalid format. Please use: CODE BONUS_DAYS MAX_USES [VALIDITY_DAYS]",
"admin_promo_invalid_code_format": "Code must be 330 alphanumeric characters.", "admin_promo_invalid_code_format": "Code must be 330 alphanumeric characters.",
"admin_promo_invalid_bonus_days": "Bonus days must be a positive number.",
"admin_promo_invalid_max_activations": "Max activations must be a positive number.",
"admin_promo_invalid_bonus_or_activations": "Bonus days and max uses must be positive numbers.", "admin_promo_invalid_bonus_or_activations": "Bonus days and max uses must be positive numbers.",
"admin_promo_invalid_validity_days": "Validity period (in days) must be a positive number.", "admin_promo_invalid_validity_days": "Validity period (in days) must be a positive number.",
"admin_promo_invalid_values": "Invalid values. {error}", "admin_promo_invalid_values": "Invalid values. {error}",
@@ -226,11 +251,11 @@
"admin_log_user_not_found": "User \"{input}\" not found in bot database.", "admin_log_user_not_found": "User \"{input}\" not found in bot database.",
"admin_user_logs_title": "Logs for {user_display} (page {current_page}/{total_pages}):", "admin_user_logs_title": "Logs for {user_display} (page {current_page}/{total_pages}):",
"sync_started": "🔄 Starting data sync with panel...", "sync_started_simple": "🔄 Starting synchronization...",
"sync_failed": " Sync with panel failed. Details: {details}", "sync_success_simple": " Synchronization completed successfully",
"sync_completed": " Sync with panel completed. Status: {status}. Details: {details}", "sync_failed_simple": " Synchronization failed",
"sync_completed_details": "Checked: {total_checked} entries.\nUsers synced/updated: {users_synced}.\nSubscriptions synced/updated: {subs_synced}.", "sync_errors_simple": "⚠️ Synchronization completed with errors ({errors_count} errors)",
"sync_completed_with_errors_details": "Checked: {total_checked} entries.\nUsers synced/updated: {users_synced}.\nSubscriptions synced/updated: {subs_synced}.\nErrors: {errors_count}.\n\nFirst errors:\n{error_details_preview}", "sync_critical_error": "❌ Critical synchronization error",
"no_errors_placeholder": "none", "no_errors_placeholder": "none",
"admin_sync_initiated_from_panel": "Sync initiated...", "admin_sync_initiated_from_panel": "Sync initiated...",
"admin_panel_user_creation_failed": "❌ Failed to create panel user for TG ID {user_id}. Panel unreachable?", "admin_panel_user_creation_failed": "❌ Failed to create panel user for TG ID {user_id}. Panel unreachable?",
@@ -314,6 +339,7 @@
"log_payment_received": "{provider_emoji} <b>Payment Received</b>\n\n👤 User: {user_display}\n💰 Amount: <b>{amount} {currency}</b>\n📅 Period: <b>{months} mo.</b>\n🏦 Provider: {payment_provider}\n🕐 Time: {timestamp}", "log_payment_received": "{provider_emoji} <b>Payment Received</b>\n\n👤 User: {user_display}\n💰 Amount: <b>{amount} {currency}</b>\n📅 Period: <b>{months} mo.</b>\n🏦 Provider: {payment_provider}\n🕐 Time: {timestamp}",
"log_promo_activation": "🎁 <b>Promo Code Activated</b>\n\n👤 User: {user_display}\n🏷 Code: <code>{promo_code}</code>\n🎯 Bonus: <b>+{bonus_days}d</b>\n🕐 Time: {timestamp}", "log_promo_activation": "🎁 <b>Promo Code Activated</b>\n\n👤 User: {user_display}\n🏷 Code: <code>{promo_code}</code>\n🎯 Bonus: <b>+{bonus_days}d</b>\n🕐 Time: {timestamp}",
"log_trial_activation": "🆓 <b>Trial Activated</b>\n\n👤 User: {user_display}\n⏰ Valid until: <b>{end_date}</b>\n🕐 Time: {timestamp}", "log_trial_activation": "🆓 <b>Trial Activated</b>\n\n👤 User: {user_display}\n⏰ Valid until: <b>{end_date}</b>\n🕐 Time: {timestamp}",
"log_panel_sync": "{status_emoji} <b>Panel Synchronization</b>\n\n📊 Status: <b>{status}</b>\n👥 Users processed: <b>{users_processed}</b>\n📋 Subscriptions synced: <b>{subs_synced}</b>\n🕐 Time: {timestamp}\n\n📝 Details:\n{details}",
"log_suspicious_promo": "⚠️ <b>Suspicious Promo Code Attempt</b>\n\n👤 User: {user_display}\n🆔 ID: <code>{user_id}</code>\n📝 Input: <pre>{suspicious_input}</pre>\n🕐 Time: {timestamp}", "log_suspicious_promo": "⚠️ <b>Suspicious Promo Code Attempt</b>\n\n👤 User: {user_display}\n🆔 ID: <code>{user_id}</code>\n📝 Input: <pre>{suspicious_input}</pre>\n🕐 Time: {timestamp}",
"admin_general_cancel_operation": "Operation cancelled ❌", "admin_general_cancel_operation": "Operation cancelled ❌",
+42 -6
View File
@@ -16,6 +16,7 @@
"choose_language": "Выберите язык / Select language:", "choose_language": "Выберите язык / Select language:",
"language_set_alert": "Язык изменен!", "language_set_alert": "Язык изменен!",
"error_occurred_try_again": "Произошла ошибка, попробуйте снова.", "error_occurred_try_again": "Произошла ошибка, попробуйте снова.",
"error_try_again": "Попробуйте еще раз.",
"error_displaying_menu": "Ошибка отображения меню.", "error_displaying_menu": "Ошибка отображения меню.",
"main_menu_unknown_action": "Неизвестное действие.", "main_menu_unknown_action": "Неизвестное действие.",
@@ -122,6 +123,28 @@
"admin_stats_recent_payments_header": "Последние платежи:", "admin_stats_recent_payments_header": "Последние платежи:",
"admin_stats_payment_item": "{status_emoji} {amount} {currency} от {user_info} ({p_status}) [{p_date}]", "admin_stats_payment_item": "{status_emoji} {amount} {currency} от {user_info} ({p_status}) [{p_date}]",
"admin_stats_no_payments_found": "Платежей пока нет.", "admin_stats_no_payments_found": "Платежей пока нет.",
"admin_view_payments_button": "💰 Платежи",
"admin_payments_header": "💰 <b>Все платежи</b>",
"admin_no_payments_found": "Платежи не найдены.",
"admin_export_payments_csv": "📊 Экспорт CSV",
"admin_refresh_payments": "🔄 Обновить",
"admin_no_payments_to_export": "Нет платежей для экспорта.",
"admin_payments_export_success": "📊 Экспорт платежей завершен!\nВсего записей: {count}",
"admin_export_sent": "Файл отправлен!",
"admin_csv_payment_id": "ID",
"admin_csv_user_id": "User ID",
"admin_csv_username": "Логин",
"admin_csv_first_name": "Имя",
"admin_csv_amount": "Сумма",
"admin_csv_currency": "Валюта",
"admin_csv_provider": "Платежная система",
"admin_csv_status": "Статус",
"admin_csv_description": "Описание",
"admin_csv_months": "Месяцев",
"admin_csv_created_at": "Дата создания",
"admin_csv_provider_payment_id": "ID платежа в системе",
"admin_stats_last_sync_header": "Последняя синхронизация с панелью:", "admin_stats_last_sync_header": "Последняя синхронизация с панелью:",
"admin_stats_sync_time": "Время", "admin_stats_sync_time": "Время",
"admin_stats_sync_status": "Статус", "admin_stats_sync_status": "Статус",
@@ -131,7 +154,7 @@
"admin_sync_status_never_run": "Синхронизация с панелью еще не проводилась.", "admin_sync_status_never_run": "Синхронизация с панелью еще не проводилась.",
"admin_broadcast_enter_message": "Введите сообщение для рассылки (HTML поддерживается):", "admin_broadcast_enter_message": "Введите сообщение для рассылки (HTML поддерживается):",
"admin_broadcast_confirm_prompt": "Вы собираетесь отправить следующее сообщение (первые 200 символов):\n\n{message_preview}\n\nПодтверждаете отправку?", "admin_broadcast_confirm_prompt": "Вы собираетесь отправить следующее сообщение:\n\n{message_preview}\n\nПодтверждаете отправку?",
"confirm_broadcast_send_button": "✅ Отправить", "confirm_broadcast_send_button": "✅ Отправить",
"cancel_broadcast_button": "❌ Отмена", "cancel_broadcast_button": "❌ Отмена",
"admin_broadcast_sending_started": "Начинаю рассылку...", "admin_broadcast_sending_started": "Начинаю рассылку...",
@@ -145,11 +168,23 @@
"admin_promo_create_prompt": "Введите детали промокода в формате: КОД ДНИ_БОНУСА МАКС_АКТИВАЦИЙ [СРОК_ДЕЙСТВИЯ_В_ДНЯХ_ОТ_СЕЙЧАС]\nПример: <code>{example_format}</code>\n(Срок действия необязателен, по умолчанию - бессрочный)", "admin_promo_create_prompt": "Введите детали промокода в формате: КОД ДНИ_БОНУСА МАКС_АКТИВАЦИЙ [СРОК_ДЕЙСТВИЯ_В_ДНЯХ_ОТ_СЕЙЧАС]\nПример: <code>{example_format}</code>\n(Срок действия необязателен, по умолчанию - бессрочный)",
"admin_promo_invalid_format": "Неверный формат ввода. Пожалуйста, используйте: КОД ДНИ_БОНУСА МАКС_АКТИВАЦИЙ [ДНИ_ДЕЙСТВИЯ]", "admin_promo_invalid_format": "Неверный формат ввода. Пожалуйста, используйте: КОД ДНИ_БОНУСА МАКС_АКТИВАЦИЙ [ДНИ_ДЕЙСТВИЯ]",
"admin_promo_invalid_code_format": "Код должен быть от 3 до 30 символов и содержать только буквы и цифры.", "admin_promo_invalid_code_format": "Код должен быть от 3 до 30 символов и содержать только буквы и цифры.",
"admin_promo_invalid_bonus_days": "Количество бонусных дней должно быть положительным числом.",
"admin_promo_invalid_max_activations": "Максимальное количество активаций должно быть положительным числом.",
"admin_promo_invalid_bonus_or_activations": "Количество бонусных дней и максимальных активаций должны быть положительными числами.", "admin_promo_invalid_bonus_or_activations": "Количество бонусных дней и максимальных активаций должны быть положительными числами.",
"admin_promo_invalid_validity_days": "Срок действия промокода (в днях) должен быть положительным числом.", "admin_promo_invalid_validity_days": "Срок действия промокода (в днях) должен быть положительным числом.",
"admin_promo_invalid_values": "Неверные значения. {error}", "admin_promo_invalid_values": "Неверные значения. {error}",
"admin_promo_invalid_format_general": "Ошибка парсинга деталей промокода. Проверьте формат.", "admin_promo_invalid_format_general": "Ошибка парсинга деталей промокода. Проверьте формат.",
"admin_promo_created_success": "✅ Промокод <code>{code}</code> успешно создан!\nБонус: {bonus_days} дней\nМакс. активаций: {max_activations}\nДействителен: {valid_until_str}", "admin_promo_created_success": "✅ Промокод <code>{code}</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": "📊 Создано: <b>{created}</b> из <b>{total}</b>",
"admin_bulk_promo_settings": "🎁 Бонусные дни: <b>{bonus_days}</b>\n📊 Макс. активаций: <b>{max_activations}</b>\n⏰ Срок действия: <b>{validity}</b>",
"admin_promo_list_page_info": "Страница {current}/{total} ({count} промокодов)",
"admin_queue_status_button": "📊 Статус очередей",
"admin_queue_status_title": "📊 Статус очередей сообщений",
"admin_queue_status_info": "📤 <b>Очереди сообщений:</b>\n\n👥 <b>Пользователи (25 сообщ/сек):</b>\n 📋 В очереди: {user_queue_size}\n 🔄 Обрабатывается: {user_processing}\n 📈 Отправлено за минуту: {user_recent}\n\n📢 <b>Группы/каналы (15 сообщ/мин):</b>\n 📋 В очереди: {group_queue_size}\n 🔄 Обрабатывается: {group_processing}\n 📈 Отправлено за минуту: {group_recent}",
"admin_promo_creation_failed_duplicate": "❌ Ошибка: Промокод <code>{code}</code> уже существует.", "admin_promo_creation_failed_duplicate": "❌ Ошибка: Промокод <code>{code}</code> уже существует.",
"admin_promo_creation_failed": "❌ Не удалось создать промокод. Пожалуйста, попробуйте позже.", "admin_promo_creation_failed": "❌ Не удалось создать промокод. Пожалуйста, попробуйте позже.",
"admin_active_promos_list_header": "Активные промокоды:", "admin_active_promos_list_header": "Активные промокоды:",
@@ -226,11 +261,11 @@
"admin_log_user_not_found": "Пользователь по запросу \"{input}\" не найден в базе данных бота.", "admin_log_user_not_found": "Пользователь по запросу \"{input}\" не найден в базе данных бота.",
"admin_user_logs_title": "Логи пользователя {user_display} (стр. {current_page}/{total_pages}):", "admin_user_logs_title": "Логи пользователя {user_display} (стр. {current_page}/{total_pages}):",
"sync_started": "🔄 Начинаю синхронизацию данных с панелью...", "sync_started_simple": "🔄 Начинаю синхронизацию...",
"sync_failed": "❌ Ошибка синхронизации с панелью. Детали: {details}", "sync_success_simple": "Синхронизация успешно завершена",
"sync_completed": " Синхронизация с панелью завершена. Статус: {status}. Детали: {details}", "sync_failed_simple": " Синхронизация завершилась с ошибкой",
"sync_completed_details": "Проверено: {total_checked} записей.\nПользователей синхронизировано/обновлено: {users_synced}.\nПодписок синхронизировано/обновлено: {subs_synced}.", "sync_errors_simple": "⚠️ Синхронизация завершена с ошибками ({errors_count} ошибок)",
"sync_completed_with_errors_details": "Проверено: {total_checked} записей.\nПользователей синхронизировано/обновлено: {users_synced}.\nПодписок синхронизировано/обновлено: {subs_synced}.\nОшибок: {errors_count}.\n\nПервые ошибки:\n{error_details_preview}", "sync_critical_error": "❌ Критическая ошибка синхронизации",
"no_errors_placeholder": "нет", "no_errors_placeholder": "нет",
"admin_sync_initiated_from_panel": "Синхронизация запущена...", "admin_sync_initiated_from_panel": "Синхронизация запущена...",
"admin_panel_user_creation_failed": "❌ Не удалось создать пользователя на панели для TG ID {user_id}. Панель недоступна?", "admin_panel_user_creation_failed": "❌ Не удалось создать пользователя на панели для TG ID {user_id}. Панель недоступна?",
@@ -313,6 +348,7 @@
"log_payment_received": "{provider_emoji} <b>Получен платеж</b>\n\n👤 Пользователь: {user_display}\n💰 Сумма: <b>{amount} {currency}</b>\n📅 Период: <b>{months} мес.</b>\n🏦 Провайдер: {payment_provider}\n🕐 Время: {timestamp}", "log_payment_received": "{provider_emoji} <b>Получен платеж</b>\n\n👤 Пользователь: {user_display}\n💰 Сумма: <b>{amount} {currency}</b>\n📅 Период: <b>{months} мес.</b>\n🏦 Провайдер: {payment_provider}\n🕐 Время: {timestamp}",
"log_promo_activation": "🎁 <b>Активирован промокод</b>\n\n👤 Пользователь: {user_display}\n🏷 Код: <code>{promo_code}</code>\n🎯 Бонус: <b>+{bonus_days} дн.</b>\n🕐 Время: {timestamp}", "log_promo_activation": "🎁 <b>Активирован промокод</b>\n\n👤 Пользователь: {user_display}\n🏷 Код: <code>{promo_code}</code>\n🎯 Бонус: <b>+{bonus_days} дн.</b>\n🕐 Время: {timestamp}",
"log_trial_activation": "🆓 <b>Активирован триал</b>\n\n👤 Пользователь: {user_display}\n⏰ Действует до: <b>{end_date}</b>\n🕐 Время: {timestamp}", "log_trial_activation": "🆓 <b>Активирован триал</b>\n\n👤 Пользователь: {user_display}\n⏰ Действует до: <b>{end_date}</b>\n🕐 Время: {timestamp}",
"log_panel_sync": "{status_emoji} <b>Синхронизация с панелью</b>\n\n📊 Статус: <b>{status}</b>\n👥 Обработано пользователей: <b>{users_processed}</b>\n📋 Синхронизировано подписок: <b>{subs_synced}</b>\n🕐 Время: {timestamp}\n\n📝 Детали:\n{details}",
"log_suspicious_promo": "⚠️ <b>Подозрительная попытка ввода промокода</b>\n\n👤 Пользователь: {user_display}\n🆔 ID: <code>{user_id}</code>\n📝 Ввод: <pre>{suspicious_input}</pre>\n🕐 Время: {timestamp}", "log_suspicious_promo": "⚠️ <b>Подозрительная попытка ввода промокода</b>\n\n👤 Пользователь: {user_display}\n🆔 ID: <code>{user_id}</code>\n📝 Ввод: <pre>{suspicious_input}</pre>\n🕐 Время: {timestamp}",
"admin_general_cancel_operation": "Операция отменена ❌", "admin_general_cancel_operation": "Операция отменена ❌",