From 0df3018703799dcf47a2d845ed9611a004e339dc Mon Sep 17 00:00:00 2001 From: gy9vin Date: Thu, 25 Dec 2025 23:01:49 +0300 Subject: [PATCH] =?UTF-8?q?feat(nalogo):=20=D1=81=D0=B8=D1=81=D1=82=D0=B5?= =?UTF-8?q?=D0=BC=D0=B0=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D0=B8=20?= =?UTF-8?q?=D1=87=D0=B5=D0=BA=D0=BE=D0=B2=20=D1=81=20=D0=BE=D1=82=D0=BB?= =?UTF-8?q?=D0=BE=D0=B6=D0=B5=D0=BD=D0=BD=D0=BE=D0=B9=20=D0=BE=D1=82=D0=BF?= =?UTF-8?q?=D1=80=D0=B0=D0=B2=D0=BA=D0=BE=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Реализована отказоустойчивая система отправки чеков в налоговую: - Добавлен NalogoQueueService для фоновой обработки очереди чеков - При недоступности nalog.ru (503) чеки сохраняются в Redis - Автоматическая повторная отправка с настраиваемым интервалом - Защита от DDoS: задержка между чеками (NALOGO_QUEUE_RECEIPT_DELAY) - Уведомления админам в топик при проблемах и успешной разгрузке Изменения в файлах: - app/services/nalogo_queue_service.py: новый фоновый сервис - app/services/nalogo_service.py: методы очереди, определение 503 - app/utils/cache.py: lpush/rpop/llen/lrange для Redis List - app/handlers/admin/monitoring.py: статистика чеков в админке - app/config.py: NALOGO_QUEUE_* и ADMIN_NOTIFICATIONS_NALOG_TOPIC_ID - main.py: интеграция запуска/остановки сервиса Новые ENV переменные: - ADMIN_NOTIFICATIONS_NALOG_TOPIC_ID - NALOGO_QUEUE_CHECK_INTERVAL (300с) - NALOGO_QUEUE_RECEIPT_DELAY (3с) - NALOGO_QUEUE_MAX_ATTEMPTS (10) --- .env.example | 4 + app/config.py | 6 + app/handlers/admin/monitoring.py | 138 ++++++++++++- app/services/nalogo_queue_service.py | 285 +++++++++++++++++++++++++++ app/services/nalogo_service.py | 93 ++++++++- app/services/payment/yookassa.py | 6 +- app/utils/cache.py | 52 ++++- main.py | 33 ++++ 8 files changed, 606 insertions(+), 11 deletions(-) create mode 100644 app/services/nalogo_queue_service.py diff --git a/.env.example b/.env.example index b92c2722..66d36d5a 100644 --- a/.env.example +++ b/.env.example @@ -14,6 +14,7 @@ ADMIN_NOTIFICATIONS_ENABLED=true ADMIN_NOTIFICATIONS_CHAT_ID=-1001234567890 # Замени на ID твоего канала (-100) - ПРЕФИКС ЗАКРЫТОГО КАНАЛА! ВСТАВИТЬ СВОЙ ID СРАЗУ ПОСЛЕ (-100) БЕЗ ПРОБЕЛОВ! ADMIN_NOTIFICATIONS_TOPIC_ID=123 # Опционально: ID топика ADMIN_NOTIFICATIONS_TICKET_TOPIC_ID=126 # Опционально: ID топика для тикетов +ADMIN_NOTIFICATIONS_NALOG_TOPIC_ID=133 # Опционально: ID топика для уведомлений о чеках NaloGO # Автоматические отчеты ADMIN_REPORTS_ENABLED=false ADMIN_REPORTS_CHAT_ID= # Опционально: чат для отчетов (по умолчанию ADMIN_NOTIFICATIONS_CHAT_ID) @@ -317,6 +318,9 @@ NALOGO_INN= # ИНН самозанятого NALOGO_PASSWORD= # Пароль от личного кабинета налоговой NALOGO_DEVICE_ID= # Опционально: ID устройства для авторизации NALOGO_STORAGE_PATH=./nalogo_tokens.json # Путь к файлу с токенами +NALOGO_QUEUE_CHECK_INTERVAL=300 # Интервал проверки очереди чеков (секунды) +NALOGO_QUEUE_RECEIPT_DELAY=3 # Задержка между отправкой чеков (секунды) +NALOGO_QUEUE_MAX_ATTEMPTS=10 # Максимум попыток отправки одного чека # ===== НАСТРОЙКИ ОПИСАНИЙ ПЛАТЕЖЕЙ ===== # Эти настройки позволяют изменить описания платежей, diff --git a/app/config.py b/app/config.py index eaa75cf5..1beebf6d 100644 --- a/app/config.py +++ b/app/config.py @@ -45,6 +45,12 @@ class Settings(BaseSettings): ADMIN_NOTIFICATIONS_CHAT_ID: Optional[str] = None ADMIN_NOTIFICATIONS_TOPIC_ID: Optional[int] = None ADMIN_NOTIFICATIONS_TICKET_TOPIC_ID: Optional[int] = None + ADMIN_NOTIFICATIONS_NALOG_TOPIC_ID: Optional[int] = None + + # Настройки очереди чеков NaloGO + NALOGO_QUEUE_CHECK_INTERVAL: int = 300 # Интервал проверки очереди (секунды) + NALOGO_QUEUE_RECEIPT_DELAY: int = 3 # Задержка между отправкой чеков (секунды) + NALOGO_QUEUE_MAX_ATTEMPTS: int = 10 # Максимум попыток отправки чека ADMIN_REPORTS_ENABLED: bool = False ADMIN_REPORTS_CHAT_ID: Optional[str] = None diff --git a/app/handlers/admin/monitoring.py b/app/handlers/admin/monitoring.py index d7dc7444..7c45e1f7 100644 --- a/app/handlers/admin/monitoring.py +++ b/app/handlers/admin/monitoring.py @@ -10,6 +10,7 @@ from aiogram.exceptions import TelegramBadRequest from app.config import settings from app.database.database import get_db from app.services.monitoring_service import monitoring_service +from app.services.nalogo_queue_service import nalogo_queue_service from app.utils.decorators import admin_required from app.utils.pagination import paginate_list from app.keyboards.admin import get_monitoring_keyboard, get_admin_main_keyboard @@ -911,11 +912,36 @@ async def monitoring_statistics_callback(callback: CallbackQuery): • Уведомления: {'🟢 Вкл' if getattr(settings, 'ENABLE_NOTIFICATIONS', True) else '🔴 Выкл'} • Автооплата: {', '.join(map(str, settings.get_autopay_warning_days()))} дней """ - + + # Добавляем информацию о чеках NaloGO + if settings.is_nalogo_enabled(): + nalogo_status = await nalogo_queue_service.get_status() + queue_len = nalogo_status.get("queue_length", 0) + total_amount = nalogo_status.get("total_amount", 0) + running = nalogo_status.get("running", False) + + nalogo_section = f""" +🧾 Чеки NaloGO: +• Сервис: {'🟢 Работает' if running else '🔴 Остановлен'} +• В очереди: {queue_len} чек(ов)""" + if queue_len > 0: + nalogo_section += f"\n• На сумму: {total_amount:,.2f} ₽" + text += nalogo_section + from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton - keyboard = InlineKeyboardMarkup(inline_keyboard=[ - [InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_monitoring")] - ]) + + buttons = [] + # Кнопка обработки очереди чеков если есть что обрабатывать + if settings.is_nalogo_enabled(): + nalogo_status = await nalogo_queue_service.get_status() + if nalogo_status.get("queue_length", 0) > 0: + buttons.append([InlineKeyboardButton( + text=f"🧾 Отправить чеки ({nalogo_status['queue_length']} шт.)", + callback_data="admin_mon_nalogo_force_process" + )]) + + buttons.append([InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_monitoring")]) + keyboard = InlineKeyboardMarkup(inline_keyboard=buttons) await callback.message.edit_text(text, parse_mode="HTML", reply_markup=keyboard) break @@ -925,6 +951,110 @@ async def monitoring_statistics_callback(callback: CallbackQuery): await callback.answer(f"❌ Ошибка получения статистики: {str(e)}", show_alert=True) +@router.callback_query(F.data == "admin_mon_nalogo_force_process") +@admin_required +async def nalogo_force_process_callback(callback: CallbackQuery): + """Принудительная отправка чеков из очереди.""" + try: + await callback.answer("🔄 Запускаю обработку очереди чеков...", show_alert=False) + + result = await nalogo_queue_service.force_process() + + if "error" in result: + await callback.answer(f"❌ {result['error']}", show_alert=True) + return + + message = result.get("message", "Готово") + processed = result.get("processed", 0) + remaining = result.get("remaining", 0) + + if processed > 0: + text = f"✅ Обработано: {processed} чек(ов)" + if remaining > 0: + text += f"\n⏳ Осталось в очереди: {remaining}" + else: + if remaining > 0: + text = f"⚠️ Сервис nalog.ru недоступен\n⏳ В очереди: {remaining} чек(ов)" + else: + text = "📭 Очередь пуста" + + await callback.answer(text, show_alert=True) + + # Обновляем страницу статистики + from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton + + # Перезагружаем статистику + async for db in get_db(): + from app.database.crud.subscription import get_subscriptions_statistics + sub_stats = await get_subscriptions_statistics(db) + mon_status = await monitoring_service.get_monitoring_status(db) + + week_ago = datetime.now() - timedelta(days=7) + week_logs = await monitoring_service.get_monitoring_logs(db, limit=1000) + week_logs = [log for log in week_logs if log['created_at'] >= week_ago] + week_success = sum(1 for log in week_logs if log['is_success']) + week_errors = len(week_logs) - week_success + + stats_text = f""" +📊 Статистика мониторинга + +📱 Подписки: +• Всего: {sub_stats['total_subscriptions']} +• Активных: {sub_stats['active_subscriptions']} +• Тестовых: {sub_stats['trial_subscriptions']} +• Платных: {sub_stats['paid_subscriptions']} + +📈 За сегодня: +• Успешных операций: {mon_status['stats_24h']['successful']} +• Ошибок: {mon_status['stats_24h']['failed']} +• Успешность: {mon_status['stats_24h']['success_rate']}% + +📊 За неделю: +• Всего событий: {len(week_logs)} +• Успешных: {week_success} +• Ошибок: {week_errors} +• Успешность: {round(week_success/len(week_logs)*100, 1) if week_logs else 0}% + +🔧 Система: +• Интервал: {settings.MONITORING_INTERVAL} мин +• Уведомления: {'🟢 Вкл' if getattr(settings, 'ENABLE_NOTIFICATIONS', True) else '🔴 Выкл'} +• Автооплата: {', '.join(map(str, settings.get_autopay_warning_days()))} дней +""" + + if settings.is_nalogo_enabled(): + nalogo_status = await nalogo_queue_service.get_status() + queue_len = nalogo_status.get("queue_length", 0) + total_amount = nalogo_status.get("total_amount", 0) + running = nalogo_status.get("running", False) + + nalogo_section = f""" +🧾 Чеки NaloGO: +• Сервис: {'🟢 Работает' if running else '🔴 Остановлен'} +• В очереди: {queue_len} чек(ов)""" + if queue_len > 0: + nalogo_section += f"\n• На сумму: {total_amount:,.2f} ₽" + stats_text += nalogo_section + + buttons = [] + if settings.is_nalogo_enabled(): + nalogo_status = await nalogo_queue_service.get_status() + if nalogo_status.get("queue_length", 0) > 0: + buttons.append([InlineKeyboardButton( + text=f"🧾 Отправить чеки ({nalogo_status['queue_length']} шт.)", + callback_data="admin_mon_nalogo_force_process" + )]) + + buttons.append([InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_monitoring")]) + keyboard = InlineKeyboardMarkup(inline_keyboard=buttons) + + await callback.message.edit_text(stats_text, parse_mode="HTML", reply_markup=keyboard) + break + + except Exception as e: + logger.error(f"Ошибка принудительной обработки чеков: {e}") + await callback.answer(f"❌ Ошибка: {str(e)}", show_alert=True) + + def get_monitoring_logs_keyboard(current_page: int, total_pages: int): from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton diff --git a/app/services/nalogo_queue_service.py b/app/services/nalogo_queue_service.py new file mode 100644 index 00000000..57fb3813 --- /dev/null +++ b/app/services/nalogo_queue_service.py @@ -0,0 +1,285 @@ +"""Фоновый сервис для обработки очереди чеков NaloGO. + +При временной недоступности сервиса nalog.ru (503), чеки сохраняются в Redis +и отправляются позже этим сервисом. +""" + +import asyncio +import logging +from datetime import datetime, timedelta +from typing import Optional + +from aiogram import Bot + +from app.config import settings +from app.services.nalogo_service import NaloGoService + +logger = logging.getLogger(__name__) + + +class NalogoQueueService: + """Сервис фоновой обработки очереди чеков NaloGO.""" + + def __init__(self, nalogo_service: Optional[NaloGoService] = None): + self._nalogo_service = nalogo_service + self._bot: Optional[Bot] = None + self._task: Optional[asyncio.Task] = None + self._running = False + self._last_notification_time: Optional[datetime] = None + self._notification_cooldown = timedelta(hours=1) # Не чаще раза в час + self._had_pending_receipts = False # Флаг для отслеживания успешной разгрузки + + def set_nalogo_service(self, service: NaloGoService) -> None: + """Установить сервис NaloGO.""" + self._nalogo_service = service + + def set_bot(self, bot: Bot) -> None: + """Установить бота для отправки уведомлений.""" + self._bot = bot + + def is_running(self) -> bool: + """Проверка, запущен ли сервис.""" + return self._running and self._task is not None and not self._task.done() + + @property + def _check_interval(self) -> int: + """Интервал проверки очереди в секундах.""" + return getattr(settings, "NALOGO_QUEUE_CHECK_INTERVAL", 300) + + @property + def _receipt_delay(self) -> int: + """Задержка между отправкой чеков в секундах.""" + return getattr(settings, "NALOGO_QUEUE_RECEIPT_DELAY", 3) + + @property + def _max_attempts(self) -> int: + """Максимальное количество попыток отправки чека.""" + return getattr(settings, "NALOGO_QUEUE_MAX_ATTEMPTS", 10) + + async def start(self) -> None: + """Запустить фоновую обработку очереди.""" + if not self._nalogo_service or not self._nalogo_service.configured: + logger.info("NaloGO не настроен, сервис очереди чеков не запущен") + return + + if self.is_running(): + logger.warning("Сервис очереди чеков уже запущен") + return + + self._running = True + self._task = asyncio.create_task(self._process_queue_loop()) + logger.info( + f"Сервис очереди чеков NaloGO запущен " + f"(интервал: {self._check_interval}с, задержка между чеками: {self._receipt_delay}с)" + ) + + async def stop(self) -> None: + """Остановить фоновую обработку.""" + self._running = False + if self._task and not self._task.done(): + self._task.cancel() + try: + await self._task + except asyncio.CancelledError: + pass + self._task = None + logger.info("Сервис очереди чеков NaloGO остановлен") + + async def _send_admin_notification(self, message: str, skip_cooldown: bool = False) -> None: + """Отправить уведомление админам о чеках.""" + if not self._bot: + return + + chat_id = settings.get_admin_notifications_chat_id() + if not chat_id: + return + + topic_id = settings.ADMIN_NOTIFICATIONS_NALOG_TOPIC_ID + + # Проверяем cooldown (можно пропустить для важных уведомлений) + if not skip_cooldown: + now = datetime.now() + if self._last_notification_time: + if now - self._last_notification_time < self._notification_cooldown: + logger.debug("Уведомление о чеках пропущено (cooldown)") + return + + try: + await self._bot.send_message( + chat_id=chat_id, + message_thread_id=topic_id, + text=message, + parse_mode="HTML", + ) + self._last_notification_time = datetime.now() + logger.info("Отправлено уведомление о чеках NaloGO") + except Exception as error: + logger.error(f"Ошибка отправки уведомления о чеках: {error}") + + async def _process_queue_loop(self) -> None: + """Основной цикл обработки очереди.""" + while self._running: + try: + await self._process_pending_receipts() + except Exception as error: + logger.error(f"Ошибка в цикле обработки очереди чеков: {error}") + + await asyncio.sleep(self._check_interval) + + async def _process_pending_receipts(self) -> None: + """Обработать все ожидающие чеки в очереди.""" + if not self._nalogo_service: + return + + queue_length = await self._nalogo_service.get_queue_length() + if queue_length == 0: + return + + logger.info(f"Начинаем обработку очереди чеков: {queue_length} шт.") + self._had_pending_receipts = True + + processed = 0 + failed = 0 + skipped = 0 + total_processed_amount = 0.0 + service_unavailable = False + + while True: + receipt_data = await self._nalogo_service.pop_receipt_from_queue() + if not receipt_data: + break + + attempts = receipt_data.get("attempts", 0) + payment_id = receipt_data.get("payment_id", "unknown") + amount = receipt_data.get("amount", 0) + + # Проверяем количество попыток + if attempts >= self._max_attempts: + logger.error( + f"Чек {payment_id} превысил лимит попыток ({self._max_attempts}), " + f"удален из очереди" + ) + skipped += 1 + continue + + # Пытаемся отправить чек + try: + receipt_uuid = await self._nalogo_service.create_receipt( + name=receipt_data.get("name", "Интернет-сервис - Пополнение баланса"), + amount=amount, + quantity=receipt_data.get("quantity", 1), + client_info=receipt_data.get("client_info"), + payment_id=payment_id, + queue_on_failure=False, # Не добавлять в очередь повторно автоматически + ) + + if receipt_uuid: + processed += 1 + total_processed_amount += amount + logger.info( + f"Чек из очереди успешно создан: {receipt_uuid} " + f"(payment_id={payment_id}, попытка {attempts + 1})" + ) + else: + # Вернуть в очередь с увеличенным счетчиком попыток + await self._nalogo_service.requeue_receipt(receipt_data) + failed += 1 + service_unavailable = True + logger.warning( + f"Не удалось создать чек из очереди (payment_id={payment_id}), " + f"возвращен в очередь (попытка {attempts + 1}/{self._max_attempts})" + ) + # Если сервис недоступен, прекращаем попытки до следующего цикла + break + + except Exception as error: + await self._nalogo_service.requeue_receipt(receipt_data) + failed += 1 + logger.error( + f"Ошибка при создании чека из очереди (payment_id={payment_id}): {error}" + ) + # Прекращаем попытки при ошибке + break + + # Задержка между чеками чтобы не долбить API + await asyncio.sleep(self._receipt_delay) + + if processed > 0 or failed > 0 or skipped > 0: + logger.info( + f"Обработка очереди завершена: " + f"успешно={processed}, неудачно={failed}, пропущено={skipped}" + ) + + # Проверяем остаток в очереди + remaining = await self._nalogo_service.get_queue_length() + + # Отправляем уведомление если есть проблемы + if service_unavailable or failed > 0: + if remaining > 0: + queued = await self._nalogo_service.get_queued_receipts() + total_queued_amount = sum(r.get("amount", 0) for r in queued) + + message = ( + f"⚠️ Проблема с отправкой чеков NaloGO\n\n" + f"Сервис nalog.ru временно недоступен.\n\n" + f"📋 В очереди: {remaining} чек(ов)\n" + f"💰 На сумму: {total_queued_amount:,.2f} ₽\n\n" + f"Чеки будут отправлены автоматически когда сервис восстановится." + ) + await self._send_admin_notification(message) + + # Уведомление об успешной разгрузке очереди + elif remaining == 0 and self._had_pending_receipts and processed > 0: + self._had_pending_receipts = False + message = ( + f"✅ Очередь чеков NaloGO разгружена\n\n" + f"Все отложенные чеки успешно отправлены!\n\n" + f"📋 Отправлено: {processed} чек(ов)\n" + f"💰 На сумму: {total_processed_amount:,.2f} ₽" + ) + await self._send_admin_notification(message, skip_cooldown=True) + + async def force_process(self) -> dict: + """Принудительно обработать очередь (для ручного запуска).""" + if not self._nalogo_service: + return {"error": "NaloGO сервис не настроен"} + + queue_length = await self._nalogo_service.get_queue_length() + if queue_length == 0: + return {"message": "Очередь пуста", "processed": 0} + + await self._process_pending_receipts() + new_length = await self._nalogo_service.get_queue_length() + + return { + "message": "Обработка завершена", + "was_in_queue": queue_length, + "remaining": new_length, + "processed": queue_length - new_length, + } + + async def get_status(self) -> dict: + """Получить статус сервиса и очереди.""" + queue_length = 0 + total_amount = 0.0 + queued_receipts = [] + + if self._nalogo_service: + queue_length = await self._nalogo_service.get_queue_length() + if queue_length > 0: + queued_receipts = await self._nalogo_service.get_queued_receipts() + total_amount = sum(r.get("amount", 0) for r in queued_receipts) + + return { + "running": self.is_running(), + "check_interval_seconds": self._check_interval, + "receipt_delay_seconds": self._receipt_delay, + "queue_length": queue_length, + "total_amount": total_amount, + "max_attempts": self._max_attempts, + "queued_receipts": queued_receipts[:10], # Показываем только первые 10 + } + + +# Глобальный экземпляр сервиса +nalogo_queue_service = NalogoQueueService() diff --git a/app/services/nalogo_service.py b/app/services/nalogo_service.py index dacb9c29..6a4d0289 100644 --- a/app/services/nalogo_service.py +++ b/app/services/nalogo_service.py @@ -1,4 +1,5 @@ import logging +from datetime import datetime from typing import Optional, Dict, Any from decimal import Decimal @@ -6,9 +7,12 @@ from nalogo import Client from nalogo.dto.income import IncomeClient, IncomeType from app.config import settings +from app.utils.cache import cache logger = logging.getLogger(__name__) +NALOGO_QUEUE_KEY = "nalogo:receipt_queue" + class NaloGoService: """Сервис для работы с API NaloGO (налоговая служба самозанятых).""" @@ -49,6 +53,45 @@ class NaloGoService: ) self.configured = False + @staticmethod + def _is_service_unavailable(error: Exception) -> bool: + """Проверяет, является ли ошибка временной недоступностью сервиса (503).""" + error_str = str(error).lower() + return ( + "503" in error_str + or "service temporarily unavailable" in error_str + or "service unavailable" in error_str + or "ведутся работы" in error_str + or ("health" in error_str and "false" in error_str) + ) + + async def _queue_receipt( + self, + name: str, + amount: float, + quantity: int, + client_info: Optional[Dict[str, Any]], + payment_id: Optional[str] = None, + ) -> bool: + """Добавить чек в очередь для отложенной отправки.""" + receipt_data = { + "name": name, + "amount": amount, + "quantity": quantity, + "client_info": client_info, + "payment_id": payment_id, + "created_at": datetime.now().isoformat(), + "attempts": 0, + } + success = await cache.lpush(NALOGO_QUEUE_KEY, receipt_data) + if success: + queue_len = await cache.llen(NALOGO_QUEUE_KEY) + logger.info( + f"Чек добавлен в очередь (payment_id={payment_id}, " + f"сумма={amount}₽, в очереди: {queue_len})" + ) + return success + async def authenticate(self) -> bool: """Аутентификация в сервисе NaloGO.""" if not self.configured: @@ -60,10 +103,24 @@ class NaloGoService: logger.info("Успешная аутентификация в NaloGO") return True except Exception as error: - logger.error("Ошибка аутентификации в NaloGO: %s", error, exc_info=True) + if self._is_service_unavailable(error): + logger.warning( + "NaloGO временно недоступен (техработы): %s", + str(error)[:200] + ) + else: + logger.error("Ошибка аутентификации в NaloGO: %s", error, exc_info=True) return False - async def create_receipt(self, name: str, amount: float, quantity: int = 1, client_info: Optional[Dict[str, Any]] = None) -> Optional[str]: + async def create_receipt( + self, + name: str, + amount: float, + quantity: int = 1, + client_info: Optional[Dict[str, Any]] = None, + payment_id: Optional[str] = None, + queue_on_failure: bool = True, + ) -> Optional[str]: """Создание чека о доходе. Args: @@ -71,6 +128,8 @@ class NaloGoService: amount: Сумма в рублях quantity: Количество client_info: Информация о клиенте (опционально) + payment_id: ID платежа для логирования + queue_on_failure: Добавить в очередь при временной недоступности Returns: UUID чека или None при ошибке @@ -84,6 +143,9 @@ class NaloGoService: if not hasattr(self.client, '_access_token') or not self.client._access_token: auth_success = await self.authenticate() if not auth_success: + # Если сервис недоступен — добавляем в очередь + if queue_on_failure: + await self._queue_receipt(name, amount, quantity, client_info, payment_id) return None income_api = self.client.income() @@ -114,5 +176,30 @@ class NaloGoService: return None except Exception as error: - logger.error("Ошибка создания чека в NaloGO: %s", error, exc_info=True) + if self._is_service_unavailable(error): + logger.warning( + "NaloGO временно недоступен, чек будет отправлен позже " + f"(payment_id={payment_id}, сумма={amount}₽)" + ) + if queue_on_failure: + await self._queue_receipt(name, amount, quantity, client_info, payment_id) + else: + logger.error("Ошибка создания чека в NaloGO: %s", error, exc_info=True) return None + + async def get_queue_length(self) -> int: + """Получить количество чеков в очереди.""" + return await cache.llen(NALOGO_QUEUE_KEY) + + async def get_queued_receipts(self) -> list: + """Получить список чеков в очереди (без удаления).""" + return await cache.lrange(NALOGO_QUEUE_KEY) + + async def pop_receipt_from_queue(self) -> Optional[Dict[str, Any]]: + """Извлечь следующий чек из очереди.""" + return await cache.rpop(NALOGO_QUEUE_KEY) + + async def requeue_receipt(self, receipt_data: Dict[str, Any]) -> bool: + """Вернуть чек обратно в очередь (при неудачной отправке).""" + receipt_data["attempts"] = receipt_data.get("attempts", 0) + 1 + return await cache.lpush(NALOGO_QUEUE_KEY, receipt_data) diff --git a/app/services/payment/yookassa.py b/app/services/payment/yookassa.py index beed9443..8f9ad9f0 100644 --- a/app/services/payment/yookassa.py +++ b/app/services/payment/yookassa.py @@ -1029,13 +1029,13 @@ class YooKassaPaymentMixin: receipt_uuid = await self.nalogo_service.create_receipt( name=receipt_name, amount=amount_rubles, - quantity=1 + quantity=1, + payment_id=payment.yookassa_payment_id, ) if receipt_uuid: logger.info(f"Чек NaloGO создан для платежа {payment.yookassa_payment_id}: {receipt_uuid}") - else: - logger.warning(f"Не удалось создать чек NaloGO для платежа {payment.yookassa_payment_id}") + # При временной недоступности чек добавляется в очередь автоматически except Exception as error: logger.error( diff --git a/app/utils/cache.py b/app/utils/cache.py index 88eaa1cf..5bbc916f 100644 --- a/app/utils/cache.py +++ b/app/utils/cache.py @@ -159,7 +159,7 @@ class CacheService: async def get_hash(self, name: str, key: str = None) -> Optional[Union[dict, str]]: if not self._connected: return None - + try: if key: value = await self.redis_client.hget(name, key) @@ -171,6 +171,56 @@ class CacheService: logger.error(f"Ошибка получения хеша {name}: {e}") return None + async def lpush(self, key: str, value: Any) -> bool: + """Добавить элемент в начало списка (очереди).""" + if not self._connected: + return False + + try: + serialized = json.dumps(value, default=str) + await self.redis_client.lpush(key, serialized) + return True + except Exception as e: + logger.error(f"Ошибка добавления в очередь {key}: {e}") + return False + + async def rpop(self, key: str) -> Optional[Any]: + """Извлечь элемент из конца списка (FIFO очередь).""" + if not self._connected: + return None + + try: + value = await self.redis_client.rpop(key) + if value: + return json.loads(value) + return None + except Exception as e: + logger.error(f"Ошибка извлечения из очереди {key}: {e}") + return None + + async def llen(self, key: str) -> int: + """Получить длину списка (очереди).""" + if not self._connected: + return 0 + + try: + return await self.redis_client.llen(key) + except Exception as e: + logger.error(f"Ошибка получения длины очереди {key}: {e}") + return 0 + + async def lrange(self, key: str, start: int = 0, end: int = -1) -> list: + """Получить элементы списка без удаления.""" + if not self._connected: + return [] + + try: + items = await self.redis_client.lrange(key, start, end) + return [json.loads(item) for item in items] + except Exception as e: + logger.error(f"Ошибка чтения очереди {key}: {e}") + return [] + cache = CacheService() diff --git a/main.py b/main.py index fa714e12..1e0e4cf4 100644 --- a/main.py +++ b/main.py @@ -34,6 +34,7 @@ from app.services.external_admin_service import ensure_external_admin_token from app.services.broadcast_service import broadcast_service from app.services.referral_contest_service import referral_contest_service from app.services.contest_rotation_service import contest_rotation_service +from app.services.nalogo_queue_service import nalogo_queue_service from app.utils.startup_timeline import StartupTimeline from app.utils.timezone import TimezoneAwareFormatter @@ -271,6 +272,11 @@ async def main(): payment_service = PaymentService(bot) auto_payment_verification_service.set_payment_service(payment_service) + # Настройка сервиса очереди чеков NaloGO + if payment_service.nalogo_service: + nalogo_queue_service.set_nalogo_service(payment_service.nalogo_service) + nalogo_queue_service.set_bot(bot) + verification_providers: list[str] = [] auto_verification_active = False async with timeline.stage( @@ -331,6 +337,27 @@ async def main(): if auto_verification_active: stage.log("Фоновая автопроверка запущена") + async with timeline.stage( + "Очередь чеков NaloGO", + "🧾", + success_message="Сервис очереди чеков запущен", + ) as stage: + if settings.is_nalogo_enabled(): + try: + await nalogo_queue_service.start() + if nalogo_queue_service.is_running(): + queue_len = await payment_service.nalogo_service.get_queue_length() + if queue_len > 0: + stage.log(f"В очереди ожидает {queue_len} чек(ов)") + stage.success("Фоновая обработка чеков активна") + else: + stage.skip("Сервис не запущен") + except Exception as e: + stage.warning(f"Ошибка запуска очереди чеков: {e}") + logger.error(f"❌ Ошибка запуска очереди чеков NaloGO: {e}") + else: + stage.skip("NaloGO отключен настройками") + async with timeline.stage( "Внешняя админка", "🛡️", @@ -646,6 +673,12 @@ async def main(): except Exception as e: logger.error(f"Ошибка остановки ротации игр: {e}") + logger.info("ℹ️ Остановка очереди чеков NaloGO...") + try: + await nalogo_queue_service.stop() + except Exception as e: + logger.error(f"Ошибка остановки очереди чеков NaloGO: {e}") + logger.info("ℹ️ Остановка сервиса бекапов...") try: await backup_service.stop_auto_backup()