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")