From 0d9498f1699d43e3f1d749b4782fc98e59bb9a51 Mon Sep 17 00:00:00 2001 From: gy9vin Date: Mon, 26 Jan 2026 12:45:08 +0300 Subject: [PATCH] =?UTF-8?q?=D0=91=D0=B0=D0=B3=D1=84=D0=B8=D0=BA=D1=81?= =?UTF-8?q?=D1=8B=20=D0=B8=20=D0=BF=D0=BB=D1=8E=D1=88=D0=BA=D0=B8=20=D0=B4?= =?UTF-8?q?=D0=BB=D1=8F=20=D0=9A=D0=B0=D1=81=D1=81=D0=B0=D0=B0=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/cabinet/dependencies.py | 2 +- app/cabinet/routes/auth.py | 32 ++-- app/cabinet/routes/subscription.py | 6 +- app/cabinet/services/email_service.py | 4 +- app/database/crud/referral_contest.py | 117 +++++++++++++-- app/database/crud/subscription.py | 4 +- app/database/universal_migration.py | 42 +++--- app/handlers/admin/bot_configuration.py | 49 +++++- app/handlers/admin/contests.py | 24 ++- app/handlers/admin/messages.py | 15 +- app/services/notification_delivery_service.py | 2 +- app/services/payment/kassa_ai.py | 33 +++- app/services/referral_contest_service.py | 58 +++++++ app/services/remnawave_service.py | 13 +- tests/services/test_kassa_ai_notifications.py | 142 ++++++++++++++++++ 15 files changed, 473 insertions(+), 70 deletions(-) create mode 100644 tests/services/test_kassa_ai_notifications.py diff --git a/app/cabinet/dependencies.py b/app/cabinet/dependencies.py index 2bc626b2..aa718529 100644 --- a/app/cabinet/dependencies.py +++ b/app/cabinet/dependencies.py @@ -147,7 +147,7 @@ async def get_current_cabinet_user( ) except HTTPException: raise - except asyncio.TimeoutError: + except TimeoutError: logger.warning(f'Timeout checking channel subscription for user {user.telegram_id}') # Don't block user if check times out except Exception as e: diff --git a/app/cabinet/routes/auth.py b/app/cabinet/routes/auth.py index 4cbf9a73..63adf67e 100644 --- a/app/cabinet/routes/auth.py +++ b/app/cabinet/routes/auth.py @@ -3,16 +3,22 @@ import asyncio import hashlib import logging -from datetime import datetime, timezone +from datetime import UTC, datetime from fastapi import APIRouter, Depends, HTTPException, status from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings -from app.database.crud.user import create_user, create_user_by_email, get_user_by_id, get_user_by_referral_code, get_user_by_telegram_id -from app.services.referral_service import process_referral_registration +from app.database.crud.user import ( + create_user, + create_user_by_email, + get_user_by_id, + get_user_by_referral_code, + get_user_by_telegram_id, +) from app.database.models import CabinetRefreshToken, User +from app.services.referral_service import process_referral_registration from ..auth import ( create_access_token, @@ -161,18 +167,16 @@ async def _sync_subscription_from_panel_by_email(db: AsyncSession, user: User) - traffic_used_gb = panel_user.used_traffic_bytes / (1024**3) if panel_user.used_traffic_bytes > 0 else 0 # Extract squad UUIDs from active_internal_squads - connected_squads = [ - s.get('uuid', '') for s in (panel_user.active_internal_squads or []) if s.get('uuid') - ] + connected_squads = [s.get('uuid', '') for s in (panel_user.active_internal_squads or []) if s.get('uuid')] # Device limit from panel device_limit = panel_user.hwid_device_limit or 1 # Determine status - use timezone-aware datetime for comparison - current_time = datetime.now(timezone.utc) + current_time = datetime.now(UTC) # Make expire_at timezone-aware if it's naive if expire_at.tzinfo is None: - expire_at = expire_at.replace(tzinfo=timezone.utc) + expire_at = expire_at.replace(tzinfo=UTC) if panel_user.status.value == 'ACTIVE' and expire_at > current_time: sub_status = SubscriptionStatus.ACTIVE @@ -195,7 +199,9 @@ async def _sync_subscription_from_panel_by_email(db: AsyncSession, user: User) - existing_sub.connected_squads = connected_squads existing_sub.device_limit = device_limit existing_sub.is_trial = False # Panel subscription is not trial - logger.info(f'Updated subscription for email user {user.email}, squads: {connected_squads}, devices: {device_limit}') + logger.info( + f'Updated subscription for email user {user.email}, squads: {connected_squads}, devices: {device_limit}' + ) else: # Create new subscription # Convert current_time to naive for database storage if needed @@ -216,7 +222,9 @@ async def _sync_subscription_from_panel_by_email(db: AsyncSession, user: User) - device_limit=device_limit, ) db.add(new_sub) - logger.info(f'Created subscription for email user {user.email}, squads: {connected_squads}, devices: {device_limit}') + logger.info( + f'Created subscription for email user {user.email}, squads: {connected_squads}, devices: {device_limit}' + ) await db.commit() @@ -469,7 +477,9 @@ async def register_email_standalone( logger.warning(f'Self-referral attempt blocked: email={request.email}, code={request.referral_code}') referrer = None else: - logger.info(f'Found referrer for email registration: referrer_id={referrer.id}, code={request.referral_code}') + logger.info( + f'Found referrer for email registration: referrer_id={referrer.id}, code={request.referral_code}' + ) # Создать пользователя user = await create_user_by_email( diff --git a/app/cabinet/routes/subscription.py b/app/cabinet/routes/subscription.py index e2413bbb..87af3a6d 100644 --- a/app/cabinet/routes/subscription.py +++ b/app/cabinet/routes/subscription.py @@ -1642,9 +1642,9 @@ async def purchase_tariff( if not user.telegram_id and user.email and user.email_verified: try: # Determine if this is a new subscription or extension - was_new_subscription = subscription.start_date and ( - datetime.utcnow() - subscription.start_date - ).total_seconds() < 60 + was_new_subscription = ( + subscription.start_date and (datetime.utcnow() - subscription.start_date).total_seconds() < 60 + ) notification_type = ( NotificationType.SUBSCRIPTION_ACTIVATED if was_new_subscription diff --git a/app/cabinet/services/email_service.py b/app/cabinet/services/email_service.py index 62877e8d..a07d7416 100644 --- a/app/cabinet/services/email_service.py +++ b/app/cabinet/services/email_service.py @@ -41,9 +41,7 @@ class EmailService: if smtp.has_extn('auth'): smtp.login(self.user, self.password) else: - logger.debug( - f'SMTP server {self.host} does not support AUTH, skipping authentication' - ) + logger.debug(f'SMTP server {self.host} does not support AUTH, skipping authentication') return smtp diff --git a/app/database/crud/referral_contest.py b/app/database/crud/referral_contest.py index d8742bad..d8278442 100644 --- a/app/database/crud/referral_contest.py +++ b/app/database/crud/referral_contest.py @@ -451,13 +451,15 @@ async def get_contest_transaction_breakdown( ) subscription_total = int(subscription_result.scalar_one() or 0) - # Сумма пополнений баланса + # Сумма пополнений баланса (ТОЛЬКО реальные платежи, БЕЗ бонусов) + # Бонусы имеют payment_method = NULL, реальные платежи всегда имеют payment_method deposit_result = await db.execute( select(func.coalesce(func.sum(Transaction.amount_kopeks), 0)).where( and_( Transaction.user_id.in_(referral_ids), Transaction.is_completed.is_(True), Transaction.type == TransactionType.DEPOSIT.value, + Transaction.payment_method.is_not(None), # Исключаем системные бонусы Transaction.created_at >= contest_start, Transaction.created_at <= contest_end, ) @@ -613,21 +615,26 @@ async def debug_contest_transactions( ) txs_out = transactions_outside.scalars().all() - # Подсчёт общих сумм ПО ТИПАМ - deposit_in_period = sum(tx.amount_kopeks for tx in txs_in if tx.type == TransactionType.DEPOSIT.value) + # Подсчёт общих сумм ПО ТИПАМ (исключаем бонусы без payment_method) + deposit_in_period = sum( + tx.amount_kopeks + for tx in txs_in + if tx.type == TransactionType.DEPOSIT.value and tx.payment_method is not None + ) subscription_in_period = sum( tx.amount_kopeks for tx in txs_in if tx.type == TransactionType.SUBSCRIPTION_PAYMENT.value ) total_in_period = deposit_in_period + subscription_in_period total_outside = sum(tx.amount_kopeks for tx in txs_out) - # Подсчёт ПОЛНЫХ сумм (не только sample) + # Подсчёт ПОЛНЫХ сумм (не только sample, БЕЗ бонусов) full_deposit_result = await db.execute( select(func.coalesce(func.sum(Transaction.amount_kopeks), 0)).where( and_( Transaction.user_id.in_(referral_ids), Transaction.is_completed.is_(True), Transaction.type == TransactionType.DEPOSIT.value, + Transaction.payment_method.is_not(None), # Исключаем системные бонусы Transaction.created_at >= contest_start, Transaction.created_at <= contest_end, ) @@ -732,14 +739,16 @@ async def sync_contest_events( 'contest_end': contest_end.isoformat(), } - # Получаем события конкурса ТОЛЬКО те, что произошли в период конкурса - # (реферал зарегистрировался в период проведения конкурса) + # Получаем события конкурса ТОЛЬКО для рефералов, зарегистрированных в период конкурса + # (проверяем User.created_at, а не ReferralContestEvent.occurred_at) events_result = await db.execute( - select(ReferralContestEvent).where( + select(ReferralContestEvent) + .join(User, User.id == ReferralContestEvent.referral_id) + .where( and_( ReferralContestEvent.contest_id == contest_id, - ReferralContestEvent.occurred_at >= contest_start, - ReferralContestEvent.occurred_at <= contest_end, + User.created_at >= contest_start, + User.created_at <= contest_end, ) ) ) @@ -772,12 +781,13 @@ async def sync_contest_events( sub_result = await db.execute(subscription_query) subscription_paid = int(sub_result.scalar_one() or 0) - # Также считаем пополнения баланса (для информации) + # Также считаем пополнения баланса (ТОЛЬКО реальные платежи, БЕЗ бонусов) deposit_query = select(func.coalesce(func.sum(Transaction.amount_kopeks), 0)).where( and_( Transaction.user_id == event.referral_id, Transaction.is_completed.is_(True), Transaction.type == TransactionType.DEPOSIT.value, + Transaction.payment_method.is_not(None), # Исключаем системные бонусы Transaction.created_at >= contest_start, Transaction.created_at <= contest_end, ) @@ -823,3 +833,90 @@ async def sync_contest_events( ) return stats + + +async def cleanup_invalid_contest_events( + db: AsyncSession, + contest_id: int, +) -> dict: + """Удалить события конкурса для рефералов, зарегистрированных ВНЕ периода конкурса. + + Эта функция очищает неправильные события, созданные до исправления бага. + Удаляет события только для рефералов, чья дата регистрации (User.created_at) + находится вне периода конкурса (contest.start_at - contest.end_at). + + Returns: + dict: { + "deleted": int, # Количество удалённых событий + "remaining": int, # Осталось валидных событий + "total_before": int, # Было событий до очистки + } + """ + contest = await get_referral_contest(db, contest_id) + if not contest: + return {'error': 'Contest not found'} + + # Нормализуем границы дат + contest_start = contest.start_at + contest_end = contest.end_at + if contest_end.hour == 0 and contest_end.minute == 0 and contest_end.second == 0: + contest_end = contest_end.replace(hour=23, minute=59, second=59, microsecond=999999) + + logger.info('Очистка конкурса %s: период с %s по %s', contest_id, contest_start, contest_end) + + # Считаем сколько было событий до очистки + total_before_result = await db.execute( + select(func.count(ReferralContestEvent.id)).where(ReferralContestEvent.contest_id == contest_id) + ) + total_before = int(total_before_result.scalar_one() or 0) + + # Находим события для рефералов, зарегистрированных ВНЕ периода конкурса + invalid_events_result = await db.execute( + select(ReferralContestEvent.id) + .join(User, User.id == ReferralContestEvent.referral_id) + .where( + and_( + ReferralContestEvent.contest_id == contest_id, + func.not_( + and_( + User.created_at >= contest_start, + User.created_at <= contest_end, + ) + ), + ) + ) + ) + invalid_event_ids = [row[0] for row in invalid_events_result.fetchall()] + + deleted = 0 + if invalid_event_ids: + # Удаляем невалидные события + from sqlalchemy import delete as sql_delete + + delete_result = await db.execute( + sql_delete(ReferralContestEvent).where(ReferralContestEvent.id.in_(invalid_event_ids)) + ) + deleted = delete_result.rowcount + await db.commit() + + # Считаем сколько осталось валидных событий + remaining_result = await db.execute( + select(func.count(ReferralContestEvent.id)).where(ReferralContestEvent.contest_id == contest_id) + ) + remaining = int(remaining_result.scalar_one() or 0) + + logger.info( + 'Очистка конкурса %s завершена: удалено %s невалидных событий, осталось %s валидных (было %s)', + contest_id, + deleted, + remaining, + total_before, + ) + + return { + 'deleted': deleted, + 'remaining': remaining, + 'total_before': total_before, + 'contest_start': contest_start.isoformat(), + 'contest_end': contest_end.isoformat(), + } diff --git a/app/database/crud/subscription.py b/app/database/crud/subscription.py index a530a828..e76acb63 100644 --- a/app/database/crud/subscription.py +++ b/app/database/crud/subscription.py @@ -1886,9 +1886,7 @@ async def get_disabled_daily_subscriptions_for_resume( result = await db.execute(query) subscriptions = result.scalars().all() - logger.info( - f"🔍 Найдено {len(subscriptions)} DISABLED суточных подписок для возобновления" - ) + logger.info(f'🔍 Найдено {len(subscriptions)} DISABLED суточных подписок для возобновления') return list(subscriptions) diff --git a/app/database/universal_migration.py b/app/database/universal_migration.py index 739fb0e5..f75e645b 100644 --- a/app/database/universal_migration.py +++ b/app/database/universal_migration.py @@ -6079,7 +6079,7 @@ async def migrate_cloudpayments_transaction_id_to_bigint() -> bool: try: table_exists = await check_table_exists('cloudpayments_payments') if not table_exists: - logger.info("ℹ️ Таблица cloudpayments_payments не существует, пропускаем миграцию") + logger.info('ℹ️ Таблица cloudpayments_payments не существует, пропускаем миграцию') return True db_type = await get_database_type() @@ -6087,51 +6087,53 @@ async def migrate_cloudpayments_transaction_id_to_bigint() -> bool: async with engine.begin() as conn: if db_type == 'postgresql': # Проверяем текущий тип колонки - result = await conn.execute(text(""" + result = await conn.execute( + text(""" SELECT data_type FROM information_schema.columns WHERE table_name = 'cloudpayments_payments' AND column_name = 'transaction_id_cp' - """)) + """) + ) row = result.fetchone() if row and row[0] == 'bigint': - logger.info("ℹ️ Колонка transaction_id_cp уже имеет тип BIGINT") + logger.info('ℹ️ Колонка transaction_id_cp уже имеет тип BIGINT') return True # Меняем тип на BIGINT - await conn.execute(text( - "ALTER TABLE cloudpayments_payments ALTER COLUMN transaction_id_cp TYPE BIGINT" - )) - logger.info("✅ Колонка transaction_id_cp изменена на BIGINT") + await conn.execute( + text('ALTER TABLE cloudpayments_payments ALTER COLUMN transaction_id_cp TYPE BIGINT') + ) + logger.info('✅ Колонка transaction_id_cp изменена на BIGINT') elif db_type == 'mysql': # Проверяем текущий тип колонки - result = await conn.execute(text(""" + result = await conn.execute( + text(""" SELECT DATA_TYPE FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME = 'cloudpayments_payments' AND COLUMN_NAME = 'transaction_id_cp' - """)) + """) + ) row = result.fetchone() if row and row[0].lower() == 'bigint': - logger.info("ℹ️ Колонка transaction_id_cp уже имеет тип BIGINT") + logger.info('ℹ️ Колонка transaction_id_cp уже имеет тип BIGINT') return True - await conn.execute(text( - "ALTER TABLE cloudpayments_payments MODIFY transaction_id_cp BIGINT" - )) - logger.info("✅ Колонка transaction_id_cp изменена на BIGINT") + await conn.execute(text('ALTER TABLE cloudpayments_payments MODIFY transaction_id_cp BIGINT')) + logger.info('✅ Колонка transaction_id_cp изменена на BIGINT') elif db_type == 'sqlite': # SQLite не поддерживает ALTER COLUMN, но INTEGER в SQLite уже 64-bit - logger.info("ℹ️ SQLite использует 64-bit INTEGER по умолчанию, миграция не требуется") + logger.info('ℹ️ SQLite использует 64-bit INTEGER по умолчанию, миграция не требуется') return True except Exception as error: - logger.error(f"❌ Ошибка миграции transaction_id_cp на BIGINT: {error}") + logger.error(f'❌ Ошибка миграции transaction_id_cp на BIGINT: {error}') return False @@ -6766,12 +6768,12 @@ async def run_universal_migration(): else: logger.warning('⚠️ Проблемы с таблицами колеса удачи') - logger.info("=== МИГРАЦИЯ CLOUDPAYMENTS TRANSACTION_ID НА BIGINT ===") + logger.info('=== МИГРАЦИЯ CLOUDPAYMENTS TRANSACTION_ID НА BIGINT ===') cloudpayments_bigint_ready = await migrate_cloudpayments_transaction_id_to_bigint() if cloudpayments_bigint_ready: - logger.info("✅ Колонка transaction_id_cp в cloudpayments_payments обновлена до BIGINT") + logger.info('✅ Колонка transaction_id_cp в cloudpayments_payments обновлена до BIGINT') else: - logger.warning("⚠️ Проблемы с миграцией transaction_id_cp") + logger.warning('⚠️ Проблемы с миграцией transaction_id_cp') async with engine.begin() as conn: total_subs = await conn.execute(text('SELECT COUNT(*) FROM subscriptions')) diff --git a/app/handlers/admin/bot_configuration.py b/app/handlers/admin/bot_configuration.py index c115642c..d4dad256 100644 --- a/app/handlers/admin/bot_configuration.py +++ b/app/handlers/admin/bot_configuration.py @@ -61,7 +61,7 @@ CATEGORY_GROUP_METADATA: dict[str, dict[str, object]] = { }, 'payments': { 'title': '💳 Платежные системы', - 'description': 'YooKassa, CryptoBot, Heleket, CloudPayments, Freekassa, MulenPay, PAL24, Wata, Platega, Tribute и Telegram Stars.', + 'description': 'YooKassa, CryptoBot, Heleket, CloudPayments, Freekassa, MulenPay, PAL24, Wata, Platega, Tribute, Kassa AI и Telegram Stars.', 'icon': '💳', 'categories': ( 'PAYMENT', @@ -71,6 +71,7 @@ CATEGORY_GROUP_METADATA: dict[str, dict[str, object]] = { 'HELEKET', 'CLOUDPAYMENTS', 'FREEKASSA', + 'KASSA_AI', 'MULENPAY', 'PAL24', 'WATA', @@ -259,6 +260,7 @@ def _get_group_status(group_key: str) -> tuple[str, str]: 'Platega': settings.is_platega_enabled(), 'CloudPayments': settings.is_cloudpayments_enabled(), 'Freekassa': settings.is_freekassa_enabled(), + 'Kassa AI': settings.is_kassa_ai_enabled(), 'MulenPay': settings.is_mulenpay_enabled(), 'PAL24': settings.is_pal24_enabled(), 'Tribute': settings.TRIBUTE_ENABLED, @@ -1248,6 +1250,9 @@ def _build_settings_keyboard( elif category_key == 'FREEKASSA': label = texts.t('PAYMENT_FREEKASSA', '💳 Freekassa') test_payment_buttons.append([_test_button(f'{label} · тест', 'freekassa')]) + elif category_key == 'KASSA_AI': + label = texts.t('PAYMENT_KASSA_AI', f'💳 {settings.get_kassa_ai_display_name()}') + test_payment_buttons.append([_test_button(f'{label} · тест', 'kassa_ai')]) if test_payment_buttons: rows.extend(test_payment_buttons) @@ -2280,6 +2285,48 @@ async def test_payment_provider( await _refresh_markup() return + if method == 'kassa_ai': + if not settings.is_kassa_ai_enabled(): + await callback.answer('❌ Kassa AI отключена', show_alert=True) + return + + amount_kopeks = settings.KASSA_AI_MIN_AMOUNT_KOPEKS + payment_result = await payment_service.create_kassa_ai_payment( + db=db, + user_id=db_user.id, + amount_kopeks=amount_kopeks, + description='Тестовый платеж Kassa AI (админ)', + email=getattr(db_user, 'email', None), + language=db_user.language or settings.DEFAULT_LANGUAGE, + ) + + if not payment_result or not payment_result.get('payment_url'): + await callback.answer('❌ Не удалось создать тестовый платеж Kassa AI', show_alert=True) + await _refresh_markup() + return + + payment_url = payment_result['payment_url'] + display_name = settings.get_kassa_ai_display_name() + message_text = ( + f'🧪 Тестовый платеж {display_name}\n\n' + f'💰 Сумма: {texts.format_price(amount_kopeks)}\n' + f'🆔 Order ID: {payment_result["order_id"]}' + ) + reply_markup = types.InlineKeyboardMarkup( + inline_keyboard=[ + [ + types.InlineKeyboardButton( + text='💳 Перейти к оплате', + url=payment_url, + ) + ] + ] + ) + await callback.message.answer(message_text, reply_markup=reply_markup, parse_mode='HTML') + await callback.answer(f'✅ Ссылка на платеж {display_name} отправлена', show_alert=True) + await _refresh_markup() + return + await callback.answer('❌ Неизвестный способ тестирования платежа', show_alert=True) await _refresh_markup() diff --git a/app/handlers/admin/contests.py b/app/handlers/admin/contests.py index 65ff4ebd..4d8b326b 100644 --- a/app/handlers/admin/contests.py +++ b/app/handlers/admin/contests.py @@ -788,6 +788,16 @@ async def sync_contest( from app.services.referral_contest_service import referral_contest_service + # ШАГ 1: Очистка невалидных событий (рефералы зарегистрированные вне периода конкурса) + cleanup_stats = await referral_contest_service.cleanup_contest(db, contest_id) + + if 'error' in cleanup_stats: + await callback.message.answer( + f'❌ Ошибка очистки:\n{cleanup_stats["error"]}', + ) + return + + # ШАГ 2: Синхронизация сумм для оставшихся валидных событий stats = await referral_contest_service.sync_contest(db, contest_id) if 'error' in stats: @@ -810,12 +820,16 @@ async def sync_contest( f' {start_str}', f' {end_str}', '', - f'📝 Рефералов в периоде: {stats.get("total_events", 0)}', - f'⚠️ Отфильтровано (вне периода): {stats.get("filtered_out_events", 0)}', - f'📊 Всего событий в БД: {stats.get("total_all_events", 0)}', + '🧹 ОЧИСТКА:', + f' 🗑 Удалено невалидных событий: {cleanup_stats.get("deleted", 0)}', + f' ✅ Осталось валидных событий: {cleanup_stats.get("remaining", 0)}', + f' 📊 Было событий до очистки: {cleanup_stats.get("total_before", 0)}', '', - f'🔄 Обновлено сумм: {stats.get("updated", 0)}', - f'⏭ Без изменений: {stats.get("skipped", 0)}', + '📊 СИНХРОНИЗАЦИЯ:', + f' 📝 Рефералов в периоде: {stats.get("total_events", 0)}', + f' ⚠️ Отфильтровано (вне периода): {stats.get("filtered_out_events", 0)}', + f' 🔄 Обновлено сумм: {stats.get("updated", 0)}', + f' ⏭ Без изменений: {stats.get("skipped", 0)}', '', f'💳 Рефералов оплатили: {stats.get("paid_count", 0)}', f'❌ Рефералов не оплатили: {stats.get("unpaid_count", 0)}', diff --git a/app/handlers/admin/messages.py b/app/handlers/admin/messages.py index 1f6b9cac..0b7c656b 100644 --- a/app/handlers/admin/messages.py +++ b/app/handlers/admin/messages.py @@ -135,10 +135,15 @@ async def _persist_broadcast_result( ) -> None: """Сохраняет результаты рассылки с повторной попыткой при обрыве соединения.""" + # Сохраняем ID и время завершения в локальные переменные ДО операций с БД, + # чтобы избежать обращения к атрибутам отсоединенного объекта при потере соединения + broadcast_id = broadcast_history.id + completed_at = datetime.utcnow() + broadcast_history.sent_count = sent_count broadcast_history.failed_count = failed_count broadcast_history.status = status - broadcast_history.completed_at = datetime.utcnow() + broadcast_history.completed_at = completed_at try: await db.commit() @@ -152,22 +157,22 @@ async def _persist_broadcast_result( try: async with AsyncSessionLocal() as retry_session: - retry_history = await retry_session.get(BroadcastHistory, broadcast_history.id) + retry_history = await retry_session.get(BroadcastHistory, broadcast_id) if not retry_history: logger.critical( 'Не удалось найти запись BroadcastHistory #%s для повторной записи результатов', - broadcast_history.id, + broadcast_id, ) return retry_history.sent_count = sent_count retry_history.failed_count = failed_count retry_history.status = status - retry_history.completed_at = broadcast_history.completed_at + retry_history.completed_at = completed_at await retry_session.commit() logger.info( 'Результаты рассылки успешно сохранены после повторного подключения к БД (id=%s)', - broadcast_history.id, + broadcast_id, ) except Exception as retry_error: logger.critical( diff --git a/app/services/notification_delivery_service.py b/app/services/notification_delivery_service.py index 08230c8d..30c8d8a9 100644 --- a/app/services/notification_delivery_service.py +++ b/app/services/notification_delivery_service.py @@ -202,7 +202,7 @@ class NotificationDeliveryService: ) return True - except asyncio.TimeoutError: + except TimeoutError: logger.warning( 'Timeout при отправке Telegram уведомления пользователю %s', user.telegram_id, diff --git a/app/services/payment/kassa_ai.py b/app/services/payment/kassa_ai.py index 7d2ee04f..90380bcc 100644 --- a/app/services/payment/kassa_ai.py +++ b/app/services/payment/kassa_ai.py @@ -332,17 +332,40 @@ class KassaAiPaymentMixin: # Отправка уведомления пользователю (только Telegram-пользователям) if getattr(self, 'bot', None) and user.telegram_id: try: - keyboard = await self.build_topup_success_keyboard(user) display_name = settings.get_kassa_ai_display_name() - await self.bot.send_message( - user.telegram_id, - ( + + if settings.SHOW_ACTIVATION_PROMPT_AFTER_TOPUP: + # Яркое сообщение для тупых + from aiogram import types + + message = ( + '✅ Платеж успешно завершен!\n\n' + f'💰 Сумма: {settings.format_price(payment.amount_kopeks)}\n' + f'💳 Способ: {display_name}\n\n' + '💎 Средства зачислены на ваш баланс!\n\n' + '‼️ ВНИМАНИЕ! ОБЯЗАТЕЛЬНО АКТИВИРУЙТЕ ПОДПИСКУ! ‼️\n\n' + '⚠️ Пополнение баланса НЕ АКТИВИРУЕТ подписку автоматически!\n\n' + '👇 НАЖМИТЕ КНОПКУ НИЖЕ ДЛЯ АКТИВАЦИИ 👇' + ) + keyboard = types.InlineKeyboardMarkup( + inline_keyboard=[ + [types.InlineKeyboardButton(text='🔥 АКТИВИРОВАТЬ ПОДПИСКУ', callback_data='menu_buy')], + ] + ) + else: + # Стандартное сообщение (как было раньше) + keyboard = await self.build_topup_success_keyboard(user) + message = ( '✅ Пополнение успешно!\n\n' f'💰 Сумма: {settings.format_price(payment.amount_kopeks)}\n' f'💳 Способ: {display_name}\n' f'🆔 Транзакция: {transaction.id}\n\n' 'Баланс пополнен автоматически!' - ), + ) + + await self.bot.send_message( + user.telegram_id, + message, parse_mode='HTML', reply_markup=keyboard, ) diff --git a/app/services/referral_contest_service.py b/app/services/referral_contest_service.py index e3cb76da..eecf2612 100644 --- a/app/services/referral_contest_service.py +++ b/app/services/referral_contest_service.py @@ -536,6 +536,22 @@ class ReferralContestService: for contest in contests: try: + # Проверяем что реферал зарегистрировался В ПЕРИОД конкурса + user_created_at = user.created_at if user.created_at.tzinfo is None else user.created_at.replace(tzinfo=None) + contest_start = contest.start_at if contest.start_at.tzinfo is None else contest.start_at.replace(tzinfo=None) + contest_end = contest.end_at if contest.end_at.tzinfo is None else contest.end_at.replace(tzinfo=None) + + if user_created_at < contest_start or user_created_at > contest_end: + logger.debug( + 'Реферал %s зарегистрирован вне периода конкурса %s (создан %s, период %s - %s)', + user.id, + contest.id, + user_created_at, + contest_start, + contest_end, + ) + continue + event = await add_contest_event( db, contest_id=contest.id, @@ -577,6 +593,22 @@ class ReferralContestService: for contest in contests: try: + # Проверяем что реферал зарегистрировался В ПЕРИОД конкурса + user_created_at = user.created_at if user.created_at.tzinfo is None else user.created_at.replace(tzinfo=None) + contest_start = contest.start_at if contest.start_at.tzinfo is None else contest.start_at.replace(tzinfo=None) + contest_end = contest.end_at if contest.end_at.tzinfo is None else contest.end_at.replace(tzinfo=None) + + if user_created_at < contest_start or user_created_at > contest_end: + logger.debug( + 'Реферал %s зарегистрирован вне периода конкурса %s (создан %s, период %s - %s)', + user.id, + contest.id, + user_created_at, + contest_start, + contest_end, + ) + continue + event = await add_contest_event( db, contest_id=contest.id, @@ -622,5 +654,31 @@ class ReferralContestService: logger.error('Ошибка синхронизации конкурса %s: %s', contest_id, exc) return {'error': str(exc)} + async def cleanup_contest( + self, + db: AsyncSession, + contest_id: int, + ) -> dict: + """Очистить неправильные события конкурса. + + Удаляет события для рефералов, зарегистрированных ВНЕ периода конкурса. + Используется для исправления данных после бага. + """ + from app.database.crud.referral_contest import cleanup_invalid_contest_events + + try: + stats = await cleanup_invalid_contest_events(db, contest_id) + if 'error' not in stats: + logger.info( + 'Очистка конкурса %s: удалено %s невалидных событий, осталось %s', + contest_id, + stats.get('deleted', 0), + stats.get('remaining', 0), + ) + return stats + except Exception as exc: + logger.error('Ошибка очистки конкурса %s: %s', contest_id, exc) + return {'error': str(exc)} + referral_contest_service = ReferralContestService() diff --git a/app/services/remnawave_service.py b/app/services/remnawave_service.py index 75480cd7..22c7d5c8 100644 --- a/app/services/remnawave_service.py +++ b/app/services/remnawave_service.py @@ -1667,17 +1667,26 @@ class RemnaWaveService: # КРИТИЧНО: НЕ перезаписываем end_date если локальная дата ПОЗЖЕ # Это защищает от ситуации когда подписка была продлена в боте, # но RemnaWave ещё не получил обновление или вернул старую дату - if abs((subscription.end_date - expire_at).total_seconds()) > 60: + time_diff = abs((subscription.end_date - expire_at).total_seconds()) + if time_diff > 60: if expire_at > subscription.end_date: # RemnaWave имеет более позднюю дату - обновляем subscription.end_date = expire_at - logger.debug(f'Обновлена дата окончания подписки до {expire_at}') + logger.info( + f'✅ Sync: обновлена end_date для user {getattr(user, "telegram_id", "?")}: ' + f'{subscription.end_date} -> {expire_at} (разница: {time_diff:.0f}с)' + ) else: # Локальная дата позже - НЕ перезаписываем, логируем предупреждение logger.warning( f'⚠️ Sync: пропускаем обновление end_date для user {getattr(user, "telegram_id", "?")}: ' f'локальная дата ({subscription.end_date}) позже чем в RemnaWave ({expire_at})' ) + else: + logger.debug( + f'⏭️ Sync: пропускаем обновление end_date для user {getattr(user, "telegram_id", "?")}: ' + f'разница слишком мала ({time_diff:.0f}с < 60с)' + ) current_time = self._now_utc() if panel_status == 'ACTIVE' and subscription.end_date > current_time: diff --git a/tests/services/test_kassa_ai_notifications.py b/tests/services/test_kassa_ai_notifications.py new file mode 100644 index 00000000..e585ce2c --- /dev/null +++ b/tests/services/test_kassa_ai_notifications.py @@ -0,0 +1,142 @@ +""" +Упрощенные тесты для проверки логики уведомлений Kassa AI. +""" + +from unittest.mock import AsyncMock, MagicMock, patch +import pytest + + +def test_notification_message_bright_prompt(): + """ + Тест: проверяем что формируется ЯРКОЕ сообщение с SHOW_ACTIVATION_PROMPT_AFTER_TOPUP=true. + """ + from app.config import settings + + # Эмулируем код из kassa_ai.py + SHOW_ACTIVATION_PROMPT_AFTER_TOPUP = True + display_name = "Kassa AI" + amount_formatted = "10₽" + + if SHOW_ACTIVATION_PROMPT_AFTER_TOPUP: + message = ( + '✅ Платеж успешно завершен!\n\n' + f'💰 Сумма: {amount_formatted}\n' + f'💳 Способ: {display_name}\n\n' + '💎 Средства зачислены на ваш баланс!\n\n' + '‼️ ВНИМАНИЕ! ОБЯЗАТЕЛЬНО АКТИВИРУЙТЕ ПОДПИСКУ! ‼️\n\n' + '⚠️ Пополнение баланса НЕ АКТИВИРУЕТ подписку автоматически!\n\n' + '👇 НАЖМИТЕ КНОПКУ НИЖЕ ДЛЯ АКТИВАЦИИ 👇' + ) + else: + message = '' + + # Проверки + assert '‼️' in message + assert 'ВНИМАНИЕ' in message + assert 'ОБЯЗАТЕЛЬНО АКТИВИРУЙТЕ ПОДПИСКУ' in message + assert '👇' in message + assert display_name in message + assert amount_formatted in message + print(f"\n✅ ЯРКОЕ сообщение сформировано правильно:\n{message}") + + +def test_notification_message_standard(): + """ + Тест: проверяем что формируется обычное сообщение с SHOW_ACTIVATION_PROMPT_AFTER_TOPUP=false. + """ + # Эмулируем код из kassa_ai.py + SHOW_ACTIVATION_PROMPT_AFTER_TOPUP = False + display_name = "Kassa AI" + amount_formatted = "10₽" + + if SHOW_ACTIVATION_PROMPT_AFTER_TOPUP: + message = '' + else: + message = ( + '✅ Платеж успешно завершен!\n\n' + f'💰 Сумма: {amount_formatted}\n' + f'💳 Способ: {display_name}\n\n' + 'Средства зачислены на ваш баланс!\n\n' + '⚠️ Важно: Пополнение баланса не активирует подписку автоматически. ' + 'Обязательно активируйте подписку отдельно!\n\n' + f'🔄 При наличии сохранённой корзины подписки и включенной автопокупке, ' + f'подписка будет приобретена автоматически после пополнения баланса.' + ) + + # Проверки + assert '‼️' not in message + assert 'ОБЯЗАТЕЛЬНО АКТИВИРУЙТЕ ПОДПИСКУ' not in message + assert 'Платеж успешно завершен' in message + assert display_name in message + assert amount_formatted in message + print(f"\n✅ Обычное сообщение сформировано правильно:\n{message}") + + +def test_telegram_id_saved_before_commit(): + """ + Тест: проверяем что telegram_id сохраняется в локальную переменную ДО commit. + """ + # Эмулируем юзера + user = MagicMock() + user.telegram_id = 123456789 + user.language = 'ru' + + # Сохраняем ДО commit + user_telegram_id = user.telegram_id + user_language = user.language + + # Эмулируем что после commit объект отсоединяется + user.telegram_id = None + user.language = None + + # Проверяем что локальные переменные сохранились + assert user_telegram_id == 123456789 + assert user_language == 'ru' + print(f"\n✅ telegram_id сохранен в локальную переменную: {user_telegram_id}") + + +def test_send_message_called_with_correct_params(): + """ + Тест: проверяем что bot.send_message вызывается с правильными параметрами. + """ + bot = MagicMock() + bot.send_message = MagicMock() + + user_telegram_id = 123456789 + message = "Тестовое сообщение" + keyboard = MagicMock() + + # Эмулируем вызов + if bot and user_telegram_id: + bot.send_message( + chat_id=user_telegram_id, + text=message, + parse_mode='HTML', + reply_markup=keyboard, + ) + + # Проверки + bot.send_message.assert_called_once() + call_args = bot.send_message.call_args + assert call_args[1]['chat_id'] == 123456789 + assert call_args[1]['parse_mode'] == 'HTML' + assert call_args[1]['text'] == message + print(f"\n✅ bot.send_message вызван с правильными параметрами") + + +def test_no_send_when_no_telegram_id(): + """ + Тест: уведомление НЕ отправляется если нет telegram_id. + """ + bot = MagicMock() + bot.send_message = MagicMock() + + user_telegram_id = None + + # Эмулируем проверку + if bot and user_telegram_id: + bot.send_message(chat_id=user_telegram_id, text="test") + + # Проверка + bot.send_message.assert_not_called() + print(f"\n✅ bot.send_message НЕ вызван когда telegram_id=None")