Merge pull request #2421 from BEDOLAGA-DEV/main

w
This commit is contained in:
Egor
2026-01-26 19:07:15 +03:00
committed by GitHub
19 changed files with 519 additions and 81 deletions
+27
View File
@@ -0,0 +1,27 @@
name: Lint
on:
push:
branches: ['**']
pull_request:
branches: ['**']
jobs:
lint:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: astral-sh/setup-uv@v5
- uses: actions/setup-python@v5
with:
python-version: '3.13'
- run: uv sync --group dev
- name: Check formatting
run: uv run ruff format --check .
- name: Check linting
run: uv run ruff check .
+7
View File
@@ -31,6 +31,13 @@ docker-compose.override.yml
# Разрешаем .gitignore чтобы он попал в репозиторий
!.gitignore
# Разрешаем .github/ (workflows, pre-commit и т.д.)
!.github/
!.github/**
# Разрешаем Makefile
!Makefile
# Внутри разрешенных папок игнорируем служебные файлы
app/__pycache__/
app/**/__pycache__/
+1 -3
View File
@@ -45,7 +45,5 @@ help: ## Показать список доступных команд
@echo ""
@echo "📘 Команды Makefile:"
@echo ""
@grep -E '^[a-zA-Z0-9_-]+:.*?##' $(MAKEFILE_LIST) | \
sed -E 's/:.*?## /| /' | \
awk -F'|' '{printf " \033[36m%-16s\033[0m %s\n", $$1, $$2}'
@awk -F':.*## ' '/^[a-zA-Z0-9_-]+:.*## / {printf " \033[36m%-16s\033[0m %s\n", $$1, $$2}' $(MAKEFILE_LIST)
@echo ""
+1 -1
View File
@@ -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:
+21 -11
View File
@@ -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(
+3 -3
View File
@@ -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
+1 -3
View File
@@ -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
+107 -10
View File
@@ -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(),
}
+1 -3
View File
@@ -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)
+22 -20
View File
@@ -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'))
+5 -1
View File
@@ -341,7 +341,11 @@ class RemnaWaveAPI:
connector = aiohttp.TCPConnector(**connector_kwargs)
session_kwargs = {'timeout': aiohttp.ClientTimeout(total=60, connect=10), 'headers': headers, 'connector': connector}
session_kwargs = {
'timeout': aiohttp.ClientTimeout(total=60, connect=10),
'headers': headers,
'connector': connector,
}
if cookies:
session_kwargs['cookies'] = cookies
+48 -1
View File
@@ -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'🧪 <b>Тестовый платеж {display_name}</b>\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()
+19 -5
View File
@@ -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' <code>{start_str}</code>',
f' <code>{end_str}</code>',
'',
f'📝 Рефералов в периоде: <b>{stats.get("total_events", 0)}</b>',
f'⚠️ Отфильтровано (вне периода): <b>{stats.get("filtered_out_events", 0)}</b>',
f'📊 Всего событий в БД: <b>{stats.get("total_all_events", 0)}</b>',
'🧹 <b>ОЧИСТКА:</b>',
f' 🗑 Удалено невалидных событий: <b>{cleanup_stats.get("deleted", 0)}</b>',
f' ✅ Осталось валидных событий: <b>{cleanup_stats.get("remaining", 0)}</b>',
f' 📊 Было событий до очистки: <b>{cleanup_stats.get("total_before", 0)}</b>',
'',
f'🔄 Обновлено сумм: <b>{stats.get("updated", 0)}</b>',
f'⏭ Без изменений: <b>{stats.get("skipped", 0)}</b>',
'📊 <b>СИНХРОНИЗАЦИЯ:</b>',
f' 📝 Рефералов в периоде: <b>{stats.get("total_events", 0)}</b>',
f' ⚠️ Отфильтровано (вне периода): <b>{stats.get("filtered_out_events", 0)}</b>',
f' 🔄 Обновлено сумм: <b>{stats.get("updated", 0)}</b>',
f' ⏭ Без изменений: <b>{stats.get("skipped", 0)}</b>',
'',
f'💳 Рефералов оплатили: <b>{stats.get("paid_count", 0)}</b>',
f'❌ Рефералов не оплатили: <b>{stats.get("unpaid_count", 0)}</b>',
+10 -5
View File
@@ -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(
@@ -202,7 +202,7 @@ class NotificationDeliveryService:
)
return True
except asyncio.TimeoutError:
except TimeoutError:
logger.warning(
'Timeout при отправке Telegram уведомления пользователю %s',
user.telegram_id,
+28 -5
View File
@@ -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 = (
'✅ <b>Платеж успешно завершен!</b>\n\n'
f'💰 Сумма: {settings.format_price(payment.amount_kopeks)}\n'
f'💳 Способ: {display_name}\n\n'
'💎 Средства зачислены на ваш баланс!\n\n'
'‼️ <b>ВНИМАНИЕ! ОБЯЗАТЕЛЬНО АКТИВИРУЙТЕ ПОДПИСКУ!</b> ‼️\n\n'
'⚠️ Пополнение баланса <b>НЕ АКТИВИРУЕТ</b> подписку автоматически!\n\n'
'👇 <b>НАЖМИТЕ КНОПКУ НИЖЕ ДЛЯ АКТИВАЦИИ</b> 👇'
)
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='🔥 АКТИВИРОВАТЬ ПОДПИСКУ', callback_data='menu_buy')],
]
)
else:
# Стандартное сообщение (как было раньше)
keyboard = await self.build_topup_success_keyboard(user)
message = (
'✅ <b>Пополнение успешно!</b>\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,
)
+58
View File
@@ -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()
+17 -9
View File
@@ -1174,10 +1174,7 @@ class RemnaWaveService:
user.remnawave_uuid: user for user in bot_users if getattr(user, 'remnawave_uuid', None)
}
# Index users by email for email-only sync
bot_users_by_email = {
user.email.lower(): user for user in bot_users
if user.email and user.email_verified
}
bot_users_by_email = {user.email.lower(): user for user in bot_users if user.email and user.email_verified}
# Also index email-only users by their remnawave_uuid for sync
email_users_count = sum(1 for u in bot_users if u.telegram_id is None)
if email_users_count > 0:
@@ -1203,8 +1200,7 @@ class RemnaWaveService:
# Email-only пользователи из панели (без telegram_id, но с email)
panel_users_email_only = [
user for user in panel_users
if user.get('telegramId') is None and user.get('email')
user for user in panel_users if user.get('telegramId') is None and user.get('email')
]
if panel_users_email_only:
logger.info(f'📧 Пользователей в панели с Email (без Telegram): {len(panel_users_email_only)}')
@@ -1667,17 +1663,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:
@@ -1834,7 +1839,10 @@ class RemnaWaveService:
telegram_id=user.telegram_id,
email=user.email,
description=settings.format_remnawave_user_description(
full_name=user.full_name, username=user.username, telegram_id=user.telegram_id, email=user.email
full_name=user.full_name,
username=user.username,
telegram_id=user.telegram_id,
email=user.email,
),
active_internal_squads=sub.connected_squads,
)
@@ -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 = (
'✅ <b>Платеж успешно завершен!</b>\n\n'
f'💰 Сумма: {amount_formatted}\n'
f'💳 Способ: {display_name}\n\n'
'💎 Средства зачислены на ваш баланс!\n\n'
'‼️ <b>ВНИМАНИЕ! ОБЯЗАТЕЛЬНО АКТИВИРУЙТЕ ПОДПИСКУ!</b> ‼️\n\n'
'⚠️ Пополнение баланса <b>НЕ АКТИВИРУЕТ</b> подписку автоматически!\n\n'
'👇 <b>НАЖМИТЕ КНОПКУ НИЖЕ ДЛЯ АКТИВАЦИИ</b> 👇'
)
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 = (
'✅ <b>Платеж успешно завершен!</b>\n\n'
f'💰 Сумма: {amount_formatted}\n'
f'💳 Способ: {display_name}\n\n'
'Средства зачислены на ваш баланс!\n\n'
'⚠️ <b>Важно:</b> Пополнение баланса не активирует подписку автоматически. '
'Обязательно активируйте подписку отдельно!\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")