Merge pull request #2481 from BEDOLAGA-DEV/main

w
This commit is contained in:
Egor
2026-02-01 16:50:38 +03:00
committed by GitHub
16 changed files with 2881 additions and 267 deletions
+169 -1
View File
@@ -1015,6 +1015,10 @@ async def show_virtual_participants(
text=' Добавить',
callback_data=f'admin_contest_vp_add_{contest_id}',
),
types.InlineKeyboardButton(
text='🎭 Массовка',
callback_data=f'admin_contest_vp_mass_{contest_id}',
),
],
]
if vps:
@@ -1160,7 +1164,10 @@ async def delete_virtual_participant_handler(
lines.append('Пока нет виртуальных участников.')
rows = [
[types.InlineKeyboardButton(text=' Добавить', callback_data=f'admin_contest_vp_add_{contest_id}')],
[
types.InlineKeyboardButton(text=' Добавить', callback_data=f'admin_contest_vp_add_{contest_id}'),
types.InlineKeyboardButton(text='🎭 Массовка', callback_data=f'admin_contest_vp_mass_{contest_id}'),
],
]
if vps:
for v in vps:
@@ -1180,6 +1187,164 @@ async def delete_virtual_participant_handler(
)
@admin_required
@error_handler
async def start_mass_virtual_participants(
callback: types.CallbackQuery,
db_user,
db: AsyncSession,
state: FSMContext,
):
"""Начинает массовое создание виртуальных участников (массовка)."""
contest_id = int(callback.data.split('_')[-1])
await state.set_state(AdminStates.adding_mass_virtual_count)
await state.update_data(mass_vp_contest_id=contest_id)
text = """
🎭 <b>Массовка — массовое создание виртуальных участников</b>
<i>Для чего это нужно?</i>
Виртуальные участники (призраки) позволяют создать видимость активности в конкурсе. Они отображаются в таблице лидеров наравне с реальными участниками, но помечаются значком 👻.
Это помогает:
• Мотивировать реальных участников соревноваться
• Задать планку для участия
• Сделать конкурс более живым
<b>Введите количество призраков для создания:</b>
<i>(от 1 до 50)</i>
"""
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data=f'admin_contest_vp_{contest_id}')],
]
),
)
await callback.answer()
@admin_required
@error_handler
async def process_mass_virtual_count(
message: types.Message,
db_user,
db: AsyncSession,
state: FSMContext,
):
"""Обрабатывает количество призраков для массового создания."""
try:
count = int(message.text.strip())
if count < 1 or count > 50:
await message.answer(
'❌ Введите число от 1 до 50:',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data='admin_contests_ref')],
]
),
)
return
except ValueError:
await message.answer(
'❌ Введите корректное число от 1 до 50:',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data='admin_contests_ref')],
]
),
)
return
await state.update_data(mass_vp_count=count)
await state.set_state(AdminStates.adding_mass_virtual_referrals)
data = await state.get_data()
contest_id = data.get('mass_vp_contest_id')
await message.answer(
f'✅ Будет создано <b>{count}</b> призраков.\n\n'
f'<b>Введите количество рефералов у каждого:</b>\n'
f'<i>(от 1 до 100)</i>',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data=f'admin_contest_vp_{contest_id}')],
]
),
)
@admin_required
@error_handler
async def process_mass_virtual_referrals(
message: types.Message,
db_user,
db: AsyncSession,
state: FSMContext,
):
"""Создаёт массовку призраков с рандомными именами."""
import random
import string
try:
referrals_count = int(message.text.strip())
if referrals_count < 1 or referrals_count > 100:
await message.answer('❌ Введите число от 1 до 100:')
return
except ValueError:
await message.answer('❌ Введите корректное число от 1 до 100:')
return
data = await state.get_data()
contest_id = data.get('mass_vp_contest_id')
ghost_count = data.get('mass_vp_count', 1)
await state.clear()
# Генерируем и создаём призраков
created = []
for _ in range(ghost_count):
# Рандомное имя до 5 символов (буквы + цифры)
name_length = random.randint(3, 5)
name = ''.join(random.choices(string.ascii_letters + string.digits, k=name_length))
vp = await add_virtual_participant(db, contest_id, name, referrals_count)
created.append(vp)
# Показываем результат
text = f"""
✅ <b>Массовка создана!</b>
📊 <b>Результат:</b>
• Создано призраков: {len(created)}
• Рефералов у каждого: {referrals_count}
• Всего виртуальных рефералов: {len(created) * referrals_count}
👻 <b>Созданные призраки:</b>
"""
for vp in created[:10]:
text += f'{vp.display_name}{vp.referral_count} реф.\n'
if len(created) > 10:
text += f'<i>... и ещё {len(created) - 10}</i>\n'
await message.answer(
text,
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[
types.InlineKeyboardButton(
text='👻 К списку призраков', callback_data=f'admin_contest_vp_{contest_id}'
)
],
[types.InlineKeyboardButton(text='⬅️ К конкурсу', callback_data=f'admin_contest_view_{contest_id}')],
]
),
)
@admin_required
@error_handler
async def start_edit_virtual_participant(
@@ -1282,7 +1447,10 @@ def register_handlers(dp: Dispatcher):
dp.callback_query.register(start_add_virtual_participant, F.data.startswith('admin_contest_vp_add_'))
dp.callback_query.register(delete_virtual_participant_handler, F.data.startswith('admin_contest_vp_del_'))
dp.callback_query.register(start_edit_virtual_participant, F.data.startswith('admin_contest_vp_edit_'))
dp.callback_query.register(start_mass_virtual_participants, F.data.startswith('admin_contest_vp_mass_'))
dp.callback_query.register(show_virtual_participants, F.data.regexp(r'^admin_contest_vp_\d+$'))
dp.message.register(process_virtual_participant_name, AdminStates.adding_virtual_participant_name)
dp.message.register(process_virtual_participant_count, AdminStates.adding_virtual_participant_count)
dp.message.register(process_edit_virtual_participant_count, AdminStates.editing_virtual_participant_count)
dp.message.register(process_mass_virtual_count, AdminStates.adding_mass_virtual_count)
dp.message.register(process_mass_virtual_referrals, AdminStates.adding_mass_virtual_referrals)
+250 -1
View File
@@ -42,6 +42,10 @@ def _method_display(method: PaymentMethod) -> str:
return 'CryptoBot'
if method == PaymentMethod.TELEGRAM_STARS:
return 'Telegram Stars'
if method == PaymentMethod.KASSA_AI:
return settings.get_kassa_ai_display_name()
if method == PaymentMethod.FREEKASSA:
return settings.get_freekassa_display_name()
return method.value
@@ -144,6 +148,18 @@ def _status_info(
}
return mapping.get(status, ('', texts.t('ADMIN_PAYMENT_STATUS_UNKNOWN', '❓ Unknown')))
if record.method == PaymentMethod.KASSA_AI:
mapping = {
'pending': ('', texts.t('ADMIN_PAYMENT_STATUS_PENDING', '⏳ Pending')),
'created': ('', texts.t('ADMIN_PAYMENT_STATUS_PENDING', '⏳ Pending')),
'processing': ('', texts.t('ADMIN_PAYMENT_STATUS_PROCESSING', '⌛ Processing')),
'success': ('', texts.t('ADMIN_PAYMENT_STATUS_PAID', '✅ Paid')),
'paid': ('', texts.t('ADMIN_PAYMENT_STATUS_PAID', '✅ Paid')),
'canceled': ('', texts.t('ADMIN_PAYMENT_STATUS_CANCELED', '❌ Cancelled')),
'error': ('', texts.t('ADMIN_PAYMENT_STATUS_FAILED', '❌ Failed')),
}
return mapping.get(status, ('', texts.t('ADMIN_PAYMENT_STATUS_UNKNOWN', '❓ Unknown')))
return '', texts.t('ADMIN_PAYMENT_STATUS_UNKNOWN', '❓ Unknown')
@@ -168,7 +184,9 @@ def _is_checkable(record: PendingPayment) -> bool:
if record.method == PaymentMethod.CRYPTOBOT:
return status in {'active'}
if record.method == PaymentMethod.FREEKASSA:
return status in {'pending', ''}
return status in {'pending', 'created', ''}
if record.method == PaymentMethod.KASSA_AI:
return status in {'pending', 'created', 'processing', ''}
return False
@@ -184,6 +202,7 @@ def _build_list_keyboard(
page: int,
total_pages: int,
language: str,
has_checkable: bool = False,
) -> InlineKeyboardMarkup:
buttons: list[list[InlineKeyboardButton]] = []
texts = get_texts(language)
@@ -204,6 +223,28 @@ def _build_list_keyboard(
]
)
# Кнопка "Проверить все" если есть что проверять
if has_checkable:
buttons.append(
[
InlineKeyboardButton(
text=texts.t('ADMIN_PAYMENTS_CHECK_ALL', '🔄 Проверить все'),
callback_data='admin_payments_check_all',
)
]
)
# Кнопка экспорта если есть платежи
if records:
buttons.append(
[
InlineKeyboardButton(
text=texts.t('ADMIN_PAYMENTS_EXPORT', '📥 Выгрузить в файл'),
callback_data='admin_payments_export',
)
]
)
if total_pages > 1:
navigation_row: list[InlineKeyboardButton] = []
if page > 1:
@@ -485,11 +526,22 @@ async def show_payments_overview(
lines = [header, '', description]
# Проверяем есть ли платежи для массовой проверки
checkable_records = [r for r in records if _is_checkable(r) and not r.is_paid]
has_checkable = len(checkable_records) > 0
if page_records:
for idx, record in enumerate(page_records, start=start_index + 1):
lines.extend(_build_record_lines(record, index=idx, texts=texts, language=db_user.language))
lines.append('')
lines.append(notice)
if has_checkable:
lines.append('')
lines.append(
texts.t('ADMIN_PAYMENTS_CHECKABLE_COUNT', '🔄 Доступно для проверки: {count}').format(
count=len(checkable_records)
)
)
else:
empty_text = texts.t('ADMIN_PAYMENTS_EMPTY', 'No pending top-ups in the last 24 hours.')
lines.append('')
@@ -500,6 +552,7 @@ async def show_payments_overview(
page=page,
total_pages=total_pages,
language=db_user.language,
has_checkable=has_checkable,
)
await callback.message.edit_text(
@@ -550,28 +603,42 @@ async def manual_check_payment(
db_user: User,
db: AsyncSession,
) -> None:
import logging
logger = logging.getLogger(__name__)
logger.info('manual_check_payment called: %s', callback.data)
parsed = _parse_method_and_id(callback.data, prefix='admin_payment_check_')
if not parsed:
logger.warning('Failed to parse: %s', callback.data)
await callback.answer('❌ Invalid payment reference', show_alert=True)
return
method, payment_id = parsed
logger.info('Checking payment: method=%s, id=%s', method, payment_id)
record = await get_payment_record(db, method, payment_id)
texts = get_texts(db_user.language)
if not record:
logger.warning('Payment not found: method=%s, id=%s', method, payment_id)
await callback.answer(texts.t('ADMIN_PAYMENT_NOT_FOUND', 'Payment not found.'), show_alert=True)
return
logger.info('Record found: status=%s, is_paid=%s', record.status, record.is_paid)
if not _is_checkable(record):
logger.info('Payment not checkable: method=%s, status=%s', method, record.status)
await callback.answer(
texts.t('ADMIN_PAYMENT_CHECK_NOT_AVAILABLE', 'Manual check is not available for this invoice.'),
show_alert=True,
)
return
logger.info('Running manual check...')
payment_service = PaymentService(callback.bot)
updated = await run_manual_check(db, method, payment_id, payment_service)
logger.info('Check result: updated=%s', updated is not None)
if not updated:
await callback.answer(
@@ -597,7 +664,189 @@ async def manual_check_payment(
await callback.answer(message, show_alert=True)
@admin_required
@error_handler
async def check_all_payments(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession,
) -> None:
"""Массовая проверка всех ожидающих платежей."""
import logging
logger = logging.getLogger(__name__)
logger.info('check_all_payments called')
texts = get_texts(db_user.language)
# Получаем все ожидающие платежи
records = await list_recent_pending_payments(db)
logger.info('Found %d total records', len(records))
checkable_records = [r for r in records if _is_checkable(r) and not r.is_paid]
logger.info('Found %d checkable records', len(checkable_records))
if not checkable_records:
await callback.answer(
texts.t('ADMIN_PAYMENTS_NO_CHECKABLE', 'Нет платежей для проверки'),
show_alert=True,
)
return
await callback.answer(
texts.t('ADMIN_PAYMENTS_CHECKING_ALL', '🔄 Проверяю {count} платежей...').format(count=len(checkable_records)),
)
payment_service = PaymentService(callback.bot)
checked = 0
confirmed = 0
failed = 0
for record in checkable_records:
try:
logger.info('Checking %s payment id=%s', record.method.value, record.local_id)
updated = await run_manual_check(db, record.method, record.local_id, payment_service)
checked += 1
logger.info('Check result: is_paid=%s', updated.is_paid if updated else None)
if updated and updated.is_paid and not record.is_paid:
confirmed += 1
except Exception as e:
logger.error('Check failed for %s id=%s: %s', record.method.value, record.local_id, e, exc_info=True)
failed += 1
logger.info('Check complete: checked=%d, confirmed=%d, failed=%d', checked, confirmed, failed)
# Показываем результат
result_lines = [
texts.t('ADMIN_PAYMENTS_CHECK_ALL_RESULT', '🔄 <b>Результат проверки</b>'),
'',
texts.t('ADMIN_PAYMENTS_CHECK_ALL_CHECKED', '✅ Проверено: {count}').format(count=checked),
texts.t('ADMIN_PAYMENTS_CHECK_ALL_CONFIRMED', '💰 Подтверждено: {count}').format(count=confirmed),
]
if failed:
result_lines.append(texts.t('ADMIN_PAYMENTS_CHECK_ALL_FAILED', '❌ Ошибок: {count}').format(count=failed))
# Перезагружаем список платежей
records = await list_recent_pending_payments(db)
total = len(records)
total_pages = max(1, (total + PAGE_SIZE - 1) // PAGE_SIZE)
page_records = records[:PAGE_SIZE]
checkable_records = [r for r in records if _is_checkable(r) and not r.is_paid]
result_lines.append('')
result_lines.append(texts.t('ADMIN_PAYMENTS_TITLE', '💳 <b>Top-up verification</b>'))
if page_records:
result_lines.append('')
for idx, record in enumerate(page_records, start=1):
result_lines.extend(_build_record_lines(record, index=idx, texts=texts, language=db_user.language))
result_lines.append('')
keyboard = _build_list_keyboard(
page_records,
page=1,
total_pages=total_pages,
language=db_user.language,
has_checkable=len(checkable_records) > 0,
)
logger.info('Updating message with results...')
try:
await callback.message.edit_text(
'\n'.join(result_lines),
parse_mode='HTML',
reply_markup=keyboard,
)
logger.info('Message updated successfully')
except Exception as e:
logger.error('Failed to update message: %s', e, exc_info=True)
@admin_required
@error_handler
async def export_payments(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession,
) -> None:
"""Экспорт данных платежей в JSON файл."""
import json
from aiogram.types import BufferedInputFile
texts = get_texts(db_user.language)
records = await list_recent_pending_payments(db)
if not records:
await callback.answer(
texts.t('ADMIN_PAYMENTS_EXPORT_EMPTY', 'Нет платежей для экспорта'),
show_alert=True,
)
return
# Формируем данные для экспорта
export_data = []
for record in records:
payment = record.payment
user = record.user
payment_data = {
'id': record.local_id,
'method': record.method.value,
'method_display': _method_display(record.method),
'identifier': record.identifier,
'amount_kopeks': record.amount_kopeks,
'amount_rubles': record.amount_kopeks / 100,
'status': record.status,
'is_paid': record.is_paid,
'created_at': record.created_at.isoformat() if record.created_at else None,
'expires_at': record.expires_at.isoformat() if record.expires_at else None,
'user': {
'id': user.id,
'telegram_id': user.telegram_id,
'username': user.username,
'full_name': user.full_name,
},
}
# Добавляем специфичные поля в зависимости от метода
if hasattr(payment, 'order_id'):
payment_data['order_id'] = payment.order_id
if hasattr(payment, 'payment_url'):
payment_data['payment_url'] = payment.payment_url
if hasattr(payment, 'callback_payload'):
payment_data['callback_payload'] = payment.callback_payload
export_data.append(payment_data)
# Создаём JSON файл
json_content = json.dumps(export_data, ensure_ascii=False, indent=2, default=str)
file_bytes = json_content.encode('utf-8')
# Отправляем файл
from datetime import datetime
filename = f'payments_export_{datetime.now().strftime("%Y%m%d_%H%M%S")}.json'
await callback.message.answer_document(
document=BufferedInputFile(file_bytes, filename=filename),
caption=texts.t(
'ADMIN_PAYMENTS_EXPORT_CAPTION',
'📥 Экспорт платежей\n\n📊 Всего записей: {count}\n💰 Оплачено: {paid}\n⏳ Ожидают: {pending}',
).format(
count=len(export_data),
paid=sum(1 for r in export_data if r['is_paid']),
pending=sum(1 for r in export_data if not r['is_paid']),
),
)
await callback.answer(texts.t('ADMIN_PAYMENTS_EXPORT_SUCCESS', '✅ Файл отправлен'))
def register_handlers(dp: Dispatcher) -> None:
dp.callback_query.register(check_all_payments, F.data == 'admin_payments_check_all')
dp.callback_query.register(export_payments, F.data == 'admin_payments_export')
dp.callback_query.register(manual_check_payment, F.data.startswith('admin_payment_check_'))
dp.callback_query.register(
show_payment_details,
+802
View File
@@ -83,6 +83,7 @@ async def show_referral_statistics(callback: types.CallbackQuery, db_user: User,
keyboard_rows = [
[types.InlineKeyboardButton(text='🔄 Обновить', callback_data='admin_referrals')],
[types.InlineKeyboardButton(text='👥 Топ рефереров', callback_data='admin_referrals_top')],
[types.InlineKeyboardButton(text='🔍 Диагностика логов', callback_data='admin_referral_diagnostics')],
]
# Кнопка заявок на вывод (если функция включена)
@@ -650,11 +651,812 @@ async def process_test_referral_earning(message: types.Message, db_user: User, d
)
def _get_period_dates(period: str) -> tuple[datetime.datetime, datetime.datetime]:
"""Возвращает начальную и конечную даты для заданного периода."""
now = datetime.datetime.now()
today = now.replace(hour=0, minute=0, second=0, microsecond=0)
if period == 'today':
start_date = today
end_date = today + datetime.timedelta(days=1)
elif period == 'yesterday':
start_date = today - datetime.timedelta(days=1)
end_date = today
elif period == 'week':
start_date = today - datetime.timedelta(days=7)
end_date = today + datetime.timedelta(days=1)
elif period == 'month':
start_date = today - datetime.timedelta(days=30)
end_date = today + datetime.timedelta(days=1)
else:
# По умолчанию — сегодня
start_date = today
end_date = today + datetime.timedelta(days=1)
return start_date, end_date
def _get_period_display_name(period: str) -> str:
"""Возвращает человекочитаемое название периода."""
names = {'today': 'сегодня', 'yesterday': 'вчера', 'week': '7 дней', 'month': '30 дней'}
return names.get(period, 'сегодня')
async def _show_diagnostics_for_period(callback: types.CallbackQuery, db: AsyncSession, state: FSMContext, period: str):
"""Внутренняя функция для отображения диагностики за указанный период."""
try:
await callback.answer('Анализирую логи...')
from app.services.referral_diagnostics_service import referral_diagnostics_service
# Сохраняем период в state
await state.update_data(diagnostics_period=period)
from app.states import AdminStates
await state.set_state(AdminStates.referral_diagnostics_period)
# Получаем даты периода
start_date, end_date = _get_period_dates(period)
# Анализируем логи
report = await referral_diagnostics_service.analyze_period(db, start_date, end_date)
# Формируем отчёт
period_display = _get_period_display_name(period)
text = f"""
🔍 <b>Диагностика рефералов — {period_display}</b>
<b>📊 Статистика переходов:</b>
• Всего кликов по реф-ссылкам: {report.total_ref_clicks}
• Уникальных пользователей: {report.unique_users_clicked}
• Потерянных рефералов: {len(report.lost_referrals)}
"""
if report.lost_referrals:
text += '\n<b>❌ Потерянные рефералы:</b>\n'
text += '<i>(пришли по ссылке, но реферер не засчитался)</i>\n\n'
for i, lost in enumerate(report.lost_referrals[:15], 1):
# Статус пользователя
if not lost.registered:
status = '⚠️ Не в БД'
elif not lost.has_referrer:
status = '❌ Без реферера'
else:
status = f'⚡ Другой реферер (ID{lost.current_referrer_id})'
# Имя или ID
user_name = lost.username or lost.full_name or f'ID{lost.telegram_id}'
if lost.username:
user_name = f'@{lost.username}'
# Ожидаемый реферер
referrer_info = ''
if lost.expected_referrer_name:
referrer_info = f'{lost.expected_referrer_name}'
elif lost.expected_referrer_id:
referrer_info = f' → ID{lost.expected_referrer_id}'
# Время
time_str = lost.click_time.strftime('%H:%M')
text += f'{i}. {user_name}{status}\n'
text += f' <code>{lost.referral_code}</code>{referrer_info} ({time_str})\n'
if len(report.lost_referrals) > 15:
text += f'\n<i>... и ещё {len(report.lost_referrals) - 15}</i>\n'
else:
text += '\n✅ <b>Все рефералы засчитаны!</b>\n'
# Информация о логах
log_path = referral_diagnostics_service.log_path
log_exists = log_path.exists()
log_size = log_path.stat().st_size if log_exists else 0
text += f'\n<i>📂 {log_path.name}'
if log_exists:
text += f' ({log_size / 1024:.0f} KB)'
text += f' | Строк: {report.lines_in_period}'
else:
text += ' (не найден!)'
text += '</i>'
# Кнопки: только "Сегодня" (текущий лог) и "Загрузить файл" (старые логи)
keyboard_rows = [
[
types.InlineKeyboardButton(text='📅 Сегодня (текущий лог)', callback_data='admin_ref_diag:today'),
],
[types.InlineKeyboardButton(text='📤 Загрузить лог-файл', callback_data='admin_ref_diag_upload')],
[types.InlineKeyboardButton(text='🔍 Проверить бонусы (по БД)', callback_data='admin_ref_check_bonuses')],
[
types.InlineKeyboardButton(
text='🏆 Синхронизировать с конкурсом', callback_data='admin_ref_sync_contest'
)
],
]
# Кнопки действий (только если есть потерянные рефералы)
if report.lost_referrals:
keyboard_rows.append(
[types.InlineKeyboardButton(text='📋 Предпросмотр исправлений', callback_data='admin_ref_fix_preview')]
)
keyboard_rows.extend(
[
[types.InlineKeyboardButton(text='🔄 Обновить', callback_data=f'admin_ref_diag:{period}')],
[types.InlineKeyboardButton(text='⬅️ К статистике', callback_data='admin_referrals')],
]
)
keyboard = types.InlineKeyboardMarkup(inline_keyboard=keyboard_rows)
await callback.message.edit_text(text, reply_markup=keyboard)
except Exception as e:
logger.error(f'Ошибка в _show_diagnostics_for_period: {e}', exc_info=True)
await callback.answer('Ошибка при анализе логов', show_alert=True)
@admin_required
@error_handler
async def show_referral_diagnostics(callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext):
"""Показывает диагностику реферальной системы по логам."""
# Определяем период из callback_data или используем "today" по умолчанию
if ':' in callback.data:
period = callback.data.split(':')[1]
else:
period = 'today'
await _show_diagnostics_for_period(callback, db, state, period)
@admin_required
@error_handler
async def preview_referral_fixes(callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext):
"""Показывает предпросмотр исправлений потерянных рефералов."""
try:
await callback.answer('Анализирую...')
# Получаем период из state
state_data = await state.get_data()
period = state_data.get('diagnostics_period', 'today')
from app.services.referral_diagnostics_service import DiagnosticReport, referral_diagnostics_service
# Проверяем, работаем ли с загруженным файлом
if period == 'uploaded_file':
# Используем сохранённый отчёт из загруженного файла (десериализуем)
report_data = state_data.get('uploaded_file_report')
if not report_data:
await callback.answer('Отчёт загруженного файла не найден', show_alert=True)
return
report = DiagnosticReport.from_dict(report_data)
period_display = 'загруженный файл'
else:
# Получаем даты периода
start_date, end_date = _get_period_dates(period)
# Анализируем логи
report = await referral_diagnostics_service.analyze_period(db, start_date, end_date)
period_display = _get_period_display_name(period)
if not report.lost_referrals:
await callback.answer('Нет потерянных рефералов для исправления', show_alert=True)
return
# Запускаем предпросмотр исправлений
fix_report = await referral_diagnostics_service.fix_lost_referrals(db, report.lost_referrals, apply=False)
# Формируем отчёт
text = f"""
📋 <b>Предпросмотр исправлений — {period_display}</b>
<b>📊 Что будет сделано:</b>
• Исправлено рефералов: {fix_report.users_fixed}
• Бонусов рефералам: {settings.format_price(fix_report.bonuses_to_referrals)}
• Бонусов рефереам: {settings.format_price(fix_report.bonuses_to_referrers)}
• Ошибок: {fix_report.errors}
<b>🔍 Детали:</b>
"""
# Показываем первые 10 деталей
for i, detail in enumerate(fix_report.details[:10], 1):
user_name = detail.username or detail.full_name or f'ID{detail.telegram_id}'
if detail.username:
user_name = f'@{detail.username}'
if detail.error:
text += f'{i}. {user_name} — ❌ {detail.error}\n'
else:
text += f'{i}. {user_name}\n'
if detail.referred_by_set:
text += f' • Реферер: {detail.referrer_name or f"ID{detail.referrer_id}"}\n'
if detail.had_first_topup:
text += f' • Первое пополнение: {settings.format_price(detail.topup_amount_kopeks)}\n'
if detail.bonus_to_referral_kopeks > 0:
text += f' • Бонус рефералу: {settings.format_price(detail.bonus_to_referral_kopeks)}\n'
if detail.bonus_to_referrer_kopeks > 0:
text += f' • Бонус рефереру: {settings.format_price(detail.bonus_to_referrer_kopeks)}\n'
if len(fix_report.details) > 10:
text += f'\n<i>... и ещё {len(fix_report.details) - 10}</i>\n'
text += '\n⚠️ <b>Внимание!</b> Это только предпросмотр. Нажмите "Применить", чтобы выполнить исправления.'
# Кнопка назад зависит от источника
back_button_text = '⬅️ К диагностике'
back_button_callback = f'admin_ref_diag:{period}' if period != 'uploaded_file' else 'admin_referral_diagnostics'
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='✅ Применить исправления', callback_data='admin_ref_fix_apply')],
[types.InlineKeyboardButton(text=back_button_text, callback_data=back_button_callback)],
]
)
await callback.message.edit_text(text, reply_markup=keyboard)
except Exception as e:
logger.error(f'Ошибка в preview_referral_fixes: {e}', exc_info=True)
await callback.answer('Ошибка при создании предпросмотра', show_alert=True)
@admin_required
@error_handler
async def apply_referral_fixes(callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext):
"""Применяет исправления потерянных рефералов."""
try:
await callback.answer('Применяю исправления...')
# Получаем период из state
state_data = await state.get_data()
period = state_data.get('diagnostics_period', 'today')
from app.services.referral_diagnostics_service import DiagnosticReport, referral_diagnostics_service
# Проверяем, работаем ли с загруженным файлом
if period == 'uploaded_file':
# Используем сохранённый отчёт из загруженного файла (десериализуем)
report_data = state_data.get('uploaded_file_report')
if not report_data:
await callback.answer('Отчёт загруженного файла не найден', show_alert=True)
return
report = DiagnosticReport.from_dict(report_data)
period_display = 'загруженный файл'
else:
# Получаем даты периода
start_date, end_date = _get_period_dates(period)
# Анализируем логи
report = await referral_diagnostics_service.analyze_period(db, start_date, end_date)
period_display = _get_period_display_name(period)
if not report.lost_referrals:
await callback.answer('Нет потерянных рефералов для исправления', show_alert=True)
return
# Применяем исправления
fix_report = await referral_diagnostics_service.fix_lost_referrals(db, report.lost_referrals, apply=True)
# Формируем отчёт
text = f"""
✅ <b>Исправления применены — {period_display}</b>
<b>📊 Результаты:</b>
• Исправлено рефералов: {fix_report.users_fixed}
• Бонусов рефералам: {settings.format_price(fix_report.bonuses_to_referrals)}
• Бонусов рефереам: {settings.format_price(fix_report.bonuses_to_referrers)}
• Ошибок: {fix_report.errors}
<b>🔍 Детали:</b>
"""
# Показываем первые 10 успешных деталей
success_count = 0
for detail in fix_report.details:
if not detail.error and success_count < 10:
success_count += 1
user_name = detail.username or detail.full_name or f'ID{detail.telegram_id}'
if detail.username:
user_name = f'@{user_name}'
text += f'{success_count}. {user_name}\n'
if detail.referred_by_set:
text += f' • Реферер: {detail.referrer_name or f"ID{detail.referrer_id}"}\n'
if detail.bonus_to_referral_kopeks > 0:
text += f' • Бонус рефералу: {settings.format_price(detail.bonus_to_referral_kopeks)}\n'
if detail.bonus_to_referrer_kopeks > 0:
text += f' • Бонус рефереру: {settings.format_price(detail.bonus_to_referrer_kopeks)}\n'
if fix_report.users_fixed > 10:
text += f'\n<i>... и ещё {fix_report.users_fixed - 10} исправлений</i>\n'
# Показываем ошибки
if fix_report.errors > 0:
text += '\n<b>❌ Ошибки:</b>\n'
error_count = 0
for detail in fix_report.details:
if detail.error and error_count < 5:
error_count += 1
user_name = detail.username or detail.full_name or f'ID{detail.telegram_id}'
text += f'{user_name}: {detail.error}\n'
if fix_report.errors > 5:
text += f'<i>... и ещё {fix_report.errors - 5} ошибок</i>\n'
# Кнопки зависят от источника
keyboard_rows = []
if period != 'uploaded_file':
keyboard_rows.append(
[types.InlineKeyboardButton(text='🔄 Обновить диагностику', callback_data=f'admin_ref_diag:{period}')]
)
keyboard_rows.append([types.InlineKeyboardButton(text='⬅️ К статистике', callback_data='admin_referrals')])
keyboard = types.InlineKeyboardMarkup(inline_keyboard=keyboard_rows)
await callback.message.edit_text(text, reply_markup=keyboard)
# Очищаем сохранённый отчёт из state
if period == 'uploaded_file':
await state.update_data(uploaded_file_report=None)
except Exception as e:
logger.error(f'Ошибка в apply_referral_fixes: {e}', exc_info=True)
await callback.answer('Ошибка при применении исправлений', show_alert=True)
# =============================================================================
# Проверка бонусов по БД
# =============================================================================
@admin_required
@error_handler
async def check_missing_bonuses(callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext):
"""Проверяет по БД — всем ли рефералам начислены бонусы."""
from app.services.referral_diagnostics_service import (
referral_diagnostics_service,
)
await callback.answer('🔍 Проверяю бонусы...')
try:
report = await referral_diagnostics_service.check_missing_bonuses(db)
# Сохраняем отчёт в state для последующего применения
await state.update_data(missing_bonuses_report=report.to_dict())
text = f"""
🔍 <b>Проверка бонусов по БД</b>
📊 <b>Статистика:</b>
• Всего рефералов: {report.total_referrals_checked}
• С пополнением ≥ минимума: {report.referrals_with_topup}
• <b>Без бонусов: {len(report.missing_bonuses)}</b>
"""
if report.missing_bonuses:
text += f"""
💰 <b>Требуется начислить:</b>
• Рефералам: {report.total_missing_to_referrals / 100:.0f}
• Рефереерам: {report.total_missing_to_referrers / 100:.0f}
• <b>Итого: {(report.total_missing_to_referrals + report.total_missing_to_referrers) / 100:.0f}₽</b>
👤 <b>Список ({len(report.missing_bonuses)} чел.):</b>
"""
for i, mb in enumerate(report.missing_bonuses[:15], 1):
referral_name = mb.referral_full_name or mb.referral_username or str(mb.referral_telegram_id)
referrer_name = mb.referrer_full_name or mb.referrer_username or str(mb.referrer_telegram_id)
text += f'\n{i}. <b>{referral_name}</b>'
text += f'\n └ Пригласил: {referrer_name}'
text += f'\n └ Пополнение: {mb.first_topup_amount_kopeks / 100:.0f}'
text += f'\n └ Бонусы: {mb.referral_bonus_amount / 100:.0f}₽ + {mb.referrer_bonus_amount / 100:.0f}'
if len(report.missing_bonuses) > 15:
text += f'\n\n<i>... и ещё {len(report.missing_bonuses) - 15} чел.</i>'
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='✅ Начислить все бонусы', callback_data='admin_ref_bonus_apply')],
[types.InlineKeyboardButton(text='🔄 Обновить', callback_data='admin_ref_check_bonuses')],
[types.InlineKeyboardButton(text='⬅️ К диагностике', callback_data='admin_referral_diagnostics')],
]
)
else:
text += '\n✅ <b>Все бонусы начислены!</b>'
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='🔄 Обновить', callback_data='admin_ref_check_bonuses')],
[types.InlineKeyboardButton(text='⬅️ К диагностике', callback_data='admin_referral_diagnostics')],
]
)
await callback.message.edit_text(text, reply_markup=keyboard)
except Exception as e:
logger.error(f'Ошибка в check_missing_bonuses: {e}', exc_info=True)
await callback.answer('Ошибка при проверке бонусов', show_alert=True)
@admin_required
@error_handler
async def apply_missing_bonuses(callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext):
"""Применяет начисление пропущенных бонусов."""
from app.services.referral_diagnostics_service import (
MissingBonusReport,
referral_diagnostics_service,
)
await callback.answer('💰 Начисляю бонусы...')
try:
# Получаем сохранённый отчёт
data = await state.get_data()
report_dict = data.get('missing_bonuses_report')
if not report_dict:
await callback.answer('❌ Отчёт не найден. Обновите проверку.', show_alert=True)
return
report = MissingBonusReport.from_dict(report_dict)
if not report.missing_bonuses:
await callback.answer('✅ Нет бонусов для начисления', show_alert=True)
return
# Применяем исправления
fix_report = await referral_diagnostics_service.fix_missing_bonuses(db, report.missing_bonuses, apply=True)
text = f"""
✅ <b>Бонусы начислены!</b>
📊 <b>Результат:</b>
• Обработано: {fix_report.users_fixed} пользователей
• Начислено рефералам: {fix_report.bonuses_to_referrals / 100:.0f}
• Начислено рефереерам: {fix_report.bonuses_to_referrers / 100:.0f}
• <b>Итого: {(fix_report.bonuses_to_referrals + fix_report.bonuses_to_referrers) / 100:.0f}₽</b>
"""
if fix_report.errors > 0:
text += f'\n⚠️ Ошибок: {fix_report.errors}'
# Очищаем отчёт из state
await state.update_data(missing_bonuses_report=None)
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='🔍 Проверить снова', callback_data='admin_ref_check_bonuses')],
[types.InlineKeyboardButton(text='⬅️ К диагностике', callback_data='admin_referral_diagnostics')],
]
)
await callback.message.edit_text(text, reply_markup=keyboard)
except Exception as e:
logger.error(f'Ошибка в apply_missing_bonuses: {e}', exc_info=True)
await callback.answer('Ошибка при начислении бонусов', show_alert=True)
@admin_required
@error_handler
async def sync_referrals_with_contest(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext
):
"""Синхронизирует всех рефералов с активными конкурсами."""
from app.database.crud.referral_contest import get_contests_for_events
from app.services.referral_contest_service import referral_contest_service
await callback.answer('🏆 Синхронизирую с конкурсами...')
try:
from datetime import datetime
now_utc = datetime.utcnow()
# Получаем активные конкурсы
paid_contests = await get_contests_for_events(db, now_utc, contest_types=['referral_paid'])
reg_contests = await get_contests_for_events(db, now_utc, contest_types=['referral_registered'])
all_contests = list(paid_contests) + list(reg_contests)
if not all_contests:
await callback.message.edit_text(
'❌ <b>Нет активных конкурсов рефералов</b>\n\n'
'Создайте конкурс в разделе "Конкурсы" для синхронизации.',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='⬅️ К диагностике', callback_data='admin_referral_diagnostics')]
]
),
)
return
# Синхронизируем каждый конкурс
total_created = 0
total_updated = 0
total_skipped = 0
contest_results = []
for contest in all_contests:
stats = await referral_contest_service.sync_contest(db, contest.id)
if 'error' not in stats:
total_created += stats.get('created', 0)
total_updated += stats.get('updated', 0)
total_skipped += stats.get('skipped', 0)
contest_results.append(f'{contest.title}: +{stats.get("created", 0)} новых')
else:
contest_results.append(f'{contest.title}: ошибка')
text = f"""
🏆 <b>Синхронизация с конкурсами завершена!</b>
📊 <b>Результат:</b>
• Конкурсов обработано: {len(all_contests)}
• Новых событий добавлено: {total_created}
• Обновлено: {total_updated}
• Пропущено (уже есть): {total_skipped}
📋 <b>По конкурсам:</b>
"""
text += '\n'.join(contest_results)
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='🔄 Синхронизировать снова', callback_data='admin_ref_sync_contest')],
[types.InlineKeyboardButton(text='⬅️ К диагностике', callback_data='admin_referral_diagnostics')],
]
)
await callback.message.edit_text(text, reply_markup=keyboard)
except Exception as e:
logger.error(f'Ошибка в sync_referrals_with_contest: {e}', exc_info=True)
await callback.answer('Ошибка при синхронизации', show_alert=True)
@admin_required
@error_handler
async def request_log_file_upload(callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext):
"""Запрашивает загрузку лог-файла для анализа."""
await state.set_state(AdminStates.waiting_for_log_file)
text = """
📤 <b>Загрузка лог-файла для анализа</b>
Отправьте файл лога (расширение .log или .txt).
Файл будет проанализирован на наличие потерянных рефералов за ВСЕ время, записанное в логе.
⚠️ <b>Важно:</b>
• Файл должен быть текстовым (.log, .txt)
• Максимальный размер: 50 MB
• После анализа файл будет автоматически удалён
Если ротация логов удалила старые данные — загрузите резервную копию.
"""
keyboard = types.InlineKeyboardMarkup(
inline_keyboard=[[types.InlineKeyboardButton(text='❌ Отмена', callback_data='admin_referral_diagnostics')]]
)
await callback.message.edit_text(text, reply_markup=keyboard)
await callback.answer()
@admin_required
@error_handler
async def receive_log_file(message: types.Message, db_user: User, db: AsyncSession, state: FSMContext):
"""Получает и анализирует загруженный лог-файл."""
import tempfile
from pathlib import Path
if not message.document:
await message.answer(
'❌ Пожалуйста, отправьте файл документом.',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data='admin_referral_diagnostics')]
]
),
)
return
# Проверяем расширение файла
file_name = message.document.file_name or 'unknown'
file_ext = Path(file_name).suffix.lower()
if file_ext not in ['.log', '.txt']:
await message.answer(
f'❌ Неверный формат файла: {file_ext}\n\nПоддерживаются только текстовые файлы (.log, .txt)',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data='admin_referral_diagnostics')]
]
),
)
return
# Проверяем размер файла
max_size = 50 * 1024 * 1024 # 50 MB
if message.document.file_size > max_size:
await message.answer(
f'❌ Файл слишком большой: {message.document.file_size / 1024 / 1024:.1f} MB\n\nМаксимальный размер: 50 MB',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='❌ Отмена', callback_data='admin_referral_diagnostics')]
]
),
)
return
# Информируем о начале загрузки
status_message = await message.answer(
f'📥 Загружаю файл {file_name} ({message.document.file_size / 1024 / 1024:.1f} MB)...'
)
temp_file_path = None
try:
# Скачиваем файл во временную директорию
temp_dir = tempfile.gettempdir()
temp_file_path = str(Path(temp_dir) / f'ref_diagnostics_{message.from_user.id}_{file_name}')
# Скачиваем файл
file = await message.bot.get_file(message.document.file_id)
await message.bot.download_file(file.file_path, temp_file_path)
logger.info(f'📥 Файл загружен: {temp_file_path} ({message.document.file_size} байт)')
# Обновляем статус
await status_message.edit_text(f'🔍 Анализирую файл {file_name}...\n\nЭто может занять некоторое время.')
# Анализируем файл
from app.services.referral_diagnostics_service import referral_diagnostics_service
report = await referral_diagnostics_service.analyze_file(db, temp_file_path)
# Формируем отчёт
text = f"""
🔍 <b>Анализ лог-файла: {file_name}</b>
<b>📊 Статистика переходов:</b>
• Всего кликов по реф-ссылкам: {report.total_ref_clicks}
• Уникальных пользователей: {report.unique_users_clicked}
• Потерянных рефералов: {len(report.lost_referrals)}
• Строк в файле: {report.lines_in_period}
"""
if report.lost_referrals:
text += '\n<b>❌ Потерянные рефералы:</b>\n'
text += '<i>(пришли по ссылке, но реферер не засчитался)</i>\n\n'
for i, lost in enumerate(report.lost_referrals[:15], 1):
# Статус пользователя
if not lost.registered:
status = '⚠️ Не в БД'
elif not lost.has_referrer:
status = '❌ Без реферера'
else:
status = f'⚡ Другой реферер (ID{lost.current_referrer_id})'
# Имя или ID
user_name = lost.username or lost.full_name or f'ID{lost.telegram_id}'
if lost.username:
user_name = f'@{lost.username}'
# Ожидаемый реферер
referrer_info = ''
if lost.expected_referrer_name:
referrer_info = f'{lost.expected_referrer_name}'
elif lost.expected_referrer_id:
referrer_info = f' → ID{lost.expected_referrer_id}'
# Время
time_str = lost.click_time.strftime('%d.%m.%Y %H:%M')
text += f'{i}. {user_name}{status}\n'
text += f' <code>{lost.referral_code}</code>{referrer_info} ({time_str})\n'
if len(report.lost_referrals) > 15:
text += f'\n<i>... и ещё {len(report.lost_referrals) - 15}</i>\n'
else:
text += '\n✅ <b>Все рефералы засчитаны!</b>\n'
# Сохраняем отчёт в state для дальнейшего использования (сериализуем в dict)
await state.update_data(
diagnostics_period='uploaded_file',
uploaded_file_report=report.to_dict(),
)
# Кнопки действий
keyboard_rows = []
if report.lost_referrals:
keyboard_rows.append(
[types.InlineKeyboardButton(text='📋 Предпросмотр исправлений', callback_data='admin_ref_fix_preview')]
)
keyboard_rows.extend(
[
[types.InlineKeyboardButton(text='⬅️ К диагностике', callback_data='admin_referral_diagnostics')],
[types.InlineKeyboardButton(text='⬅️ К статистике', callback_data='admin_referrals')],
]
)
keyboard = types.InlineKeyboardMarkup(inline_keyboard=keyboard_rows)
# Удаляем статусное сообщение
await status_message.delete()
# Отправляем результат
await message.answer(text, reply_markup=keyboard)
# Очищаем состояние
await state.set_state(AdminStates.referral_diagnostics_period)
except Exception as e:
logger.error(f'❌ Ошибка при обработке файла: {e}', exc_info=True)
try:
await status_message.edit_text(
f'❌ <b>Ошибка при анализе файла</b>\n\n'
f'Файл: {file_name}\n'
f'Ошибка: {e!s}\n\n'
f'Проверьте, что файл является текстовым логом бота.',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[
types.InlineKeyboardButton(
text='🔄 Попробовать снова', callback_data='admin_ref_diag_upload'
)
],
[
types.InlineKeyboardButton(
text='⬅️ К диагностике', callback_data='admin_referral_diagnostics'
)
],
]
),
)
except:
await message.answer(
f'❌ Ошибка при анализе файла: {e!s}',
reply_markup=types.InlineKeyboardMarkup(
inline_keyboard=[
[types.InlineKeyboardButton(text='⬅️ Назад', callback_data='admin_referral_diagnostics')]
]
),
)
finally:
# Удаляем временный файл
if temp_file_path and Path(temp_file_path).exists():
try:
Path(temp_file_path).unlink()
logger.info(f'🗑️ Временный файл удалён: {temp_file_path}')
except Exception as e:
logger.error(f'Ошибка удаления временного файла: {e}')
def register_handlers(dp: Dispatcher):
dp.callback_query.register(show_referral_statistics, F.data == 'admin_referrals')
dp.callback_query.register(show_top_referrers, F.data == 'admin_referrals_top')
dp.callback_query.register(show_top_referrers_filtered, F.data.startswith('admin_top_ref:'))
dp.callback_query.register(show_referral_settings, F.data == 'admin_referrals_settings')
dp.callback_query.register(show_referral_diagnostics, F.data == 'admin_referral_diagnostics')
dp.callback_query.register(show_referral_diagnostics, F.data.startswith('admin_ref_diag:'))
dp.callback_query.register(preview_referral_fixes, F.data == 'admin_ref_fix_preview')
dp.callback_query.register(apply_referral_fixes, F.data == 'admin_ref_fix_apply')
# Загрузка лог-файла
dp.callback_query.register(request_log_file_upload, F.data == 'admin_ref_diag_upload')
dp.message.register(receive_log_file, AdminStates.waiting_for_log_file)
# Проверка бонусов по БД
dp.callback_query.register(check_missing_bonuses, F.data == 'admin_ref_check_bonuses')
dp.callback_query.register(apply_missing_bonuses, F.data == 'admin_ref_bonus_apply')
dp.callback_query.register(sync_referrals_with_contest, F.data == 'admin_ref_sync_contest')
# Хендлеры заявок на вывод
dp.callback_query.register(show_pending_withdrawal_requests, F.data == 'admin_withdrawal_requests')
+137 -188
View File
@@ -65,11 +65,8 @@ class UserFilterType(Enum):
"""Типы фильтрации пользователей."""
BALANCE = 'balance'
TRAFFIC = 'traffic'
ACTIVITY = 'activity'
SPENDING = 'spending'
PURCHASES = 'purchases'
CAMPAIGN = 'campaign'
POTENTIAL_CUSTOMERS = 'potential_customers'
@dataclass
@@ -92,34 +89,6 @@ USER_FILTER_CONFIGS: dict[UserFilterType, UserFilterConfig] = {
pagination_prefix='admin_users_balance_list',
order_param='order_by_balance',
),
UserFilterType.TRAFFIC: UserFilterConfig(
fsm_state=AdminStates.viewing_user_from_traffic_list,
title='👥 <b>Список пользователей по использованному трафику</b>',
empty_message='📶 Пользователи с трафиком не найдены',
pagination_prefix='admin_users_traffic_list',
order_param='order_by_traffic',
),
UserFilterType.ACTIVITY: UserFilterConfig(
fsm_state=AdminStates.viewing_user_from_last_activity_list,
title='👥 <b>Пользователи по активности</b>',
empty_message='🕒 Пользователи с активностью не найдены',
pagination_prefix='admin_users_activity_list',
order_param='order_by_last_activity',
),
UserFilterType.SPENDING: UserFilterConfig(
fsm_state=AdminStates.viewing_user_from_spending_list,
title='👥 <b>Пользователи по сумме трат</b>',
empty_message='💳 Пользователи с тратами не найдены',
pagination_prefix='admin_users_spending_list',
order_param='order_by_total_spent',
),
UserFilterType.PURCHASES: UserFilterConfig(
fsm_state=AdminStates.viewing_user_from_purchases_list,
title='👥 <b>Пользователи по количеству покупок</b>',
empty_message='🛒 Пользователи с покупками не найдены',
pagination_prefix='admin_users_purchases_list',
order_param='order_by_purchase_count',
),
UserFilterType.CAMPAIGN: UserFilterConfig(
fsm_state=AdminStates.viewing_user_from_campaign_list,
title='👥 <b>Пользователи по кампании регистрации</b>',
@@ -127,6 +96,13 @@ USER_FILTER_CONFIGS: dict[UserFilterType, UserFilterConfig] = {
pagination_prefix='admin_users_campaign_list',
order_param='', # использует специальный метод
),
UserFilterType.POTENTIAL_CUSTOMERS: UserFilterConfig(
fsm_state=AdminStates.viewing_user_from_potential_customers_list,
title='👥 <b>Потенциальные клиенты</b>',
empty_message='💰 Потенциальные клиенты не найдены',
pagination_prefix='admin_users_potential_customers_list',
order_param='', # использует специальный метод
),
}
@@ -173,34 +149,6 @@ def _build_user_button_text(
days_left = (user.subscription.end_date - datetime.utcnow()).days
button_text += f' | 📅 {days_left}д'
elif filter_type == UserFilterType.TRAFFIC:
if user.subscription:
sub = user.subscription
used = sub.traffic_used_gb or 0.0
if sub.traffic_limit_gb and sub.traffic_limit_gb > 0:
limit_display = f'{sub.traffic_limit_gb}'
else:
limit_display = '♾️'
traffic_display = f'{used:.1f}/{limit_display} ГБ'
else:
traffic_display = 'нет подписки'
button_text = f'{status_emoji} {sub_emoji} {user.full_name} | 📶 {traffic_display}'
if user.balance_kopeks > 0:
button_text += f' | 💰 {settings.format_price(user.balance_kopeks)}'
elif filter_type == UserFilterType.ACTIVITY:
activity_display = format_time_ago(user.last_activity, language) if user.last_activity else 'неизвестно'
button_text = f'{status_emoji} {sub_emoji} {user.full_name} | 🕒 {activity_display}'
elif filter_type in (UserFilterType.SPENDING, UserFilterType.PURCHASES):
stats = extra_data.get(user.id, {'total_spent': 0, 'purchase_count': 0}) if extra_data else {}
total_spent = stats.get('total_spent', 0)
purchases = stats.get('purchase_count', 0)
if filter_type == UserFilterType.SPENDING:
button_text = f'{status_emoji} {user.full_name} | 💳 {settings.format_price(total_spent)} | 🛒 {purchases}'
else:
button_text = f'{status_emoji} {user.full_name} | 🛒 {purchases} | 💳 {settings.format_price(total_spent)}'
elif filter_type == UserFilterType.CAMPAIGN:
info = extra_data.get(user.id, {}) if extra_data else {}
campaign_name = info.get('campaign_name') or 'Без кампании'
@@ -219,18 +167,6 @@ def _build_user_button_text(
button_text = f'{status_emoji} {sub_emoji} {short_name}'
if user.balance_kopeks > 0:
button_text += f' | 💰 {settings.format_price(user.balance_kopeks)}'
elif filter_type == UserFilterType.TRAFFIC:
if user.subscription:
sub = user.subscription
used = sub.traffic_used_gb or 0.0
if sub.traffic_limit_gb and sub.traffic_limit_gb > 0:
limit_display = f'{sub.traffic_limit_gb}'
else:
limit_display = '♾️'
traffic_display = f'{used:.1f}/{limit_display} ГБ'
else:
traffic_display = 'нет'
button_text = f'{status_emoji} {sub_emoji} {short_name} | 📶 {traffic_display}'
else:
button_text = f'{status_emoji} {short_name}'
@@ -280,10 +216,6 @@ async def _show_users_list_filtered(
await callback.answer()
return
# Для spending/purchases нужны дополнительные данные
if filter_type in (UserFilterType.SPENDING, UserFilterType.PURCHASES):
extra_data = await user_service.get_user_spending_stats_map(db, [user.id for user in users])
# Формируем текст заголовка
text = f'{config.title} (стр. {page}/{users_data["total_pages"]})\n\n'
text += 'Нажмите на пользователя для управления:'
@@ -576,38 +508,122 @@ async def show_users_ready_to_renew(
@admin_required
@error_handler
async def show_users_list_by_traffic(
async def show_potential_customers(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext, page: int = 1
):
"""Список пользователей, отсортированный по использованному трафику (убывание)."""
await _show_users_list_filtered(callback, db_user, db, state, UserFilterType.TRAFFIC, page)
"""Показывает пользователей без активной подписки с балансом >= месячной цены."""
await state.set_state(AdminStates.viewing_user_from_potential_customers_list)
texts = get_texts(db_user.language)
from app.config import PERIOD_PRICES
@admin_required
@error_handler
async def show_users_list_by_last_activity(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext, page: int = 1
):
"""Список пользователей, отсортированный по последней активности."""
await _show_users_list_filtered(callback, db_user, db, state, UserFilterType.ACTIVITY, page)
monthly_price = PERIOD_PRICES.get(30, 99000)
user_service = UserService()
users_data = await user_service.get_potential_customers(
db,
min_balance_kopeks=monthly_price,
page=page,
limit=10,
)
@admin_required
@error_handler
async def show_users_list_by_spending(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext, page: int = 1
):
"""Список пользователей, отсортированный по сумме трат (убывание)."""
await _show_users_list_filtered(callback, db_user, db, state, UserFilterType.SPENDING, page)
amount_text = settings.format_price(monthly_price)
header = texts.t(
'ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS_TITLE',
'💰 Потенциальные клиенты',
)
description = texts.t(
'ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS_DESC',
'Пользователи без активной подписки с балансом {amount} или больше.',
).format(amount=amount_text)
if not users_data['users']:
empty_text = texts.t(
'ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS_EMPTY',
'Сейчас нет пользователей, которые подходят под этот фильтр.',
)
await callback.message.edit_text(
f'{header}\n\n{description}\n\n{empty_text}',
reply_markup=get_admin_users_keyboard(db_user.language),
)
await callback.answer()
return
@admin_required
@error_handler
async def show_users_list_by_purchases(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext, page: int = 1
):
"""Список пользователей, отсортированный по количеству покупок (убывание)."""
await _show_users_list_filtered(callback, db_user, db, state, UserFilterType.PURCHASES, page)
text = f'{header}\n\n{description}\n\n'
text += 'Нажмите на пользователя для управления:'
keyboard = []
for user in users_data['users']:
subscription = user.subscription
status_emoji = '' if user.status == UserStatus.ACTIVE.value else '🚫'
subscription_emoji = ''
if subscription:
if subscription.is_trial:
subscription_emoji = '🎁'
elif subscription.is_active:
subscription_emoji = '💎'
else:
subscription_emoji = ''
button_text = (
f'{status_emoji} {subscription_emoji} {user.full_name} | 💰 {settings.format_price(user.balance_kopeks)}'
)
if len(button_text) > 60:
short_name = user.full_name
if len(short_name) > 20:
short_name = short_name[:17] + '...'
button_text = (
f'{status_emoji} {subscription_emoji} {short_name} | 💰 {settings.format_price(user.balance_kopeks)}'
)
keyboard.append(
[
types.InlineKeyboardButton(
text=button_text,
callback_data=f'admin_user_manage_{user.id}',
)
]
)
if users_data['total_pages'] > 1:
pagination_row = get_admin_pagination_keyboard(
users_data['current_page'],
users_data['total_pages'],
'admin_users_potential_customers_list',
'admin_users_potential_customers_filter',
db_user.language,
).inline_keyboard[0]
keyboard.append(pagination_row)
keyboard.extend(
[
[
types.InlineKeyboardButton(
text='🔍 Поиск',
callback_data='admin_users_search',
),
types.InlineKeyboardButton(
text='📊 Статистика',
callback_data='admin_users_stats',
),
],
[
types.InlineKeyboardButton(
text='⬅️ Назад',
callback_data='admin_users',
)
],
]
)
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard),
)
await callback.answer()
@admin_required
@@ -647,62 +663,6 @@ async def handle_users_balance_list_pagination(
await show_users_list_by_balance(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_users_traffic_list_pagination(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext
):
try:
callback_parts = callback.data.split('_')
page = int(callback_parts[-1])
await show_users_list_by_traffic(callback, db_user, db, state, page)
except (ValueError, IndexError) as e:
logger.error(f'Ошибка парсинга номера страницы: {e}')
await show_users_list_by_traffic(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_users_activity_list_pagination(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext
):
try:
callback_parts = callback.data.split('_')
page = int(callback_parts[-1])
await show_users_list_by_last_activity(callback, db_user, db, state, page)
except (ValueError, IndexError) as e:
logger.error(f'Ошибка парсинга номера страницы: {e}')
await show_users_list_by_last_activity(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_users_spending_list_pagination(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext
):
try:
callback_parts = callback.data.split('_')
page = int(callback_parts[-1])
await show_users_list_by_spending(callback, db_user, db, state, page)
except (ValueError, IndexError) as e:
logger.error(f'Ошибка парсинга номера страницы: {e}')
await show_users_list_by_spending(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_users_purchases_list_pagination(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext
):
try:
callback_parts = callback.data.split('_')
page = int(callback_parts[-1])
await show_users_list_by_purchases(callback, db_user, db, state, page)
except (ValueError, IndexError) as e:
logger.error(f'Ошибка парсинга номера страницы: {e}')
await show_users_list_by_purchases(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_users_ready_to_renew_pagination(
@@ -716,6 +676,19 @@ async def handle_users_ready_to_renew_pagination(
await show_users_ready_to_renew(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_potential_customers_pagination(
callback: types.CallbackQuery, db_user: User, db: AsyncSession, state: FSMContext
):
try:
page = int(callback.data.split('_')[-1])
await show_potential_customers(callback, db_user, db, state, page)
except (ValueError, IndexError) as e:
logger.error(f'Ошибка парсинга номера страницы: {e}')
await show_potential_customers(callback, db_user, db, state, 1)
@admin_required
@error_handler
async def handle_users_campaign_list_pagination(
@@ -1310,18 +1283,12 @@ async def show_user_management(callback: types.CallbackQuery, db_user: User, db:
current_state = await state.get_state()
if current_state == AdminStates.viewing_user_from_balance_list:
back_callback = 'admin_users_balance_filter'
elif current_state == AdminStates.viewing_user_from_traffic_list:
back_callback = 'admin_users_traffic_filter'
elif current_state == AdminStates.viewing_user_from_last_activity_list:
back_callback = 'admin_users_activity_filter'
elif current_state == AdminStates.viewing_user_from_spending_list:
back_callback = 'admin_users_spending_filter'
elif current_state == AdminStates.viewing_user_from_purchases_list:
back_callback = 'admin_users_purchases_filter'
elif current_state == AdminStates.viewing_user_from_campaign_list:
back_callback = 'admin_users_campaign_filter'
elif current_state == AdminStates.viewing_user_from_ready_to_renew_list:
back_callback = 'admin_users_ready_to_renew_filter'
elif current_state == AdminStates.viewing_user_from_potential_customers_list:
back_callback = 'admin_users_potential_customers_filter'
# Базовая клавиатура профиля
kb = get_user_management_keyboard(user.id, user.status, db_user.language, back_callback)
@@ -5483,26 +5450,14 @@ def register_handlers(dp: Dispatcher):
handle_users_balance_list_pagination, F.data.startswith('admin_users_balance_list_page_')
)
dp.callback_query.register(
handle_users_traffic_list_pagination, F.data.startswith('admin_users_traffic_list_page_')
)
dp.callback_query.register(
handle_users_activity_list_pagination, F.data.startswith('admin_users_activity_list_page_')
)
dp.callback_query.register(
handle_users_spending_list_pagination, F.data.startswith('admin_users_spending_list_page_')
)
dp.callback_query.register(
handle_users_purchases_list_pagination, F.data.startswith('admin_users_purchases_list_page_')
)
dp.callback_query.register(
handle_users_ready_to_renew_pagination, F.data.startswith('admin_users_ready_to_renew_list_page_')
)
dp.callback_query.register(
handle_potential_customers_pagination, F.data.startswith('admin_users_potential_customers_list_page_')
)
dp.callback_query.register(
handle_users_campaign_list_pagination, F.data.startswith('admin_users_campaign_list_page_')
)
@@ -5659,14 +5614,8 @@ def register_handlers(dp: Dispatcher):
dp.callback_query.register(show_users_list_by_balance, F.data == 'admin_users_balance_filter')
dp.callback_query.register(show_users_list_by_traffic, F.data == 'admin_users_traffic_filter')
dp.callback_query.register(show_users_list_by_last_activity, F.data == 'admin_users_activity_filter')
dp.callback_query.register(show_users_list_by_spending, F.data == 'admin_users_spending_filter')
dp.callback_query.register(show_users_list_by_purchases, F.data == 'admin_users_purchases_filter')
dp.callback_query.register(show_users_ready_to_renew, F.data == 'admin_users_ready_to_renew_filter')
dp.callback_query.register(show_potential_customers, F.data == 'admin_users_potential_customers_filter')
dp.callback_query.register(show_users_list_by_campaign, F.data == 'admin_users_campaign_filter')
+54 -14
View File
@@ -309,24 +309,27 @@ async def cmd_start(message: types.Message, state: FSMContext, db: AsyncSession,
logger.info(f'🚀 START: Обработка /start от {message.from_user.id}')
data = await state.get_data() or {}
had_pending_payload = 'pending_start_payload' in data
pending_start_payload = data.pop('pending_start_payload', None)
had_campaign_notification_flag = 'campaign_notification_sent' in data
campaign_notification_sent = data.pop('campaign_notification_sent', False)
state_needs_update = had_pending_payload or had_campaign_notification_flag
# ИСПРАВЛЕНИЕ БАГА: используем .get() вместо .pop() для campaign_notification_sent
# pending_start_payload обрабатывается отдельно ниже
campaign_notification_sent = data.get('campaign_notification_sent', False)
state_needs_update = False
# Получаем payload из state или Redis
pending_start_payload = data.get('pending_start_payload', None)
# Если в FSM state нет payload, пробуем получить из Redis (резервный механизм)
if not pending_start_payload:
redis_payload = await get_pending_payload_from_redis(message.from_user.id)
if redis_payload:
pending_start_payload = redis_payload
data['pending_start_payload'] = redis_payload
state_needs_update = True
logger.info(
"📦 START: Payload '%s' восстановлен из Redis (fallback)",
pending_start_payload,
)
# Очищаем Redis после получения
await delete_pending_payload_from_redis(message.from_user.id)
# НЕ удаляем Redis payload здесь - удаление только после успешной регистрации
referral_code = None
campaign = None
@@ -1167,6 +1170,12 @@ async def complete_registration_from_callback(callback: types.CallbackQuery, sta
refresh_subscription_error,
)
# ИСПРАВЛЕНИЕ БАГА: Очищаем Redis payload после успешной регистрации
await delete_pending_payload_from_redis(callback.from_user.id)
logger.info(
'🗑️ COMPLETE_FROM_CALLBACK: Redis payload удален после успешной регистрации пользователя %s', user.telegram_id
)
await state.clear()
if campaign_message:
@@ -1454,6 +1463,10 @@ async def complete_registration(message: types.Message, state: FSMContext, db: A
refresh_subscription_error,
)
# ИСПРАВЛЕНИЕ БАГА: Очищаем Redis payload после успешной регистрации
await delete_pending_payload_from_redis(message.from_user.id)
logger.info('🗑️ COMPLETE: Redis payload удален после успешной регистрации пользователя %s', user.telegram_id)
await state.clear()
if campaign_message:
@@ -1725,6 +1738,7 @@ async def required_sub_channel_check(
redis_payload = await get_pending_payload_from_redis(query.from_user.id)
if redis_payload:
pending_start_payload = redis_payload
state_data['pending_start_payload'] = redis_payload
logger.info(
"📦 CHANNEL CHECK: Payload '%s' восстановлен из Redis (fallback)",
pending_start_payload,
@@ -1780,16 +1794,31 @@ async def required_sub_channel_check(
only_active=True,
)
if campaign:
state_data['campaign_id'] = campaign.id
logger.info(
'📣 CHANNEL CHECK: Кампания %s восстановлена из payload',
campaign.id,
# Обрабатываем payload только если ещё не обработан
# (проверяем по наличию referral_code или campaign_id в state)
if not state_data.get('referral_code') and not state_data.get('campaign_id'):
campaign = await get_campaign_by_start_parameter(
db,
pending_start_payload,
only_active=True,
)
if campaign:
state_data['campaign_id'] = campaign.id
logger.info(
'📣 CHANNEL CHECK: Кампания %s восстановлена из payload',
campaign.id,
)
else:
state_data['referral_code'] = pending_start_payload
logger.info(
'🎯 CHANNEL CHECK: Payload интерпретирован как реферальный код: %s',
pending_start_payload,
)
else:
state_data['referral_code'] = pending_start_payload
logger.info(
'🎯 CHANNEL CHECK: Payload интерпретирован как реферальный код',
' CHANNEL CHECK: Реферальный код уже сохранен в state: %s',
state_data.get('referral_code') or f'campaign_id={state_data.get("campaign_id")}',
)
await state.set_data(state_data)
@@ -1829,6 +1858,12 @@ async def required_sub_channel_check(
except Exception as e:
logger.warning(f'Не удалось удалить сообщение: {e}')
# ИСПРАВЛЕНИЕ БАГА: Очищаем Redis payload ТОЛЬКО после успешной проверки подписки
# и перед показом главного меню или завершением регистрации
if pending_start_payload:
await delete_pending_payload_from_redis(query.from_user.id)
logger.info('🗑️ CHANNEL CHECK: Redis payload удален после успешной проверки подписки')
if user and user.status != UserStatus.DELETED.value:
has_active_subscription, subscription_is_active = _calculate_subscription_flags(user.subscription)
@@ -1911,6 +1946,11 @@ async def required_sub_channel_check(
)
await db.refresh(user, ['subscription'])
# ИСПРАВЛЕНИЕ БАГА: Очищаем pending_start_payload из state после создания пользователя
state_data.pop('pending_start_payload', None)
await state.set_data(state_data)
logger.info('✅ CHANNEL CHECK: pending_start_payload удален из state после создания пользователя')
# Обрабатываем реферальную регистрацию
if referrer_id:
try:
+6 -24
View File
@@ -342,36 +342,18 @@ def get_admin_users_filters_keyboard(language: str = 'ru') -> InlineKeyboardMark
callback_data='admin_users_balance_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_TRAFFIC', '📶 По трафику'),
callback_data='admin_users_traffic_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_ACTIVITY', '🕒 По активности'),
callback_data='admin_users_activity_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_SPENDING', '💳 По сумме трат'),
callback_data='admin_users_spending_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_PURCHASES', '🛒 По количеству покупок'),
callback_data='admin_users_purchases_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_RENEW_READY', '♻️ Готовы к продлению'),
callback_data='admin_users_ready_to_renew_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS', '💰 Потенциальные клиенты'),
callback_data='admin_users_potential_customers_filter',
)
],
[
InlineKeyboardButton(
text=_t(texts, 'ADMIN_USERS_FILTER_CAMPAIGN', '📢 По кампании'),
+4
View File
@@ -807,6 +807,10 @@
"ADMIN_USERS_FILTER_RENEW_READY_TITLE": "♻️ Пользователи готовы к продлению",
"ADMIN_USERS_FILTER_RENEW_READY_DESC": "Подписка истекла, а на балансе осталось {amount} или больше.",
"ADMIN_USERS_FILTER_RENEW_READY_EMPTY": "Сейчас нет пользователей, которые подходят под этот фильтр.",
"ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS": "💰 Потенциальные клиенты",
"ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS_TITLE": "💰 Потенциальные клиенты",
"ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS_DESC": "Нет подписки, но баланс достаточен для покупки.",
"ADMIN_USERS_FILTER_POTENTIAL_CUSTOMERS_EMPTY": "Нет пользователей без подписки с достаточным балансом.",
"ADMIN_USERS_FILTER_CAMPAIGN": "📢 По кампании",
"ADMIN_USERS_FILTER_PURCHASES": "🛒 По количеству покупок",
"ADMIN_USERS_FILTER_SPENDING": "💳 По сумме трат",
+2 -2
View File
@@ -396,10 +396,10 @@ class ChannelCheckerMiddleware(BaseMiddleware):
if not user or not user.subscription:
return
# НЕ реактивируем подписку заблокированных пользователей
# НЕ реактивируем подписку заблокированным пользователям
if user.status == UserStatus.BLOCKED.value:
logger.info(
'🚫 Пропуск реактивации подписки для заблокированного пользователя %s',
'🚫 Пропуск реактивации для заблокированного пользователя %s',
telegram_id,
)
return
+2 -2
View File
@@ -246,7 +246,7 @@ class KassaAiService:
}
params['signature'] = self._generate_hmac_signature(params)
logger.debug(f'KassaAI get_order_status: order_id={order_id}')
logger.info(f'KassaAI get_order_status: order_id={order_id}')
try:
async with (
@@ -259,7 +259,7 @@ class KassaAiService:
) as response,
):
text = await response.text()
logger.debug(f'KassaAI get_order_status response: {text}')
logger.info(f'KassaAI get_order_status response: {text}')
return await response.json()
except aiohttp.ClientError as e:
logger.exception(f'KassaAI API connection error: {e}')
+95 -2
View File
@@ -69,8 +69,13 @@ class KassaAiPaymentMixin:
)
return None
# Генерируем уникальный order_id
order_id = f'kai_{user_id}_{uuid.uuid4().hex[:12]}'
# Получаем telegram_id пользователя для order_id
payment_module = import_module('app.services.payment_service')
user = await payment_module.get_user_by_id(db, user_id)
tg_id = user.telegram_id if user else user_id
# Генерируем уникальный order_id с telegram_id для удобного поиска
order_id = f'k{tg_id}_{uuid.uuid4().hex[:6]}'
amount_rubles = amount_kopeks / 100
currency = settings.KASSA_AI_CURRENCY
@@ -487,3 +492,91 @@ class KassaAiPaymentMixin:
except Exception as e:
logger.exception('KassaAI: ошибка проверки статуса: %s', e)
return None
async def get_kassa_ai_payment_status(
self,
db: AsyncSession,
local_payment_id: int,
) -> dict[str, Any] | None:
"""
Проверяет статус платежа KassaAI по локальному ID через API.
Если платёж оплачен автоматически начисляет баланс.
"""
logger.info('KassaAI: checking payment status for id=%s', local_payment_id)
kassa_ai_crud = import_module('app.database.crud.kassa_ai')
payment = await kassa_ai_crud.get_kassa_ai_payment_by_id(db, local_payment_id)
if not payment:
logger.warning('KassaAI payment not found: id=%s', local_payment_id)
return None
if payment.is_paid:
return {
'payment': payment,
'status': 'success',
'is_paid': True,
}
if not settings.KASSA_AI_API_KEY:
return {
'payment': payment,
'status': payment.status or 'pending',
'is_paid': payment.is_paid,
}
try:
# Запрашиваем статус заказа в KassaAI (api.fk.life)
response = await kassa_ai_service.get_order_status(payment.order_id)
# KassaAI возвращает список заказов (как Freekassa)
orders = response.get('orders', [])
target_order = None
# Ищем наш заказ в списке
for order in orders:
order_key = str(order.get('merchant_order_id') or order.get('paymentId'))
if order_key == str(payment.order_id):
target_order = order
break
if target_order:
# Статус 1 = Оплачен (как в Freekassa)
kai_status = int(target_order.get('status', 0))
if kai_status == 1:
logger.info('KassaAI payment %s confirmed via API', payment.order_id)
callback_payload = {
'check_source': 'api',
'kai_order_data': target_order,
}
# ID заказа на стороне KassaAI
kai_intid = str(target_order.get('fk_order_id') or target_order.get('id'))
# Обновляем статус
payment = await kassa_ai_crud.update_kassa_ai_payment_status(
db=db,
payment=payment,
status='success',
is_paid=True,
kassa_ai_order_id=kai_intid,
payment_system_id=int(target_order.get('curID')) if target_order.get('curID') else None,
callback_payload=callback_payload,
)
# Финализируем (начисляем баланс)
await self._finalize_kassa_ai_payment(
db,
payment,
intid=kai_intid,
trigger='api_check',
)
except Exception as e:
logger.error('Error checking KassaAI payment status: %s', e)
return {
'payment': payment,
'status': payment.status or 'pending',
'is_paid': payment.is_paid,
}
+64 -28
View File
@@ -164,42 +164,78 @@ async def ensure_payment_method_configs(db: AsyncSession) -> None:
"""Initialize payment method configs if they don't exist yet.
Called on startup to seed defaults from env vars.
Also adds any missing methods that were added after initial setup.
"""
count_result = await db.execute(select(func.count()).select_from(PaymentMethodConfig))
count = count_result.scalar() or 0
# Get existing method IDs
existing_result = await db.execute(select(PaymentMethodConfig.method_id))
existing_method_ids = set(existing_result.scalars().all())
if count > 0:
return # Already initialized
if not existing_method_ids:
# First-time initialization
logger.info('Initializing payment method configurations from env vars...')
defaults = _get_method_defaults()
logger.info('Initializing payment method configurations from env vars...')
for idx, method_id in enumerate(DEFAULT_METHOD_ORDER):
method_def = defaults.get(method_id, {})
is_configured = method_def.get('is_configured', False)
sub_options = None
available = method_def.get('available_sub_options')
if available:
# Enable all sub-options by default
sub_options = {opt['id']: True for opt in available}
config = PaymentMethodConfig(
method_id=method_id,
sort_order=idx,
is_enabled=is_configured,
display_name=None,
sub_options=sub_options,
min_amount_kopeks=None,
max_amount_kopeks=None,
user_type_filter='all',
first_topup_filter='any',
promo_group_filter_mode='all',
)
db.add(config)
await db.commit()
logger.info(f'Payment method configurations initialized ({len(DEFAULT_METHOD_ORDER)} methods).')
return
# Add missing methods (for cases when new methods are added to code)
defaults = _get_method_defaults()
missing_methods = [m for m in DEFAULT_METHOD_ORDER if m not in existing_method_ids]
for idx, method_id in enumerate(DEFAULT_METHOD_ORDER):
method_def = defaults.get(method_id, {})
is_configured = method_def.get('is_configured', False)
sub_options = None
available = method_def.get('available_sub_options')
if available:
# Enable all sub-options by default
sub_options = {opt['id']: True for opt in available}
if missing_methods:
logger.info(f'Adding missing payment methods: {missing_methods}')
# Get max sort_order to append new methods at the end
max_order_result = await db.execute(select(func.max(PaymentMethodConfig.sort_order)))
max_order = max_order_result.scalar() or 0
config = PaymentMethodConfig(
method_id=method_id,
sort_order=idx,
is_enabled=is_configured,
display_name=None,
sub_options=sub_options,
min_amount_kopeks=None,
max_amount_kopeks=None,
user_type_filter='all',
first_topup_filter='any',
promo_group_filter_mode='all',
)
db.add(config)
for idx, method_id in enumerate(missing_methods, start=max_order + 1):
method_def = defaults.get(method_id, {})
is_configured = method_def.get('is_configured', False)
sub_options = None
available = method_def.get('available_sub_options')
if available:
sub_options = {opt['id']: True for opt in available}
await db.commit()
logger.info(f'Payment method configurations initialized ({len(DEFAULT_METHOD_ORDER)} methods).')
config = PaymentMethodConfig(
method_id=method_id,
sort_order=idx,
is_enabled=is_configured,
display_name=None,
sub_options=sub_options,
min_amount_kopeks=None,
max_amount_kopeks=None,
user_type_filter='all',
first_topup_filter='any',
promo_group_filter_mode='all',
)
db.add(config)
await db.commit()
logger.info(f'Added {len(missing_methods)} missing payment method(s).')
# ============ CRUD ============
@@ -71,6 +71,7 @@ SUPPORTED_MANUAL_CHECK_METHODS: frozenset[PaymentMethod] = frozenset(
PaymentMethod.PLATEGA,
PaymentMethod.CLOUDPAYMENTS,
PaymentMethod.FREEKASSA,
PaymentMethod.KASSA_AI,
}
)
@@ -87,6 +88,7 @@ SUPPORTED_AUTO_CHECK_METHODS: frozenset[PaymentMethod] = frozenset(
# WATA removed - API returns 429 "Use webhook polling is rate-limited".
# Payments are processed via webhook (wata_webhook.py).
PaymentMethod.FREEKASSA,
PaymentMethod.KASSA_AI,
}
)
@@ -955,6 +957,9 @@ async def run_manual_check(
elif method == PaymentMethod.FREEKASSA:
result = await payment_service.get_freekassa_payment_status(db, local_payment_id)
payment = result.get('payment') if result else None
elif method == PaymentMethod.KASSA_AI:
result = await payment_service.get_kassa_ai_payment_status(db, local_payment_id)
payment = result.get('payment') if result else None
else:
logger.warning('Manual check requested for unsupported method %s', method)
return None
File diff suppressed because it is too large Load Diff
+67 -1
View File
@@ -3,7 +3,7 @@ from datetime import datetime, timedelta
from typing import Any
from aiogram import Bot, types
from sqlalchemy import delete, func, select, update
from sqlalchemy import delete, func, or_, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import selectinload
@@ -345,6 +345,72 @@ class UserService:
'total_count': 0,
}
async def get_potential_customers(
self,
db: AsyncSession,
min_balance_kopeks: int,
page: int = 1,
limit: int = 10,
) -> dict[str, Any]:
"""Возвращает пользователей без активной подписки с достаточным балансом."""
try:
offset = (page - 1) * limit
# Фильтры: нет активной подписки И баланс >= порога
base_filters = [
User.balance_kopeks >= min_balance_kopeks,
]
# Основной запрос с LEFT JOIN для поддержки пользователей без подписки
query = (
select(User)
.options(selectinload(User.subscription))
.outerjoin(Subscription, Subscription.user_id == User.id)
.where(
*base_filters,
or_(
User.subscription == None,
~Subscription.status.in_(['active', 'trial']),
),
)
.order_by(User.balance_kopeks.desc(), User.created_at.desc())
.offset(offset)
.limit(limit)
)
result = await db.execute(query)
users = result.scalars().unique().all()
# Запрос для подсчета общего количества
count_query = (
select(func.count(User.id))
.outerjoin(Subscription, Subscription.user_id == User.id)
.where(
*base_filters,
or_(
User.subscription == None,
~Subscription.status.in_(['active', 'trial']),
),
)
)
total_count = (await db.execute(count_query)).scalar() or 0
total_pages = (total_count + limit - 1) // limit if total_count else 0
return {
'users': users,
'current_page': page,
'total_pages': total_pages,
'total_count': total_count,
}
except Exception as e:
logger.error(f'Ошибка получения потенциальных клиентов: {e}')
return {
'users': [],
'current_page': 1,
'total_pages': 1,
'total_count': 0,
}
async def get_user_spending_stats_map(self, db: AsyncSession, user_ids: list[int]) -> dict[int, dict[str, int]]:
try:
return await get_users_spending_stats(db, user_ids)
+8 -4
View File
@@ -118,6 +118,9 @@ class AdminStates(StatesGroup):
adding_virtual_participant_name = State()
adding_virtual_participant_count = State()
editing_virtual_participant_count = State()
# Массовое создание виртуальных участников (массовка)
adding_mass_virtual_count = State() # Сколько призраков создать
adding_mass_virtual_referrals = State() # По сколько рефералов у каждого
editing_daily_contest_field = State()
editing_daily_contest_value = State()
@@ -132,6 +135,10 @@ class AdminStates(StatesGroup):
# Тестовое начисление реферального дохода
test_referral_earning_input = State()
# Диагностика рефералов
referral_diagnostics_period = State()
waiting_for_log_file = State()
editing_rules_page = State()
editing_privacy_policy = State()
editing_public_offer = State()
@@ -173,12 +180,9 @@ class AdminStates(StatesGroup):
# Состояния для отслеживания источника перехода
viewing_user_from_balance_list = State()
viewing_user_from_traffic_list = State()
viewing_user_from_last_activity_list = State()
viewing_user_from_spending_list = State()
viewing_user_from_purchases_list = State()
viewing_user_from_campaign_list = State()
viewing_user_from_ready_to_renew_list = State()
viewing_user_from_potential_customers_list = State()
# Состояния для управления тарифами
creating_tariff_name = State()
+153
View File
@@ -0,0 +1,153 @@
"""
Тесты для сервиса диагностики реферальной системы.
"""
import tempfile
from datetime import datetime, timedelta
from pathlib import Path
import pytest
from app.services.referral_diagnostics_service import ReferralDiagnosticsService
@pytest.fixture
def temp_log_file():
"""Создаёт временный лог-файл для тестов."""
with tempfile.NamedTemporaryFile(mode='w', suffix='.log', delete=False) as f:
yield Path(f.name)
# Cleanup
Path(f.name).unlink(missing_ok=True)
@pytest.fixture
def sample_log_content():
"""Пример содержимого лог-файла с реферальными событиями."""
today = datetime.now().strftime('%Y-%m-%d')
return f"""
{today} 10:00:00,123 - app.handlers.start - INFO - 🔎 Найден реферальный код: <ABC123>
{today} 10:00:05,456 - app.handlers.start - INFO - Реферальный код ABC123 применен для пользователя 123456789
{today} 10:00:10,789 - app.services.referral_service - INFO - Реферальная регистрация обработана для 123456789
{today} 10:00:15,012 - app.services.referral_service - INFO - 💰 Реферал 123456789 получил бонус
{today} 11:00:00,345 - app.handlers.start - INFO - 🔎 Найден реферальный код: <XYZ999>
{today} 11:00:05,678 - app.handlers.start - INFO - Реферальный код XYZ999 применен для пользователя 987654321
{today} 12:00:00,901 - app.handlers.start - INFO - 🔎 Найден реферальный код: <TEST777>
{today} 13:00:00,234 - unrelated module - INFO - Some other log message
"""
@pytest.mark.asyncio
async def test_parse_logs_basic(temp_log_file, sample_log_content):
"""Тест базового парсинга логов."""
# Записываем тестовые данные в файл
temp_log_file.write_text(sample_log_content)
service = ReferralDiagnosticsService(log_path=str(temp_log_file))
today = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
tomorrow = today + timedelta(days=1)
events = await service._parse_logs(today, tomorrow)
# Проверяем что нашлись все события
assert len(events) >= 6, f'Expected at least 6 events, found {len(events)}'
# Проверяем типы событий
event_types = [e.event_type for e in events]
assert 'code_found' in event_types
assert 'code_applied' in event_types
assert 'registration_processed' in event_types
assert 'bonus_given' in event_types
@pytest.mark.asyncio
async def test_analyze_period_with_issues(temp_log_file, sample_log_content):
"""Тест анализа с проблемными случаями."""
temp_log_file.write_text(sample_log_content)
service = ReferralDiagnosticsService(log_path=str(temp_log_file))
today = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
tomorrow = today + timedelta(days=1)
# Используем None вместо db для базового теста парсинга
from unittest.mock import AsyncMock
mock_db = AsyncMock()
mock_db.execute.return_value.scalar_one_or_none.return_value = None
report = await service.analyze_period(mock_db, today, tomorrow)
# Проверяем статистику
# Примечание: code_found не имеет telegram_id, поэтому total_link_clicks будет 0
# Это нормально - мы считаем только события с telegram_id
assert report.total_codes_applied >= 1, 'Should have applied codes'
# Проверяем что нашлись проблемные случаи
# (987654321 применил код, но не завершил регистрацию)
assert 987654321 in report.users_applied_no_registration, (
f'Expected 987654321 in problems, got: {report.users_applied_no_registration}'
)
@pytest.mark.asyncio
async def test_empty_log_file(temp_log_file):
"""Тест работы с пустым лог-файлом."""
temp_log_file.write_text('')
service = ReferralDiagnosticsService(log_path=str(temp_log_file))
today = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
tomorrow = today + timedelta(days=1)
from unittest.mock import AsyncMock
mock_db = AsyncMock()
report = await service.analyze_period(mock_db, today, tomorrow)
# Проверяем что отчёт пустой
assert report.total_link_clicks == 0
assert report.total_codes_applied == 0
assert report.total_registrations == 0
assert len(report.events) == 0
@pytest.mark.asyncio
async def test_nonexistent_log_file():
"""Тест работы с несуществующим лог-файлом."""
service = ReferralDiagnosticsService(log_path='/nonexistent/path/to/log.log')
today = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
tomorrow = today + timedelta(days=1)
from unittest.mock import AsyncMock
mock_db = AsyncMock()
# Не должно быть исключений
report = await service.analyze_period(mock_db, today, tomorrow)
assert report.total_link_clicks == 0
assert len(report.events) == 0
@pytest.mark.asyncio
async def test_analyze_today(temp_log_file, sample_log_content):
"""Тест метода analyze_today."""
temp_log_file.write_text(sample_log_content)
service = ReferralDiagnosticsService(log_path=str(temp_log_file))
from unittest.mock import AsyncMock
mock_db = AsyncMock()
report = await service.analyze_today(mock_db)
# Проверяем что период установлен корректно
today = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
assert report.analysis_period_start.date() == today.date()