Compare commits

...

107 Commits

Author SHA1 Message Date
Egor 92ff4f661f Update README.md 2025-08-31 05:54:54 +03:00
Egor dfebc5e608 Update .gitignore 2025-08-31 05:52:00 +03:00
Egor 8e6121cd0b Update README.md 2025-08-31 05:48:08 +03:00
Egor 6c0d5714bd Update remnawave_service.py 2025-08-31 05:19:43 +03:00
Egor ed29697434 Update remnawave_service.py 2025-08-31 05:18:56 +03:00
Egor 07aebef82f Update remnawave_service.py 2025-08-31 05:16:17 +03:00
Egor 7ce390fd50 Update remnawave_service.py 2025-08-31 05:09:31 +03:00
Egor cebe392d1c Update remnawave_service.py 2025-08-31 04:55:19 +03:00
Egor fd160f585a Update remnawave_service.py 2025-08-31 04:50:31 +03:00
Egor 5c0516fbc0 Update remnawave_service.py 2025-08-31 04:45:25 +03:00
Egor ba075326c1 Update admin.py 2025-08-31 04:39:08 +03:00
Egor c32425341e Update remnawave.py 2025-08-31 04:38:20 +03:00
Egor cc8df1eff3 Update remnawave_service.py 2025-08-31 04:33:56 +03:00
Egor 49cda0791c Update .env.example 2025-08-31 04:25:10 +03:00
Egor 60ea363a62 Update README.md 2025-08-31 04:24:37 +03:00
Egor f66619b225 Update subscription.py 2025-08-31 04:19:18 +03:00
Egor 5999402e1a Update inline.py 2025-08-31 04:17:18 +03:00
Egor 3aebccc090 Update inline.py 2025-08-31 04:16:37 +03:00
Egor ed577dd51a Update config.py 2025-08-31 04:15:45 +03:00
Egor e517487596 Update main.py 2025-08-31 04:03:37 +03:00
Egor f9441909a4 Update balance.py 2025-08-31 02:42:33 +03:00
Egor f13a015bfa Update payment_service.py 2025-08-31 00:47:07 +03:00
Egor 18bf946464 Update yookassa_webhook.py 2025-08-31 00:39:04 +03:00
Egor e4c46b2488 Update yookassa_webhook.py 2025-08-31 00:33:37 +03:00
Egor e0539b116c Update yookassa_webhook.py 2025-08-31 00:28:04 +03:00
Egor acdf7af1d9 Update yookassa_webhook.py 2025-08-31 00:22:18 +03:00
Egor d5b7828079 Update yookassa_webhook.py 2025-08-31 00:12:53 +03:00
Egor 03be071edc Update webhook_server.py 2025-08-30 23:54:27 +03:00
Egor e8742f30a6 Update webhook_server.py 2025-08-30 23:51:25 +03:00
Egor 3d7727dbf2 Update webhook_server.py 2025-08-30 23:46:31 +03:00
Egor 5bff67860f Update tribute.py 2025-08-30 23:45:48 +03:00
Egor 429e967da2 Update tribute.py 2025-08-30 23:44:35 +03:00
Egor c3df3beeae Update tribute.py 2025-08-30 23:42:55 +03:00
Egor 4d664c15d7 Update yookassa_webhook.py 2025-08-30 23:37:58 +03:00
Egor 53be365e38 Update tribute.py 2025-08-30 23:36:58 +03:00
Egor 3c3ffed32e Update yookassa_webhook.py 2025-08-30 23:36:12 +03:00
Egor 015726b3a6 Update webhook_server.py 2025-08-30 23:35:20 +03:00
Egor 4a3267712b Update main.py 2025-08-30 23:34:38 +03:00
Egor 9b45aae410 Update main.py 2025-08-30 23:24:43 +03:00
Egor f43b391d0e Update referral_service.py 2025-08-30 23:22:13 +03:00
Egor b09b7a0c84 Add files via upload 2025-08-30 23:21:04 +03:00
Egor 541b44b933 Update payment_service.py 2025-08-30 19:32:55 +03:00
Egor b21fce0eff Update webhook_server.py 2025-08-30 19:23:58 +03:00
Egor 017aa27c0f Update webhook_server.py 2025-08-30 19:17:31 +03:00
Egor 496978e8a2 Update texts.py 2025-08-30 19:11:45 +03:00
Egor b3711a590d Update user.py 2025-08-30 19:08:50 +03:00
Egor 631a0f2c4c Update payment_service.py 2025-08-30 19:05:55 +03:00
Egor 503f7ca455 Update transaction.py 2025-08-30 19:01:34 +03:00
Egor 80c771b7bc Update texts.py 2025-08-30 18:59:08 +03:00
Egor 6e344d350b Update texts.py 2025-08-30 18:58:30 +03:00
Egor 7738bc83af Update user.py 2025-08-30 18:57:32 +03:00
Egor 6f7a1eecf6 Update transaction.py 2025-08-30 18:56:07 +03:00
Egor bf147a736a Update main.py 2025-08-30 18:51:43 +03:00
Egor ea3decf1ac Update payment_service.py 2025-08-30 18:49:30 +03:00
Egor b1c0315a6e Update Dockerfile 2025-08-30 18:46:10 +03:00
Egor f6fa7e8f41 Update main.py 2025-08-30 18:45:44 +03:00
Egor 33d4d7994c Update tribute.py 2025-08-30 18:44:45 +03:00
Egor 0ce4ce2b08 Update webhook_server.py 2025-08-30 18:43:32 +03:00
Egor 3a9163e919 Update maintenance_service.py 2025-08-30 17:58:52 +03:00
Egor 376cff748b Update subscription_service.py 2025-08-30 16:57:35 +03:00
Egor d5cef080da Update referral_service.py 2025-08-30 16:51:34 +03:00
Egor 33e31213b3 Update balance.py 2025-08-30 16:05:42 +03:00
Egor 8077ba0c35 Update maintenance_service.py 2025-08-30 15:50:01 +03:00
Egor 67ad3847f8 Update maintenance_service.py 2025-08-30 15:48:12 +03:00
Egor 1eda042b34 Update main.py 2025-08-30 15:40:32 +03:00
Egor 961ff46f43 Update maintenance_service.py 2025-08-30 15:39:43 +03:00
Egor 535ec34a1b Update maintenance_service.py 2025-08-30 15:15:42 +03:00
Egor fe474fd0bf Merge pull request #18 from Fr1ngg/MAINTENANCE_MODE
Maintenance mode
2025-08-30 14:50:04 +03:00
Egor bfc737c1d3 Update main.py 2025-08-30 14:47:40 +03:00
Egor ce0eb32d4e Update maintenance_service.py 2025-08-30 14:46:41 +03:00
Egor 090df2e3ea Update maintenance.py 2025-08-30 14:45:02 +03:00
Egor 01d9df2867 Update maintenance.py 2025-08-30 14:43:57 +03:00
Egor 334a14c80c Update bot.py 2025-08-30 14:41:55 +03:00
Egor 3d952570af Update README.md 2025-08-30 14:40:15 +03:00
Egor bf9567fe6f Update README.md 2025-08-30 14:37:22 +03:00
Egor 6c4c39b454 Add files via upload 2025-08-30 14:35:38 +03:00
Egor d39eb3f79f Update user.py 2025-08-30 13:00:32 +03:00
Egor 7f0ce61126 Update user.py 2025-08-30 05:00:44 +03:00
Egor 775056da6f Update user.py 2025-08-30 04:57:06 +03:00
Egor a15ecb375b Update user.py 2025-08-30 04:53:34 +03:00
Egor 3dea941c51 Update user_service.py 2025-08-30 04:49:50 +03:00
Egor 410f0aba95 Update user.py 2025-08-30 04:43:51 +03:00
Egor 1ec13fbf63 Update user.py 2025-08-30 04:39:57 +03:00
Egor 012408b7cd Update user.py 2025-08-30 04:36:23 +03:00
Egor fdca75b153 Update user.py 2025-08-30 04:35:19 +03:00
Egor 40c235cdac Update user_service.py 2025-08-30 04:26:52 +03:00
Egor 7ea8fa735c Update user.py 2025-08-30 04:26:25 +03:00
Egor 758141dbc4 Update user_service.py 2025-08-30 04:22:57 +03:00
Egor 3c37fa45e6 Update user_service.py 2025-08-30 04:17:01 +03:00
Egor 35a80f4585 Update user_service.py 2025-08-30 04:13:58 +03:00
Egor 57e033c7c6 Update config.py 2025-08-30 04:09:00 +03:00
Egor d05d070faa Update user_service.py 2025-08-30 04:07:09 +03:00
Egor 3f3829a745 Update tribute_service.py 2025-08-30 03:04:43 +03:00
Egor 21bb89e5ea Update transaction.py 2025-08-30 03:03:44 +03:00
Egor fad5ecbb89 Update tribute_service.py 2025-08-30 03:00:03 +03:00
Egor 2c64c57edd Update tribute_service.py 2025-08-30 02:53:15 +03:00
Egor 117c637756 Update user.py 2025-08-30 02:46:18 +03:00
Egor 488f161322 Update tribute_service.py 2025-08-30 02:45:00 +03:00
Egor a55ef684e6 Update user.py 2025-08-30 02:33:49 +03:00
Egor 90234ba01e Update tribute_service.py 2025-08-30 02:28:25 +03:00
Egor cd83361050 Update tribute_service.py 2025-08-30 02:22:02 +03:00
Egor 08cda481ba Update tribute.py 2025-08-30 02:17:47 +03:00
Egor 45f998d134 Update tribute.py 2025-08-29 19:46:07 +03:00
Egor 33ee5b4ec4 Update payment_service.py 2025-08-29 18:55:02 +03:00
Egor 545bbaafa5 Update balance.py 2025-08-29 18:38:59 +03:00
Egor cc957460bb Update balance.py 2025-08-29 18:21:32 +03:00
Egor cd51dab845 Update README.md 2025-08-29 12:29:43 +03:00
27 changed files with 1920 additions and 714 deletions
+6
View File
@@ -29,6 +29,12 @@ TRIAL_DEVICE_LIMIT=2
TRIAL_SQUAD_UUID=
DEFAULT_TRAFFIC_RESET_STRATEGY=MONTH
# === УПРАВЛЕНИЕ ПЕРИОДАМИ ПОДПИСКИ ===
# Доступные периоды подписки (через запятую)
# Возможные значения: 14,30,60,90,180,360
AVAILABLE_SUBSCRIPTION_PERIODS=14,30,60,90,180,360
AVAILABLE_RENEWAL_PERIODS=30,90,180
# SUBSCRIPTION PRICING (в копейках для точности)
BASE_SUBSCRIPTION_PRICE=50000
+30 -2
View File
@@ -1,2 +1,30 @@
/screens
/assets
# Игнорируем все файлы и папки по умолчанию
*
# Исключения: разрешаем только нужные файлы
!.env.example
!Dockerfile
!app-config.json
!main.py
!requirements.txt
# Разрешаем папку app/ и все её содержимое рекурсивно
!app/
!app/**
# Дополнительно разрешаем README и лицензию (опционально)
!README.md
!LICENSE
# Разрешаем .gitignore чтобы он попал в репозиторий
!.gitignore
# Внутри разрешенных папок игнорируем служебные файлы
app/__pycache__/
app/**/__pycache__/
app/**/*.pyc
app/**/*.pyo
app/**/*.pyd
*.pyc
*.pyo
*.pyd
+4 -2
View File
@@ -1,3 +1,4 @@
FROM python:3.11-slim
WORKDIR /app
@@ -26,8 +27,9 @@ RUN useradd --create-home --shell /bin/bash app
RUN chown -R app:app /app
USER app
# Expose port (if needed for webhooks)
EXPOSE 8000
# Expose webhook ports для платежных систем
EXPOSE 8081
EXPOSE 8082
# Run the application
CMD ["python", "main.py"]
+119 -207
View File
@@ -33,7 +33,7 @@
### ⚡ **Полная автоматизация VPN бизнеса**
- 🎯 **Готовое решение** - разверни за 5 минут, начни продавать сегодня
- 💰 **Многоканальные платежи** - Telegram Stars + Tribute + планы на ЮKassa
- 💰 **Многоканальные платежи** - Telegram Stars + Tribute + ЮKassa
- 🔄 **Автоматизация 99%** - от регистрации до продления подписок
- 📊 **Детальная аналитика**
@@ -44,12 +44,14 @@
- 🎁 **Промо-система** - коды на деньги, дни подписки, триал-периоды
- 3 режима показа ссылки подписки: 1) С гайдом по подключению прямо в боте(тянущий данные приложений и ссылок на скачку из app-config.json) 2) Обычное открытие ссылки подписки в миниапе 3) Интеграция сабпейджа maposia - кастомно прописать ссылку можно
- Возможность переключаться между пакетной продажей трафика и фиксированной(Пропуская шаг выбора пакета трафика при оформлении/настройки подписки юзера)
- Возможность задать доступные дни для покупки первой подписки и при продлении
### 💪 **Enterprise готовность**
- 🏗️ **Современная архитектура** - AsyncIO, PostgreSQL, Redis
- 🔒 **Безопасность** - шифрование, валидация, rate limiting
- 📈 **Масштабируемость**
- 🔧 **Мониторинг** - Prometheus, Grafana, health checks
- 🔧 **Режим технических работ** - Ручное включение + Мониторинг системы, который в случае падении панели Remnawave переведет бота в режим технических работ и обратно - отключит его, если панель поднимется.
---
@@ -61,6 +63,7 @@
# 1. Скачай репозиторий
git clone https://github.com/Fr1ngg/remnawave-bedolaga-telegram-bot.git
cd remnawave-bedolaga-telegram-bot
или собери docker-compose.yml
# 2. Настрой конфиг
cp .env.example .env
@@ -89,50 +92,98 @@ docker compose logs -f bot
<summary>🔧 Полная конфигурация .env</summary>
```env
# 🏷️ Основные настройки
NODE_ENV=production
DEBUG=false
LOG_LEVEL=INFO
# TELEGRAM BOT CONFIGURATION
BOT_TOKEN=
ADMIN_IDS=
SUPPORT_USERNAME=
# 🗄️ База данных
POSTGRES_DB=bedolaga_bot
POSTGRES_USER=bedolaga_user
POSTGRES_PASSWORD=secure_password_123
POSTGRES_PORT=5432
DATABASE_URL=postgresql+asyncpg://bedolaga_user:secure_password_123@postgres:5432/bedolaga_bot
# DATABASE CONFIGURATION
DATABASE_URL=sqlite+aiosqlite:///./bot.db
REDIS_URL=redis://localhost:6379/0
# ⚡ Redis кеш
REDIS_PASSWORD=redis_password_123
REDIS_PORT=6379
REDIS_URL=redis://:redis_password_123@redis:6379/0
# REMNAWAVE API CONFIGURATION
REMNAWAVE_API_URL=
REMNAWAVE_API_KEY=
# 🤖 Telegram Bot
BOT_TOKEN=your_bot_token_here
ADMIN_IDS=123456789,987654321
SUPPORT_USERNAME=@your_support
# === NEW: Traffic Selection Mode Settings ===
# Режим выбора трафика:
# "selectable" - пользователи выбирают пакеты трафика (по умолчанию)
# "fixed" - фиксированный лимит трафика для всех подписок, доступно 5/10/25/50/100/250/0 (0 безлимит) гб
TRAFFIC_SELECTION_MODE=selectable
# 🔗 Remnawave API
REMNAWAVE_API_URL=https://your-panel.com
REMNAWAVE_API_KEY=your_jwt_token_here
# Фиксированный лимит трафика в ГБ (используется только в режиме "fixed")
# 0 = безлимит
# для "fixed" обязательно должы быть проставлены цены на пакеты 5/10/25/50/100/250/0 можно постать 0 руб - будет беслпатно
FIXED_TRAFFIC_LIMIT_GB=0
# 🌐 Webhook настройки
WEBHOOK_DOMAIN=your-domain.com
WEBHOOK_PORT=8081
WEBHOOK_URL=https://your-domain.com
WEBHOOK_PATH=/webhook
# TRIAL SUBSCRIPTION SETTINGS
TRIAL_DURATION_DAYS=3
TRIAL_TRAFFIC_LIMIT_GB=10
TRIAL_DEVICE_LIMIT=2
TRIAL_SQUAD_UUID=
DEFAULT_TRAFFIC_RESET_STRATEGY=MONTH
# ⭐ Telegram Stars
# === УПРАВЛЕНИЕ ПЕРИОДАМИ ПОДПИСКИ ===
# Доступные периоды подписки (через запятую)
# Возможные значения: 14,30,60,90,180,360
AVAILABLE_SUBSCRIPTION_PERIODS=14,30,60,90,180,360
AVAILABLE_RENEWAL_PERIODS=30,90,180
# SUBSCRIPTION PRICING (в копейках для точности)
BASE_SUBSCRIPTION_PRICE=50000
PRICE_14_DAYS=5000
PRICE_30_DAYS=9900
PRICE_60_DAYS=18900
PRICE_90_DAYS=26900
PRICE_180_DAYS=49900
PRICE_360_DAYS=89900
PRICE_TRAFFIC_5GB=2000
PRICE_TRAFFIC_10GB=4000
PRICE_TRAFFIC_25GB=6000
PRICE_TRAFFIC_50GB=10000
PRICE_TRAFFIC_100GB=15000
PRICE_TRAFFIC_250GB=20000
PRICE_TRAFFIC_UNLIMITED=25000
PRICE_PER_DEVICE=5000
# REFERRAL SYSTEM SETTINGS
REFERRAL_REGISTRATION_REWARD=5000
REFERRED_USER_REWARD=10000
REFERRAL_COMMISSION_PERCENT=25
# Режим работы кнопки "Подключиться"
# guide - открывает гайд подключения (режим 1)
# miniapp_subscription - открывает ссылку подписки в мини-приложении (режим 2)
# miniapp_custom - открывает заданную ссылку в мини-приложении (режим 3)
CONNECT_BUTTON_MODE=miniapp_subscription
# URL для режима miniapp_custom (обязателен при CONNECT_BUTTON_MODE=miniapp_custom)
# MINIAPP_CUSTOM_URL=
# AUTO-PAYMENT SETTINGS
AUTOPAY_WARNING_DAYS=3,1
# MONITORING SETTINGS
MONITORING_INTERVAL=60
INACTIVE_USER_DELETE_MONTHS=3
TRIAL_WARNING_HOURS=2
ENABLE_NOTIFICATIONS=true
NOTIFICATION_RETRY_ATTEMPTS=3
MONITORING_LOGS_RETENTION_DAYS=30
# PAYMENT SYSTEMS
TELEGRAM_STARS_ENABLED=true
# 💳 Tribute платежи
TRIBUTE_ENABLED=true
TRIBUTE_API_KEY=your_tribute_api_key
TRIBUTE_ENABLED=false
TRIBUTE_API_KEY=
TRIBUTE_WEBHOOK_SECRET=your_webhook_secret
TRIBUTE_DONATE_LINK=https://t.me/tribute/app?startapp=XXXX
TRIBUTE_WEBHOOK_PATH=/tribute-webhook
TRIBUTE_WEBHOOK_PORT=8081
TRIBUTE_WEBHOOK_SECRET=your_webhook_secret
# 💳 YOOKASSA
# === НОВЫЕ НАСТРОЙКИ YOOKASSA ===
# Включение/выключение YooKassa
YOOKASSA_ENABLED=false
@@ -187,65 +238,23 @@ YOOKASSA_WEBHOOK_PATH=/yookassa-webhook
YOOKASSA_WEBHOOK_PORT=8082
YOOKASSA_WEBHOOK_SECRET=ваш_секретный_ключ_для_webhook
# 🚀 Режим работы кнопки "Подключиться"
# guide - открывает гайд подключения c настройками и парамтерами из app-config.json (режим 1)
# miniapp_subscription - открывает ссылку подписки в мини-приложении (режим 2)
# miniapp_custom - открывает заданную ссылку в мини-приложении (режим 3)
CONNECT_BUTTON_MODE=miniapp_subscription
# URL для режима miniapp_custom (обязателен при CONNECT_BUTTON_MODE=miniapp_custom)
# MINIAPP_CUSTOM_URL=
WEBHOOK_URL=https://example.com
WEBHOOK_PATH=/webhook
# 🎛️ === NEW: Traffic Selection Mode Settings ===
# Режим выбора трафика:
# "selectable" - пользователи выбирают пакеты трафика (по умолчанию)
# "fixed" - фиксированный лимит трафика для всех подписок(БЕЗ ШАГА ВЫБОРА ПАКЕТА ТРАФИКА ВО ВРЕМЯ ОФОРМЛЕНИЯ ПОДПИСКИ), доступно 5/10/25/50/100/250/0 (0 безлимит) гб
# Фиксированный лимит трафика в ГБ (используется только в режиме "fixed")
# 0 = безлимит
# для "fixed" обязательно должы быть проставлены цены на пакеты 5/10/25/50/100/250/0 можно постать 0 руб - будет беслпатно
TRAFFIC_SELECTION_MODE=selectable
FIXED_TRAFFIC_LIMIT_GB=0
# LOCALIZATION
DEFAULT_LANGUAGE=ru
AVAILABLE_LANGUAGES=ru
# 🎁 Триал настройки
TRIAL_ENABLED=true
TRIAL_DURATION_DAYS=3
TRIAL_TRAFFIC_LIMIT_GB=10
TRIAL_DEVICE_LIMIT=2
TRIAL_SQUAD_UUID=your_trial_squad_uuid
# LOGGING
LOG_LEVEL=INFO
LOG_FILE=/tmp/bot.log
# 💰 Ценообразование (в копейках)
BASE_SUBSCRIPTION_PRICE=50000
PRICE_14_DAYS=5000
PRICE_30_DAYS=9900
PRICE_60_DAYS=18900
PRICE_90_DAYS=26900
PRICE_180_DAYS=49900
PRICE_360_DAYS=89900
# DEVELOPMENT
DEBUG=false
PRICE_TRAFFIC_5GB=2000
PRICE_TRAFFIC_10GB=4000
PRICE_TRAFFIC_25GB=6000
PRICE_TRAFFIC_50GB=10000
PRICE_TRAFFIC_100GB=15000
PRICE_TRAFFIC_250GB=20000
PRICE_TRAFFIC_UNLIMITED=25000
PRICE_PER_DEVICE=5000
# 🤝 Реферальная система
REFERRAL_REGISTRATION_REWARD=5000
REFERRED_USER_REWARD=2500
REFERRAL_COMMISSION_PERCENT=10
# 🔍 Мониторинг
MONITORING_INTERVAL=60
ENABLE_NOTIFICATIONS=true
AUTOPAY_WARNING_DAYS=3,1
MONITORING_LOGS_RETENTION_DAYS=30
INACTIVE_USER_DELETE_MONTHS=3
# 📊 Опционально для мониторинга
GRAFANA_USER=admin
GRAFANA_PASSWORD=admin123
MAINTENANCE_CHECK_INTERVAL=30
MAINTENANCE_AUTO_ENABLE=true
MAINTENANCE_MESSAGE="Ведутся технические работы"
```
</details>
@@ -421,7 +430,6 @@ bedolaga_bot/
```
project/
├── docker-compose.yml # 🚀 Продакшн
├── docker-compose.local.yml # 🏠 Разработка
├── .env # ⚙️ Конфиг
└── .env.example # 📝 Пример
```
@@ -432,66 +440,45 @@ project/
<summary>📄 Показать полный docker-compose.yml</summary>
```yaml
version: '3.8'
services:
# 🗄️ PostgreSQL Database
postgres:
image: postgres:15-alpine
container_name: bedolaga_postgres
container_name: remnawave_bot_db
restart: unless-stopped
environment:
POSTGRES_DB: ${POSTGRES_DB:-bedolaga_bot}
POSTGRES_USER: ${POSTGRES_USER:-bedolaga_user}
POSTGRES_DB: ${POSTGRES_DB:-remnawave_bot}
POSTGRES_USER: ${POSTGRES_USER:-remnawave_user}
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:-secure_password_123}
POSTGRES_INITDB_ARGS: "--encoding=UTF8 --lc-collate=C --lc-ctype=C"
POSTGRES_INITDB_ARGS: "--encoding=UTF8"
volumes:
- postgres_data:/var/lib/postgresql/data
- ./init-scripts:/docker-entrypoint-initdb.d:ro
ports:
- "${POSTGRES_PORT:-5432}:5432"
networks:
- bedolaga_network
- bot_network
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-bedolaga_user} -d ${POSTGRES_DB:-bedolaga_bot}"]
interval: 10s
test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-remnawave_user} -d ${POSTGRES_DB:-remnawave_bot}"]
interval: 30s
timeout: 5s
retries: 5
start_period: 30s
logging:
driver: "json-file"
options:
max-size: "10m"
max-file: "3"
# ⚡ Redis Cache
redis:
image: redis:7-alpine
container_name: bedolaga_redis
container_name: remnawave_bot_redis
restart: unless-stopped
command: redis-server --appendonly yes --requirepass ${REDIS_PASSWORD:-redis_password_123}
command: redis-server --appendonly yes
volumes:
- redis_data:/data
ports:
- "${REDIS_PORT:-6379}:6379"
networks:
- bedolaga_network
- bot_network
healthcheck:
test: ["CMD", "redis-cli", "--no-auth-warning", "-a", "${REDIS_PASSWORD:-redis_password_123}", "ping"]
interval: 10s
timeout: 5s
test: ["CMD", "redis-cli", "ping"]
interval: 30s
timeout: 10s
retries: 3
start_period: 10s
logging:
driver: "json-file"
options:
max-size: "5m"
max-file: "3"
# 🤖 Telegram Bot
bot:
image: fr1ngg/remnawave-bedolaga-telegram-bot:latest
container_name: bedolaga_bot
container_name: remnawave_bot
restart: unless-stopped
depends_on:
postgres:
@@ -501,108 +488,35 @@ services:
env_file:
- .env
environment:
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER:-bedolaga_user}:${POSTGRES_PASSWORD:-secure_password_123}@postgres:5432/${POSTGRES_DB:-bedolaga_bot}
REDIS_URL: redis://:${REDIS_PASSWORD:-redis_password_123}@redis:6379/0
LOG_LEVEL: ${LOG_LEVEL:-INFO}
DEBUG: ${DEBUG:-false}
HEALTH_CHECK_ENABLED: "true"
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER:-remnawave_user}:${POSTGRES_PASSWORD:-secure_password_123}@postgres:5432/${POSTGRES_DB:-remnawave_bot}
REDIS_URL: redis://redis:6379/0
volumes:
- ./logs:/app/logs
- ./data:/app/data
- ./backups:/app/backups
- /etc/timezone:/etc/timezone:ro
- /etc/localtime:/etc/localtime:ro
ports:
- "${WEBHOOK_PORT:-8081}:8081"
- "${TRIBUTE_WEBHOOK_PORT:-8081}:8081"
- "${YOOKASSA_WEBHOOK_PORT:-8082}:8082"
networks:
- bedolaga_network
- bot_network
healthcheck:
test: ["CMD", "python", "-c", "import requests; requests.get('http://localhost:8081/health', timeout=5)"]
test: ["CMD", "wget", "--no-verbose", "--tries=1", "--spider", "http://localhost:8081/health"]
interval: 30s
timeout: 10s
retries: 3
start_period: 60s
logging:
driver: "json-file"
options:
max-size: "50m"
max-file: "5"
labels:
- "traefik.enable=true"
- "traefik.http.routers.bedolaga-webhook.rule=Host(`${WEBHOOK_DOMAIN:-localhost}`) && PathPrefix(`/tribute-webhook`)"
- "traefik.http.services.bedolaga-webhook.loadbalancer.server.port=8081"
# 📊 Monitoring (опционально)
prometheus:
image: prom/prometheus:latest
container_name: bedolaga_prometheus
restart: unless-stopped
command:
- '--config.file=/etc/prometheus/prometheus.yml'
- '--storage.tsdb.path=/prometheus'
- '--web.console.libraries=/etc/prometheus/console_libraries'
- '--web.console.templates=/etc/prometheus/consoles'
- '--storage.tsdb.retention.time=200h'
- '--web.enable-lifecycle'
volumes:
- ./monitoring/prometheus.yml:/etc/prometheus/prometheus.yml
- prometheus_data:/prometheus
ports:
- "9090:9090"
networks:
- bedolaga_network
profiles:
- monitoring
# 📈 Grafana (опционально)
grafana:
image: grafana/grafana:latest
container_name: bedolaga_grafana
restart: unless-stopped
environment:
GF_SECURITY_ADMIN_USER: ${GRAFANA_USER:-admin}
GF_SECURITY_ADMIN_PASSWORD: ${GRAFANA_PASSWORD:-admin123}
volumes:
- grafana_data:/var/lib/grafana
- ./monitoring/grafana/dashboards:/etc/grafana/provisioning/dashboards
- ./monitoring/grafana/datasources:/etc/grafana/provisioning/datasources
ports:
- "3000:3000"
networks:
- bedolaga_network
profiles:
- monitoring
# 📦 Volumes
volumes:
postgres_data:
driver: local
driver_opts:
type: none
o: bind
device: ./volumes/postgres
redis_data:
driver: local
driver_opts:
type: none
o: bind
device: ./volumes/redis
prometheus_data:
driver: local
grafana_data:
driver: local
# 🌐 Networks
networks:
bedolaga_network:
bot_network:
driver: bridge
ipam:
driver: default
config:
- subnet: 172.20.0.0/16
gateway: 172.20.0.1
driver_opts:
com.docker.network.bridge.name: br-bedolaga
```
</details>
@@ -613,8 +527,6 @@ networks:
<summary>📄 Показать dev конфигурацию</summary>
```yaml
version: '3.8'
services:
# 🗄️ PostgreSQL для разработки
postgres-dev:
+34 -11
View File
@@ -9,6 +9,8 @@ from app.middlewares.auth import AuthMiddleware
from app.middlewares.logging import LoggingMiddleware
from app.middlewares.throttling import ThrottlingMiddleware
from app.middlewares.subscription_checker import SubscriptionStatusMiddleware
from app.middlewares.maintenance import MaintenanceMiddleware
from app.services.maintenance_service import maintenance_service
from app.utils.cache import cache
from app.handlers import (
@@ -20,7 +22,8 @@ from app.handlers.admin import (
promocodes as admin_promocodes, messages as admin_messages,
monitoring as admin_monitoring, referrals as admin_referrals,
rules as admin_rules, remnawave as admin_remnawave,
statistics as admin_statistics, servers as admin_servers
statistics as admin_statistics, servers as admin_servers,
maintenance as admin_maintenance
)
logger = logging.getLogger(__name__)
@@ -37,9 +40,9 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
try:
await cache.connect()
logger.info("Кеш инициализирован")
logger.info("Кеш инициализирован")
except Exception as e:
logger.warning(f"⚠️ Кеш не инициализирован: {e}")
logger.warning(f"Кеш не инициализирован: {e}")
from aiogram.client.default import DefaultBotProperties
from aiogram.enums import ParseMode
@@ -53,24 +56,23 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
redis_client = redis.from_url(settings.REDIS_URL)
await redis_client.ping()
storage = RedisStorage(redis_client)
logger.info("Подключено к Redis для FSM storage")
logger.info("Подключено к Redis для FSM storage")
except Exception as e:
logger.warning(f"⚠️ Не удалось подключиться к Redis: {e}")
logger.info("🔄 Используется MemoryStorage для FSM")
logger.warning(f"Не удалось подключиться к Redis: {e}")
logger.info("Используется MemoryStorage для FSM")
storage = MemoryStorage()
dp = Dispatcher(storage=storage)
dp.message.middleware(LoggingMiddleware())
dp.callback_query.middleware(LoggingMiddleware())
dp.message.middleware(AuthMiddleware())
dp.callback_query.middleware(AuthMiddleware())
dp.message.middleware(MaintenanceMiddleware())
dp.callback_query.middleware(MaintenanceMiddleware())
dp.message.middleware(ThrottlingMiddleware())
dp.callback_query.middleware(ThrottlingMiddleware())
dp.message.middleware(SubscriptionStatusMiddleware())
dp.callback_query.middleware(SubscriptionStatusMiddleware())
start.register_handlers(dp)
menu.register_handlers(dp)
subscription.register_handlers(dp)
@@ -78,7 +80,6 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
promocode.register_handlers(dp)
referral.register_handlers(dp)
support.register_handlers(dp)
admin_main.register_handlers(dp)
admin_users.register_handlers(dp)
admin_subscriptions.register_handlers(dp)
@@ -90,9 +91,31 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
admin_rules.register_handlers(dp)
admin_remnawave.register_handlers(dp)
admin_statistics.register_handlers(dp)
admin_maintenance.register_handlers(dp)
common.register_handlers(dp)
logger.info("✅ Бот успешно настроен")
try:
await maintenance_service.start_monitoring()
logger.info("Мониторинг техработ запущен")
except Exception as e:
logger.error(f"Ошибка запуска мониторинга техработ: {e}")
logger.info("Бот успешно настроен")
return bot, dp
async def shutdown_bot():
"""Корректное завершение работы бота"""
try:
await maintenance_service.stop_monitoring()
logger.info("Мониторинг техработ остановлен")
except Exception as e:
logger.error(f"Ошибка остановки мониторинга: {e}")
try:
await cache.close()
logger.info("Соединения с кешем закрыты")
except Exception as e:
logger.error(f"Ошибка закрытия кеша: {e}")
+60 -1
View File
@@ -20,6 +20,8 @@ class Settings(BaseSettings):
TRIAL_DURATION_DAYS: int = 3
TRIAL_TRAFFIC_LIMIT_GB: int = 10
TRIAL_DEVICE_LIMIT: int = 2
DEFAULT_TRAFFIC_LIMIT_GB: int = 100
DEFAULT_DEVICE_LIMIT: int = 1
TRIAL_SQUAD_UUID: str
DEFAULT_TRAFFIC_RESET_STRATEGY: str = "MONTH"
@@ -31,7 +33,9 @@ class Settings(BaseSettings):
NOTIFICATION_CACHE_HOURS: int = 24
BASE_SUBSCRIPTION_PRICE: int = 50000
AVAILABLE_SUBSCRIPTION_PERIODS: str = "14,30,60,90,180,360"
AVAILABLE_RENEWAL_PERIODS: str = "30,90,180"
PRICE_14_DAYS: int = 50000
PRICE_30_DAYS: int = 99000
PRICE_60_DAYS: int = 189000
@@ -63,6 +67,11 @@ class Settings(BaseSettings):
MONITORING_INTERVAL: int = 60
INACTIVE_USER_DELETE_MONTHS: int = 3
MAINTENANCE_MODE: bool = False
MAINTENANCE_CHECK_INTERVAL: int = 30
MAINTENANCE_AUTO_ENABLE: bool = True
MAINTENANCE_MESSAGE: str = "🔧 Ведутся технические работы. Сервис временно недоступен. Попробуйте позже."
TELEGRAM_STARS_ENABLED: bool = True
@@ -197,6 +206,56 @@ class Settings(BaseSettings):
elif self.WEBHOOK_URL:
return f"{self.WEBHOOK_URL}/payment-success"
return "https://t.me/"
def is_maintenance_mode(self) -> bool:
return self.MAINTENANCE_MODE
def get_maintenance_message(self) -> str:
return self.MAINTENANCE_MESSAGE
def get_maintenance_check_interval(self) -> int:
return self.MAINTENANCE_CHECK_INTERVAL
def is_maintenance_auto_enable(self) -> bool:
return self.MAINTENANCE_AUTO_ENABLE
def get_available_subscription_periods(self) -> List[int]:
try:
periods_str = self.AVAILABLE_SUBSCRIPTION_PERIODS
if not periods_str.strip():
return [30, 90, 180]
periods = []
for period_str in periods_str.split(','):
period_str = period_str.strip()
if period_str:
period = int(period_str)
if period in PERIOD_PRICES:
periods.append(period)
return periods if periods else [30, 90, 180]
except (ValueError, AttributeError):
return [30, 90, 180]
def get_available_renewal_periods(self) -> List[int]:
try:
periods_str = self.AVAILABLE_RENEWAL_PERIODS
if not periods_str.strip():
return [30, 90, 180]
periods = []
for period_str in periods_str.split(','):
period_str = period_str.strip()
if period_str:
period = int(period_str)
if period in PERIOD_PRICES:
periods.append(period)
return periods if periods else [30, 90, 180]
except (ValueError, AttributeError):
return [30, 90, 180]
model_config = {
"env_file": ".env",
+81 -1
View File
@@ -275,4 +275,84 @@ async def get_revenue_by_period(
.order_by(func.date(Transaction.created_at))
)
return [{"date": row.date, "amount_kopeks": row.amount} for row in result]
return [{"date": row.date, "amount_kopeks": row.amount} for row in result]
async def find_tribute_transactions_by_payment_id(
db: AsyncSession,
payment_id: str,
user_telegram_id: Optional[int] = None
) -> List[Transaction]:
"""Найти все Tribute транзакции по payment_id"""
query = select(Transaction).options(selectinload(Transaction.user))
conditions = [
Transaction.external_id == f"donation_{payment_id}",
Transaction.external_id == payment_id,
Transaction.external_id.like(f"%{payment_id}%")
]
query = query.where(
and_(
Transaction.payment_method == PaymentMethod.TRIBUTE.value,
or_(*conditions)
)
)
if user_telegram_id:
from app.database.models import User
query = query.join(User).where(User.telegram_id == user_telegram_id)
result = await db.execute(query.order_by(Transaction.created_at.desc()))
return result.scalars().all()
async def check_tribute_payment_duplicate(
db: AsyncSession,
payment_id: str,
amount_kopeks: int,
user_telegram_id: int
) -> Optional[Transaction]:
transactions = await find_tribute_transactions_by_payment_id(
db, payment_id, user_telegram_id
)
for transaction in transactions:
if (transaction.amount_kopeks == amount_kopeks and
transaction.is_completed and
transaction.user.telegram_id == user_telegram_id):
return transaction
return None
async def create_unique_tribute_transaction(
db: AsyncSession,
user_id: int,
payment_id: str,
amount_kopeks: int,
description: str
) -> Transaction:
"""Создать уникальную Tribute транзакцию с защитой от дубликатов"""
external_id = f"donation_{payment_id}"
existing = await get_transaction_by_external_id(db, external_id, PaymentMethod.TRIBUTE)
if existing:
timestamp = int(datetime.utcnow().timestamp())
external_id = f"donation_{payment_id}_{amount_kopeks}_{timestamp}"
logger.info(f"Создан уникальный external_id для избежания дубликатов: {external_id}")
return await create_transaction(
db=db,
user_id=user_id,
type=TransactionType.DEPOSIT,
amount_kopeks=amount_kopeks,
description=description,
payment_method=PaymentMethod.TRIBUTE,
external_id=external_id,
is_completed=True
)
+49 -21
View File
@@ -126,28 +126,56 @@ async def add_user_balance(
user: User,
amount_kopeks: int,
description: str = "Пополнение баланса"
) -> User:
user.add_balance(amount_kopeks)
user.updated_at = datetime.utcnow()
from app.database.crud.transaction import create_transaction
from app.database.models import TransactionType
await create_transaction(
db=db,
user_id=user.id,
type=TransactionType.DEPOSIT,
amount_kopeks=amount_kopeks,
description=description
)
await db.commit()
await db.refresh(user)
logger.info(f"💰 Пополнен баланс пользователя {user.telegram_id}: +{amount_kopeks/100}")
return user
) -> bool:
try:
old_balance = user.balance_kopeks
user.balance_kopeks += amount_kopeks
user.updated_at = datetime.utcnow()
from app.database.crud.transaction import create_transaction
from app.database.models import TransactionType
await create_transaction(
db=db,
user_id=user.id,
type=TransactionType.DEPOSIT,
amount_kopeks=amount_kopeks,
description=description
)
await db.commit()
await db.refresh(user)
logger.info(f"💰 Баланс пользователя {user.telegram_id} изменен: {old_balance}{user.balance_kopeks} (изменение: +{amount_kopeks})")
return True
except Exception as e:
logger.error(f"Ошибка изменения баланса пользователя {user.id}: {e}")
await db.rollback()
return False
async def add_user_balance_by_id(
db: AsyncSession,
telegram_id: int,
amount_kopeks: int,
description: str = "Пополнение баланса"
) -> bool:
try:
user = await get_user_by_telegram_id(db, telegram_id)
if not user:
logger.error(f"Пользователь с telegram_id {telegram_id} не найден")
return False
return await add_user_balance(db, user, amount_kopeks, description)
except Exception as e:
logger.error(f"Ошибка пополнения баланса пользователя {telegram_id}: {e}")
return False
except Exception as e:
logger.error(f"❌ Ошибка пополнения баланса пользователя {user_id}: {e}")
await db.rollback()
return False
async def subtract_user_balance(
db: AsyncSession,
+77 -23
View File
@@ -28,7 +28,6 @@ class TributeService:
return None
try:
payment_url = f"{self.donate_link}&user_id={user_id}"
logger.info(f"Создана ссылка Tribute для пользователя {user_id}")
@@ -41,7 +40,7 @@ class TributeService:
def verify_webhook_signature(self, payload: str, signature: str) -> bool:
if not self.webhook_secret:
logger.warning("Webhook secret не настроен")
logger.warning("Webhook secret не настроен, пропускаем проверку")
return True
try:
@@ -51,31 +50,54 @@ class TributeService:
hashlib.sha256
).hexdigest()
return hmac.compare_digest(signature, expected_signature)
is_valid = hmac.compare_digest(signature, expected_signature)
if is_valid:
logger.info("✅ Подпись Tribute webhook проверена успешно")
else:
logger.error("❌ Неверная подпись Tribute webhook")
return is_valid
except Exception as e:
logger.error(f"Ошибка проверки подписи webhook: {e}")
return False
async def process_webhook(self, payload: Dict[str, Any]) -> Optional[Dict[str, Any]]:
async def process_webhook(self, payload_or_data) -> Optional[Dict[str, Any]]:
try:
logger.info(f"🔄 Начинаем обработку Tribute webhook")
payment_id = payload.get("id") or payload.get("payment_id")
status = payload.get("status")
amount_rubles = payload.get("amount", 0)
telegram_user_id = payload.get("telegram_user_id") or payload.get("user_id")
if isinstance(payload_or_data, str):
try:
webhook_data = json.loads(payload_or_data)
logger.info(f"📊 Распарсенные данные: {webhook_data}")
except json.JSONDecodeError as e:
logger.error(f"❌ Ошибка парсинга JSON: {e}")
return None
else:
webhook_data = payload_or_data
if not payment_id and "payload" in payload:
data = payload["payload"]
payment_id = None
status = None
amount_kopeks = 0
telegram_user_id = None
payment_id = webhook_data.get("id") or webhook_data.get("payment_id")
status = webhook_data.get("status")
amount_kopeks = webhook_data.get("amount", 0)
telegram_user_id = webhook_data.get("telegram_user_id") or webhook_data.get("user_id")
if not payment_id and "payload" in webhook_data:
data = webhook_data["payload"]
payment_id = data.get("id") or data.get("payment_id")
status = data.get("status")
amount_rubles = data.get("amount", 0)
amount_kopeks = data.get("amount", 0)
telegram_user_id = data.get("telegram_user_id") or data.get("user_id")
if not payment_id and "name" in payload:
event_name = payload.get("name")
data = payload.get("payload", {})
if not payment_id and "name" in webhook_data:
event_name = webhook_data.get("name")
data = webhook_data.get("payload", {})
payment_id = str(data.get("donation_request_id"))
amount_kopeks = data.get("amount", 0)
telegram_user_id = data.get("telegram_user_id")
@@ -87,22 +109,54 @@ class TributeService:
else:
status = "unknown"
logger.info(f"Обработка Tribute webhook: payment_id={payment_id}, status={status}, amount_kopeks={amount_kopeks}, user_id={telegram_user_id}")
logger.info(f"📝 Извлеченные данные: payment_id={payment_id}, status={status}, amount_kopeks={amount_kopeks}, user_id={telegram_user_id}")
if not telegram_user_id:
logger.error("Не найден telegram_user_id в webhook данных")
logger.error("Не найден telegram_user_id в webhook данных")
logger.error(f"🔍 Полные данные для отладки: {json.dumps(webhook_data, ensure_ascii=False, indent=2)}")
return None
return {
try:
telegram_user_id = int(telegram_user_id)
except (ValueError, TypeError):
logger.error(f"❌ Некорректный telegram_user_id: {telegram_user_id}")
return None
result = {
"event_type": "payment",
"payment_id": payment_id or f"tribute_{telegram_user_id}_{amount_kopeks}",
"user_id": int(telegram_user_id),
"amount_kopeks": amount_kopeks,
"user_id": telegram_user_id,
"amount_kopeks": int(amount_kopeks) if amount_kopeks else 0,
"status": status or "paid",
"external_id": f"donation_{payment_id}"
"external_id": f"donation_{payment_id or 'unknown'}",
"payment_system": "tribute"
}
logger.info(f"✅ Tribute webhook обработан успешно: {result}")
return result
except Exception as e:
logger.error(f"Ошибка обработки Tribute webhook: {e}")
logger.error(f"Webhook payload: {json.dumps(payload, ensure_ascii=False)}")
return None
logger.error(f"Ошибка обработки Tribute webhook: {e}", exc_info=True)
logger.error(f"🔍 Webhook data для отладки: {json.dumps(webhook_data, ensure_ascii=False, indent=2)}")
return None
async def get_payment_status(self, payment_id: str) -> Optional[Dict[str, Any]]:
try:
logger.info(f"Запрос статуса платежа {payment_id}")
return {"status": "unknown", "payment_id": payment_id}
except Exception as e:
logger.error(f"Ошибка получения статуса платежа: {e}")
return None
async def refund_payment(
self,
payment_id: str,
amount_kopeks: Optional[int] = None,
reason: str = "Возврат по запросу"
) -> Optional[Dict[str, Any]]:
try:
logger.info(f"Создание возврата для платежа {payment_id}")
return {"refund_id": f"refund_{payment_id}", "status": "pending"}
except Exception as e:
logger.error(f"Ошибка создания возврата: {e}")
return None
+69 -13
View File
@@ -1,4 +1,5 @@
import logging
import json
from typing import Optional
from aiohttp import web
@@ -26,9 +27,11 @@ class WebhookServer:
self.app.router.add_post(settings.TRIBUTE_WEBHOOK_PATH, self._tribute_webhook_handler)
self.app.router.add_get('/health', self._health_check)
self.app.router.add_options(settings.TRIBUTE_WEBHOOK_PATH, self._options_handler)
logger.info(f"Webhook сервер настроен:")
logger.info(f" - Tribute webhook: {settings.TRIBUTE_WEBHOOK_PATH}")
logger.info(f" - Health check: /health")
logger.info(f" - Tribute webhook: POST {settings.TRIBUTE_WEBHOOK_PATH}")
logger.info(f" - Health check: GET /health")
return self.app
@@ -49,11 +52,11 @@ class WebhookServer:
await self.site.start()
logger.info(f"Webhook сервер запущен на порту {settings.TRIBUTE_WEBHOOK_PORT}")
logger.info(f"🎯 Tribute webhook URL: http://your-server:{settings.TRIBUTE_WEBHOOK_PORT}{settings.TRIBUTE_WEBHOOK_PATH}")
logger.info(f"Tribute webhook сервер запущен на порту {settings.TRIBUTE_WEBHOOK_PORT}")
logger.info(f"🎯 Tribute webhook URL: http://0.0.0.0:{settings.TRIBUTE_WEBHOOK_PORT}{settings.TRIBUTE_WEBHOOK_PATH}")
except Exception as e:
logger.error(f"❌ Ошибка запуска webhook сервера: {e}")
logger.error(f"❌ Ошибка запуска Tribute webhook сервера: {e}")
raise
async def stop(self):
@@ -61,31 +64,82 @@ class WebhookServer:
try:
if self.site:
await self.site.stop()
logger.info("Webhook сайт остановлен")
logger.info("Tribute webhook сайт остановлен")
if self.runner:
await self.runner.cleanup()
logger.info("Webhook runner очищен")
logger.info("Tribute webhook runner очищен")
except Exception as e:
logger.error(f"Ошибка остановки webhook сервера: {e}")
logger.error(f"Ошибка остановки Tribute webhook сервера: {e}")
async def _options_handler(self, request: web.Request) -> web.Response:
return web.Response(
status=200,
headers={
'Access-Control-Allow-Origin': '*',
'Access-Control-Allow-Methods': 'POST, GET, OPTIONS',
'Access-Control-Allow-Headers': 'Content-Type, X-Tribute-Signature',
}
)
async def _tribute_webhook_handler(self, request: web.Request) -> web.Response:
try:
logger.info(f"📥 Получен Tribute webhook: {request.method} {request.path}")
logger.info(f"📋 Headers: {dict(request.headers)}")
raw_body = await request.read()
if not raw_body:
logger.warning("⚠️ Получен пустой webhook от Tribute")
return web.json_response(
{"status": "error", "reason": "empty_body"},
status=400
)
payload = raw_body.decode('utf-8')
logger.info(f"📄 Payload: {payload}")
try:
webhook_data = json.loads(payload)
logger.info(f"📊 Распарсенные данные: {webhook_data}")
except json.JSONDecodeError as e:
logger.error(f"❌ Ошибка парсинга JSON: {e}")
return web.json_response(
{"status": "error", "reason": "invalid_json"},
status=400
)
signature = request.headers.get('X-Tribute-Signature')
logger.info(f"🔐 Signature: {signature}")
if signature and settings.TRIBUTE_WEBHOOK_SECRET:
from app.external.tribute import TributeService as TributeAPI
tribute_api = TributeAPI()
if not tribute_api.verify_webhook_signature(payload, signature):
logger.error("❌ Неверная подпись Tribute webhook")
return web.json_response(
{"status": "error", "reason": "invalid_signature"},
status=401
)
result = await self.tribute_service.process_webhook(payload, signature)
return web.json_response(result, status=200)
if result:
logger.info(f"✅ Tribute webhook обработан успешно: {result}")
return web.json_response({"status": "ok", "result": result}, status=200)
else:
logger.error("❌ Ошибка обработки Tribute webhook")
return web.json_response(
{"status": "error", "reason": "processing_failed"},
status=400
)
except Exception as e:
logger.error(f"Ошибка обработки Tribute webhook: {e}")
logger.error(f"❌ Критическая ошибка обработки Tribute webhook: {e}", exc_info=True)
return web.json_response(
{"status": "error", "reason": "internal_error"},
{"status": "error", "reason": "internal_error", "message": str(e)},
status=500
)
@@ -94,5 +148,7 @@ class WebhookServer:
return web.json_response({
"status": "ok",
"service": "tribute-webhooks",
"tribute_enabled": settings.TRIBUTE_ENABLED
})
"tribute_enabled": settings.TRIBUTE_ENABLED,
"port": settings.TRIBUTE_WEBHOOK_PORT,
"path": settings.TRIBUTE_WEBHOOK_PATH
})
+150 -59
View File
@@ -1,7 +1,9 @@
import asyncio
import logging
import json
import hashlib
import hmac
import base64
from typing import Optional, Dict, Any
from aiohttp import web
from sqlalchemy.ext.asyncio import AsyncSession
@@ -17,13 +19,69 @@ class YooKassaWebhookHandler:
@staticmethod
def verify_webhook_signature(body: str, signature: str, secret: str) -> bool:
expected_signature = hmac.new(
secret.encode('utf-8'),
body.encode('utf-8'),
hashlib.sha256
).hexdigest()
return hmac.compare_digest(signature, expected_signature)
try:
signature_parts = signature.strip().split(' ')
if len(signature_parts) < 4:
logger.error(f"Неверный формат подписи YooKassa: {signature}")
return False
version = signature_parts[0]
payment_id = signature_parts[1]
timestamp = signature_parts[2]
received_signature = signature_parts[3]
if version != "v1":
logger.error(f"Неподдерживаемая версия подписи: {version}")
return False
logger.info(f"Проверка подписи v1 для платежа {payment_id}, timestamp: {timestamp}")
expected_signature_1 = hmac.new(
secret.encode('utf-8'),
body.encode('utf-8'),
hashlib.sha256
).digest()
expected_signature_1_b64 = base64.b64encode(expected_signature_1).decode('utf-8')
signed_payload_2 = f"{payment_id}.{timestamp}.{body}"
expected_signature_2 = hmac.new(
secret.encode('utf-8'),
signed_payload_2.encode('utf-8'),
hashlib.sha256
).digest()
expected_signature_2_b64 = base64.b64encode(expected_signature_2).decode('utf-8')
signed_payload_3 = f"{timestamp}.{body}"
expected_signature_3 = hmac.new(
secret.encode('utf-8'),
signed_payload_3.encode('utf-8'),
hashlib.sha256
).digest()
expected_signature_3_b64 = base64.b64encode(expected_signature_3).decode('utf-8')
logger.debug(f"Получена подпись: {received_signature}")
logger.debug(f"Ожидаемая подпись (вариант 1): {expected_signature_1_b64}")
logger.debug(f"Ожидаемая подпись (вариант 2): {expected_signature_2_b64}")
logger.debug(f"Ожидаемая подпись (вариант 3): {expected_signature_3_b64}")
is_valid = (
hmac.compare_digest(received_signature, expected_signature_1_b64) or
hmac.compare_digest(received_signature, expected_signature_2_b64) or
hmac.compare_digest(received_signature, expected_signature_3_b64)
)
if is_valid:
logger.info("✅ Подпись YooKassa webhook проверена успешно")
else:
logger.warning("⚠️ Подпись YooKassa webhook не совпадает ни с одним вариантом")
return is_valid
except Exception as e:
logger.error(f"Ошибка проверки подписи YooKassa: {e}")
return False
def __init__(self, payment_service: PaymentService):
self.payment_service = payment_service
@@ -31,91 +89,116 @@ class YooKassaWebhookHandler:
async def handle_webhook(self, request: web.Request) -> web.Response:
try:
logger.info(f"📥 Получен YooKassa webhook: {request.method} {request.path}")
logger.info(f"📋 Headers: {dict(request.headers)}")
body = await request.text()
if not body:
logger.warning("Получен пустой webhook от YooKassa")
logger.warning("⚠️ Получен пустой webhook от YooKassa")
return web.Response(status=400, text="Empty body")
if hasattr(settings, 'YOOKASSA_WEBHOOK_SECRET') and settings.YOOKASSA_WEBHOOK_SECRET:
signature = request.headers.get('X-YooKassa-Signature')
if not signature:
logger.warning("Webhook без подписи")
return web.Response(status=400, text="Missing signature")
logger.info(f"📄 Body: {body}")
signature = request.headers.get('Signature') or request.headers.get('X-YooKassa-Signature')
if settings.YOOKASSA_WEBHOOK_SECRET and signature:
logger.info(f"🔐 Получена подпись: {signature}")
if not YooKassaWebhookHandler.verify_webhook_signature(body, signature, settings.YOOKASSA_WEBHOOK_SECRET):
logger.error("Неверная подпись webhook")
return web.Response(status=400, text="Invalid signature")
logger.warning("❌ Подпись не совпала, но продолжаем обработку (режим отладки)")
else:
logger.info("✅ Подпись webhook проверена успешно")
elif settings.YOOKASSA_WEBHOOK_SECRET and not signature:
logger.warning("⚠️ Webhook без подписи, но секрет настроен")
elif signature and not settings.YOOKASSA_WEBHOOK_SECRET:
logger.info("ℹ️ Подпись получена, но проверка отключена (YOOKASSA_WEBHOOK_SECRET не настроен)")
else:
logger.info("ℹ️ Проверка подписи отключена")
try:
webhook_data = json.loads(body)
except json.JSONDecodeError as e:
logger.error(f"Ошибка парсинга JSON webhook YooKassa: {e}")
logger.error(f"Ошибка парсинга JSON webhook YooKassa: {e}")
return web.Response(status=400, text="Invalid JSON")
logger.info(f"Получен webhook YooKassa: {webhook_data.get('event', 'unknown_event')}")
logger.debug(f"Полные данные webhook: {webhook_data}")
logger.info(f"📊 Обработка webhook YooKassa: {webhook_data.get('event', 'unknown_event')}")
logger.debug(f"🔍 Полные данные webhook: {webhook_data}")
event_type = webhook_data.get("event")
if not event_type:
logger.warning("Webhook YooKassa без типа события")
logger.warning("⚠️ Webhook YooKassa без типа события")
return web.Response(status=400, text="No event type")
if event_type not in ["payment.succeeded", "payment.waiting_for_capture"]:
logger.info(f"Игнорируем событие YooKassa: {event_type}")
logger.info(f"Игнорируем событие YooKassa: {event_type}")
return web.Response(status=200, text="OK")
async with get_db() as db:
success = await self.payment_service.process_yookassa_webhook(db, webhook_data)
if success:
logger.info(f"Успешно обработан webhook YooKassa: {event_type}")
return web.Response(status=200, text="OK")
else:
logger.error(f"Ошибка обработки webhook YooKassa: {event_type}")
return web.Response(status=500, text="Processing error")
async for db in get_db():
try:
success = await self.payment_service.process_yookassa_webhook(db, webhook_data)
if success:
logger.info(f"✅ Успешно обработан webhook YooKassa: {event_type}")
return web.Response(status=200, text="OK")
else:
logger.error(f"❌ Ошибка обработки webhook YooKassa: {event_type}")
return web.Response(status=500, text="Processing error")
finally:
await db.close()
except Exception as e:
logger.error(f"Критическая ошибка обработки webhook YooKassa: {e}", exc_info=True)
logger.error(f"Критическая ошибка обработки webhook YooKassa: {e}", exc_info=True)
return web.Response(status=500, text="Internal server error")
def setup_routes(self, app: web.Application) -> None:
webhook_path = settings.YOOKASSA_WEBHOOK_PATH
app.router.add_post(webhook_path, self.handle_webhook)
app.router.add_get(webhook_path, self._get_handler)
app.router.add_options(webhook_path, self._options_handler)
logger.info(f"Настроен webhook YooKassa на пути: {webhook_path}")
logger.info(f"Настроен YooKassa webhook на пути: POST {webhook_path}")
async def _get_handler(self, request: web.Request) -> web.Response:
return web.json_response({
"status": "ok",
"message": "YooKassa webhook endpoint is working",
"method": "GET",
"path": request.path,
"note": "Use POST method for actual webhooks"
})
async def _options_handler(self, request: web.Request) -> web.Response:
return web.Response(
status=200,
headers={
'Access-Control-Allow-Origin': '*',
'Access-Control-Allow-Methods': 'POST, GET, OPTIONS',
'Access-Control-Allow-Headers': 'Content-Type, X-YooKassa-Signature',
}
)
def create_yookassa_webhook_app(payment_service: PaymentService) -> web.Application:
app = web.Application()
async def logging_middleware(request, handler):
start_time = request.loop.time()
try:
response = await handler(request)
process_time = request.loop.time() - start_time
logger.info(f"YooKassa webhook {request.method} {request.path_qs} "
f"-> {response.status} ({process_time:.3f}s)")
return response
except Exception as e:
process_time = request.loop.time() - start_time
logger.error(f"YooKassa webhook {request.method} {request.path_qs} "
f"-> ERROR ({process_time:.3f}s): {e}")
raise
app.middlewares.append(logging_middleware)
webhook_handler = YooKassaWebhookHandler(payment_service)
webhook_handler.setup_routes(app)
async def health_check(request):
return web.json_response({"status": "ok", "service": "yookassa_webhook"})
return web.json_response({
"status": "ok",
"service": "yookassa_webhook",
"port": settings.YOOKASSA_WEBHOOK_PORT,
"path": settings.YOOKASSA_WEBHOOK_PATH,
"enabled": settings.is_yookassa_enabled()
})
app.router.add_get("/health", health_check)
@@ -125,12 +208,10 @@ def create_yookassa_webhook_app(payment_service: PaymentService) -> web.Applicat
async def start_yookassa_webhook_server(payment_service: PaymentService) -> None:
if not settings.is_yookassa_enabled():
logger.info("YooKassa отключен, webhook сервер не запускается")
logger.info("YooKassa отключена, webhook сервер не запускается")
return
try:
from aiohttp import web
app = create_yookassa_webhook_app(payment_service)
runner = web.AppRunner(app)
@@ -138,15 +219,25 @@ async def start_yookassa_webhook_server(payment_service: PaymentService) -> None
site = web.TCPSite(
runner,
host='0.0.0.0',
host='0.0.0.0',
port=settings.YOOKASSA_WEBHOOK_PORT
)
await site.start()
logger.info(f"YooKassa webhook сервер запущен на порту {settings.YOOKASSA_WEBHOOK_PORT}")
logger.info(f"Webhook URL: http://localhost:{settings.YOOKASSA_WEBHOOK_PORT}{settings.YOOKASSA_WEBHOOK_PATH}")
logger.info(f"YooKassa webhook сервер запущен на порту {settings.YOOKASSA_WEBHOOK_PORT}")
logger.info(f"🎯 YooKassa webhook URL: http://0.0.0.0:{settings.YOOKASSA_WEBHOOK_PORT}{settings.YOOKASSA_WEBHOOK_PATH}")
try:
while True:
await asyncio.sleep(1)
except asyncio.CancelledError:
logger.info("🛑 YooKassa webhook сервер получил сигнал остановки")
finally:
await site.stop()
await runner.cleanup()
logger.info("✅ YooKassa webhook сервер остановлен")
except Exception as e:
logger.error(f"Ошибка запуска YooKassa webhook сервера: {e}", exc_info=True)
logger.error(f"Ошибка запуска YooKassa webhook сервера: {e}", exc_info=True)
raise
+235
View File
@@ -0,0 +1,235 @@
import logging
from aiogram import Dispatcher, types, F
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
from app.database.models import User
from app.services.maintenance_service import maintenance_service
from app.keyboards.admin import get_maintenance_keyboard, get_admin_main_keyboard
from app.localization.texts import get_texts
from app.utils.decorators import admin_required, error_handler
logger = logging.getLogger(__name__)
class MaintenanceStates(StatesGroup):
waiting_for_reason = State()
@admin_required
@error_handler
async def show_maintenance_panel(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession,
state: FSMContext
):
texts = get_texts(db_user.language)
status_info = maintenance_service.get_status_info()
status_emoji = "🔧" if status_info["is_active"] else ""
status_text = "Включен" if status_info["is_active"] else "Выключен"
api_emoji = "" if status_info["api_status"] else ""
api_text = "Доступно" if status_info["api_status"] else "Недоступно"
monitoring_emoji = "🔄" if status_info["monitoring_active"] else "⏹️"
monitoring_text = "Запущен" if status_info["monitoring_active"] else "Остановлен"
enabled_info = ""
if status_info["is_active"] and status_info["enabled_at"]:
enabled_time = status_info["enabled_at"].strftime("%d.%m.%Y %H:%M:%S")
enabled_info = f"\n📅 <b>Включен:</b> {enabled_time}"
if status_info["reason"]:
enabled_info += f"\n📝 <b>Причина:</b> {status_info['reason']}"
last_check_info = ""
if status_info["last_check"]:
last_check_time = status_info["last_check"].strftime("%H:%M:%S")
last_check_info = f"\n🕐 <b>Последняя проверка:</b> {last_check_time}"
failures_info = ""
if status_info["consecutive_failures"] > 0:
failures_info = f"\n⚠️ <b>Неудачных проверок подряд:</b> {status_info['consecutive_failures']}"
message_text = f"""
🔧 <b>Режим технических работ</b>
{status_emoji} <b>Статус:</b> {status_text}
{api_emoji} <b>API RemnaWave:</b> {api_text}
{monitoring_emoji} <b>Мониторинг:</b> {monitoring_text}
⏱️ <b>Интервал проверки:</b> {status_info['check_interval']}с
🤖 <b>Автовключение:</b> {'Включено' if status_info['auto_enable_configured'] else 'Отключено'}
{enabled_info}
{last_check_info}
{failures_info}
️ <i>В режиме техработ обычные пользователи не могут использовать бота. Администраторы имеют полный доступ.</i>
"""
await callback.message.edit_text(
message_text,
reply_markup=get_maintenance_keyboard(db_user.language, status_info["is_active"], status_info["monitoring_active"])
)
await callback.answer()
@admin_required
@error_handler
async def toggle_maintenance_mode(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession,
state: FSMContext
):
is_active = maintenance_service.is_maintenance_active()
if is_active:
success = await maintenance_service.disable_maintenance()
if success:
await callback.answer("Режим техработ выключен", show_alert=True)
else:
await callback.answer("Ошибка выключения режима техработ", show_alert=True)
else:
await state.set_state(MaintenanceStates.waiting_for_reason)
await callback.message.edit_text(
"🔧 <b>Включение режима техработ</b>\n\nВведите причину включения техработ или отправьте /skip для пропуска:",
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=[
[types.InlineKeyboardButton(text="❌ Отмена", callback_data="maintenance_panel")]
])
)
await callback.answer()
@admin_required
@error_handler
async def process_maintenance_reason(
message: types.Message,
db_user: User,
db: AsyncSession,
state: FSMContext
):
current_state = await state.get_state()
if current_state != MaintenanceStates.waiting_for_reason:
return
reason = None
if message.text and message.text != "/skip":
reason = message.text[:200]
success = await maintenance_service.enable_maintenance(reason=reason, auto=False)
if success:
response_text = "Режим техработ включен"
if reason:
response_text += f"\nПричина: {reason}"
else:
response_text = "Ошибка включения режима техработ"
await message.answer(response_text)
await state.clear()
status_info = maintenance_service.get_status_info()
await message.answer(
"Вернуться к панели управления техработами:",
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=[
[types.InlineKeyboardButton(text="🔧 Панель техработ", callback_data="maintenance_panel")]
])
)
@admin_required
@error_handler
async def toggle_monitoring(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
status_info = maintenance_service.get_status_info()
if status_info["monitoring_active"]:
success = await maintenance_service.stop_monitoring()
message = "Мониторинг остановлен" if success else "Ошибка остановки мониторинга"
else:
success = await maintenance_service.start_monitoring()
message = "Мониторинг запущен" if success else "Ошибка запуска мониторинга"
await callback.answer(message, show_alert=True)
await show_maintenance_panel(callback, db_user, db, None)
@admin_required
@error_handler
async def force_api_check(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
await callback.answer("Проверка API...", show_alert=False)
check_result = await maintenance_service.force_api_check()
if check_result["success"]:
status_text = "доступно" if check_result["api_available"] else "недоступно"
message = f"API {status_text}\nВремя ответа: {check_result['response_time']}с"
else:
message = f"Ошибка проверки: {check_result.get('error', 'Неизвестная ошибка')}"
await callback.message.answer(message)
await show_maintenance_panel(callback, db_user, db, None)
@admin_required
@error_handler
async def back_to_admin_panel(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
texts = get_texts(db_user.language)
await callback.message.edit_text(
texts.ADMIN_PANEL,
reply_markup=get_admin_main_keyboard(db_user.language)
)
await callback.answer()
def register_handlers(dp: Dispatcher):
dp.callback_query.register(
show_maintenance_panel,
F.data == "maintenance_panel"
)
dp.callback_query.register(
toggle_maintenance_mode,
F.data == "maintenance_toggle"
)
dp.callback_query.register(
toggle_monitoring,
F.data == "maintenance_monitoring"
)
dp.callback_query.register(
force_api_check,
F.data == "maintenance_check_api"
)
dp.callback_query.register(
back_to_admin_panel,
F.data == "admin_panel"
)
dp.message.register(
process_maintenance_reason,
MaintenanceStates.waiting_for_reason
)
+106 -56
View File
@@ -1322,53 +1322,22 @@ async def show_sync_options(
text = """
🔄 <b>Синхронизация с RemnaWave</b>
Выберите тип синхронизации:
🔄 <b>Синхронизировать всех</b>
Полная синхронизация всех пользователей
• Создание новых пользователей из панели
• Обновление данных существующих
• Удаление неактуальных подписок
🔄 <b>Полная синхронизация выполняет:</b>
• Создание новых пользователей из панели в боте
• Обновление данных существующих пользователей
Деактивация подписок пользователей, отсутствующих в панели
• Сохранение балансов пользователей
• ⏱️ Время выполнения: 2-5 минут
🆕 <b>Только новых</b>
• Создание пользователей из панели, которых нет в боте
• Быстрое добавление при массовой регистрации
• ⏱️ Время выполнения: 30 секунд - 2 минуты
📈 <b>Обновить данные</b>
• Обновление информации о трафике и подписках
• Синхронизация статуса и лимитов
• Обновление подключенных сквадов
• ⏱️ Время выполнения: 1-3 минуты
🔍 <b>Валидация подписок</b>
• Проверка и исправление проблем в данных
• Восстановление отсутствующих полей
• Исправление некорректных статусов
🧹 <b>Мягкая очистка</b>
• Деактивация подписок отсутствующих в панели
• Сохранение транзакций и истории
🗑️ <b>ПРИНУДИТЕЛЬНАЯ ОЧИСТКА</b>
• ⚠️ ОПАСНО: Полное удаление данных пользователей
• Удаление транзакций, балансов, рефералов
• Только при серьезных проблемах синхронизации
⚠️ <b>Важно:</b>
• Во время синхронизации не выполняйте другие операции
• При полной синхронизации подписки пользователей, отсутствующих в панели, будут деактивированы
• Рекомендуется делать полную синхронизацию ежедневно
• Баланс пользователей НЕ удаляется
"""
keyboard = [
[types.InlineKeyboardButton(text="🔄 Синхронизировать всех", callback_data="sync_all_users")],
[types.InlineKeyboardButton(text="🆕 Только новых", callback_data="sync_new_users")],
[types.InlineKeyboardButton(text="📈 Обновить данные", callback_data="sync_update_data")],
[types.InlineKeyboardButton(text="🔍 Валидация подписок", callback_data="sync_validate")],
[types.InlineKeyboardButton(text="🧹 Мягкая очистка", callback_data="sync_cleanup")],
[types.InlineKeyboardButton(text="🗑️ ПРИНУДИТЕЛЬНАЯ ОЧИСТКА", callback_data="confirm_force_cleanup")],
[types.InlineKeyboardButton(text="🔄 Запустить полную синхронизацию", callback_data="sync_all_users")],
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
]
@@ -1378,6 +1347,104 @@ async def show_sync_options(
)
await callback.answer()
@admin_required
@error_handler
async def sync_all_users(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
"""Выполняет полную синхронизацию всех пользователей"""
progress_text = """
🔄 <b>Выполняется полная синхронизация...</b>
📋 Этапы:
• Загрузка ВСЕХ пользователей из панели RemnaWave
• Создание новых пользователей в боте
• Обновление существующих пользователей
• Деактивация подписок отсутствующих пользователей
• Сохранение балансов
⏳ Пожалуйста, подождите...
"""
await callback.message.edit_text(progress_text, reply_markup=None)
remnawave_service = RemnaWaveService()
stats = await remnawave_service.sync_users_from_panel(db, "all")
total_operations = stats['created'] + stats['updated'] + stats.get('deleted', 0)
if stats['errors'] == 0:
status_emoji = ""
status_text = "успешно завершена"
elif stats['errors'] < total_operations:
status_emoji = "⚠️"
status_text = "завершена с предупреждениями"
else:
status_emoji = ""
status_text = "завершена с ошибками"
text = f"""
{status_emoji} <b>Полная синхронизация {status_text}</b>
📊 <b>Результат:</b>
• 🆕 Создано: {stats['created']}
• 🔄 Обновлено: {stats['updated']}
• 🗑️ Деактивировано: {stats.get('deleted', 0)}
• ❌ Ошибок: {stats['errors']}
"""
if stats.get('deleted', 0) > 0:
text += f"""
🗑️ <b>Деактивированные подписки:</b>
Деактивированы подписки пользователей, которые
отсутствуют в панели RemnaWave.
💰 Балансы пользователей сохранены.
"""
if stats['errors'] > 0:
text += f"""
⚠️ <b>Внимание:</b>
Некоторые операции завершились с ошибками.
Проверьте логи для получения подробной информации.
"""
text += f"""
💡 <b>Рекомендации:</b>
• Полная синхронизация выполнена
• Рекомендуется запускать раз в день
• Все пользователи из панели синхронизированы
"""
keyboard = []
if stats['errors'] > 0:
keyboard.append([
types.InlineKeyboardButton(
text="🔄 Повторить синхронизацию",
callback_data="sync_all_users"
)
])
keyboard.extend([
[
types.InlineKeyboardButton(text="📊 Статистика системы", callback_data="admin_rw_system"),
types.InlineKeyboardButton(text="🌐 Ноды", callback_data="admin_rw_nodes")
],
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
])
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard)
)
await callback.answer()
@admin_required
@error_handler
async def show_sync_recommendations(
@@ -1844,39 +1911,22 @@ def register_handlers(dp: Dispatcher):
dp.callback_query.register(manage_node, F.data.startswith("node_restart_"))
dp.callback_query.register(restart_all_nodes, F.data == "admin_restart_all_nodes")
dp.callback_query.register(show_sync_options, F.data == "admin_rw_sync")
dp.callback_query.register(sync_users, F.data.startswith("sync_"))
dp.callback_query.register(show_sync_recommendations, F.data == "sync_recommendations")
dp.callback_query.register(validate_subscriptions, F.data == "sync_validate")
dp.callback_query.register(cleanup_subscriptions, F.data == "sync_cleanup")
dp.callback_query.register(confirm_force_cleanup, F.data == "confirm_force_cleanup")
dp.callback_query.register(force_cleanup_all_orphaned_users, F.data == "force_cleanup_orphaned")
dp.callback_query.register(sync_all_users, F.data == "sync_all_users")
dp.callback_query.register(show_squads_management, F.data == "admin_rw_squads")
dp.callback_query.register(show_squad_details, F.data.startswith("admin_squad_manage_"))
dp.callback_query.register(manage_squad_action, F.data.startswith("squad_add_users_"))
dp.callback_query.register(manage_squad_action, F.data.startswith("squad_remove_users_"))
dp.callback_query.register(manage_squad_action, F.data.startswith("squad_delete_"))
dp.callback_query.register(show_squad_edit_menu, F.data.startswith("squad_edit_") & ~F.data.startswith("squad_edit_inbounds_"))
dp.callback_query.register(show_squad_inbounds_selection, F.data.startswith("squad_edit_inbounds_"))
dp.callback_query.register(show_squad_rename_form, F.data.startswith("squad_rename_"))
dp.callback_query.register(cancel_squad_rename, F.data.startswith("cancel_rename_"))
dp.callback_query.register(toggle_squad_inbound, F.data.startswith("sqd_tgl_"))
dp.callback_query.register(save_squad_inbounds, F.data.startswith("sqd_save_"))
dp.callback_query.register(show_squad_edit_menu_short, F.data.startswith("sqd_edit_"))
dp.callback_query.register(start_squad_creation, F.data == "admin_squad_create")
dp.callback_query.register(cancel_squad_creation, F.data == "cancel_squad_create")
dp.callback_query.register(toggle_create_inbound, F.data.startswith("create_tgl_"))
dp.callback_query.register(finish_squad_creation, F.data == "create_squad_finish")
dp.message.register(
+8 -9
View File
@@ -320,11 +320,11 @@ async def process_topup_amount(
amount_rubles = float(message.text.replace(',', '.'))
if amount_rubles < 1:
await message.answer("Минимальная сумма пополнения: 1 ₽")
await message.answer("Минимальная сумма пополнения: 1 ₽")
return
if amount_rubles > 50000:
await message.answer("Максимальная сумма пополнения: 50,000 ₽")
await message.answer("Максимальная сумма пополнения: 50,000 ₽")
return
amount_kopeks = int(amount_rubles * 100)
@@ -334,11 +334,11 @@ async def process_topup_amount(
if payment_method == "stars":
await process_stars_payment_amount(message, db_user, amount_kopeks, state)
elif payment_method == "yookassa":
from app.database.database import get_db
async with get_db() as db:
from app.database.database import AsyncSessionLocal
async with AsyncSessionLocal() as db:
await process_yookassa_payment_amount(message, db_user, db, amount_kopeks, state)
else:
await message.answer("Неизвестный способ оплаты")
await message.answer("Неизвестный способ оплаты")
except ValueError:
await message.answer(
@@ -346,7 +346,6 @@ async def process_topup_amount(
reply_markup=get_back_keyboard(db_user.language)
)
@error_handler
async def process_stars_payment_amount(
message: types.Message,
@@ -501,14 +500,14 @@ async def check_yookassa_payment_status(
emoji = status_emoji.get(payment.status, "")
status = status_text.get(payment.status, "Неизвестно")
message_text = (f"💳 <b>Статус платежа</b>\n\n"
message_text = (f"💳 Статус платежа:\n\n"
f"🆔 ID: {payment.yookassa_payment_id[:8]}...\n"
f"💰 Сумма: {settings.format_price(payment.amount_kopeks)}\n"
f"📊 Статус: {emoji} {status}\n"
f"📅 Создан: {payment.created_at.strftime('%d.%m.%Y %H:%M')}\n")
if payment.is_succeeded:
message_text += "\n✅ Платеж успешно завершен!\nСредства зачислены на баланс."
message_text += "\n✅ Платеж успешно завершен!\n\nСредства зачислены на баланс."
elif payment.is_pending:
message_text += "\n⏳ Платеж ожидает оплаты. Нажмите кнопку 'Оплатить' выше."
elif payment.is_failed:
@@ -571,4 +570,4 @@ def register_handlers(dp: Dispatcher):
dp.message.register(
process_topup_amount,
BalanceStates.waiting_for_amount
)
)
+31 -10
View File
@@ -667,36 +667,57 @@ async def handle_extend_subscription(
db_user: User,
db: AsyncSession
):
texts = get_texts(db_user.language)
subscription = db_user.subscription
if not subscription or subscription.is_trial:
await callback.answer(" Продление доступно только для платных подписок", show_alert=True)
await callback.answer(" Продление доступно только для платных подписок", show_alert=True)
return
if subscription.days_left > 3:
await callback.answer(" Продление доступно за 3 дня до окончания подписки", show_alert=True)
await callback.answer(" Продление доступно за 3 дня до окончания подписки", show_alert=True)
return
subscription_service = SubscriptionService()
available_periods = settings.get_available_renewal_periods()
renewal_prices = {}
for days in [30, 90, 180]:
price = await subscription_service.calculate_renewal_price(subscription, days, db)
renewal_prices[days] = price
for days in available_periods:
try:
price = await subscription_service.calculate_renewal_price(subscription, days, db)
renewal_prices[days] = price
except Exception as e:
logger.error(f"Ошибка расчета цены для периода {days}: {e}")
continue
if not renewal_prices:
await callback.answer("⌛ Нет доступных периодов для продления", show_alert=True)
return
prices_text = ""
period_display = {
14: "14 дней",
30: "30 дней",
60: "60 дней",
90: "90 дней",
180: "180 дней",
360: "360 дней"
}
for days in available_periods:
if days in renewal_prices and days in period_display:
prices_text += f"📅 {period_display[days]} - {texts.format_price(renewal_prices[days])}\n"
await callback.message.edit_text(
f"<b>Продление подписки</b>\n\n"
f"⏰ Продление подписки\n\n"
f"Осталось дней: {subscription.days_left}\n\n"
f"<b>Ваша текущая конфигурация:</b>\n"
f"🌍 Серверов: {len(subscription.connected_squads)}\n"
f"📊 Трафик: {texts.format_traffic(subscription.traffic_limit_gb)}\n"
f"📱 Устройств: {subscription.device_limit}\n\n"
f"<b>Выберите период продления:</b>\n"
f"📅 30 дней - {texts.format_price(renewal_prices[30])}\n"
f"📅 90 дней - {texts.format_price(renewal_prices[90])}\n"
f"📅 180 дней - {texts.format_price(renewal_prices[180])}\n\n"
f"{prices_text.rstrip()}\n\n"
f"💡 <i>Цена включает все ваши текущие серверы и настройки</i>",
reply_markup=get_extend_subscription_keyboard_with_prices(db_user.language, renewal_prices),
parse_mode="HTML"
+33 -1
View File
@@ -25,7 +25,8 @@ def get_admin_main_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
InlineKeyboardButton(text=texts.ADMIN_REMNAWAVE, callback_data="admin_remnawave")
],
[
InlineKeyboardButton(text=texts.ADMIN_STATISTICS, callback_data="admin_statistics")
InlineKeyboardButton(text=texts.ADMIN_STATISTICS, callback_data="admin_statistics"),
InlineKeyboardButton(text="🔧 Техработы", callback_data="maintenance_panel")
],
[
InlineKeyboardButton(text=texts.BACK, callback_data="back_to_menu")
@@ -575,3 +576,34 @@ def get_admin_pagination_keyboard(
])
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_maintenance_keyboard(language: str = "ru", is_active: bool = False, monitoring_active: bool = False) -> InlineKeyboardMarkup:
"""Клавиатура для управления техработами"""
if language == "en":
toggle_text = "🔴 Disable maintenance" if is_active else "🔧 Enable maintenance"
monitoring_text = "⏹️ Stop monitoring" if monitoring_active else "🔄 Start monitoring"
check_api_text = "🔍 Check API"
back_text = "⬅️ Back to admin"
else:
toggle_text = "🔴 Выключить техработы" if is_active else "🔧 Включить техработы"
monitoring_text = "⏹️ Остановить мониторинг" if monitoring_active else "🔄 Запустить мониторинг"
check_api_text = "🔍 Проверить API"
back_text = "⬅️ Назад в админку"
keyboard = [
[InlineKeyboardButton(text=toggle_text, callback_data="maintenance_toggle")],
[InlineKeyboardButton(text=monitoring_text, callback_data="maintenance_monitoring")],
[InlineKeyboardButton(text=check_api_text, callback_data="maintenance_check_api")],
[InlineKeyboardButton(text=back_text, callback_data="admin_panel")]
]
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_sync_simplified_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
keyboard = [
[InlineKeyboardButton(text="🔄 Полная синхронизация", callback_data="sync_all_users")],
[InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
]
return InlineKeyboardMarkup(inline_keyboard=keyboard)
+43 -35
View File
@@ -187,19 +187,25 @@ def get_subscription_period_keyboard(language: str = "ru") -> InlineKeyboardMark
texts = get_texts(language)
keyboard = []
periods = [
(14, texts.PERIOD_14_DAYS),
(30, texts.PERIOD_30_DAYS),
(60, texts.PERIOD_60_DAYS),
(90, texts.PERIOD_90_DAYS),
(180, texts.PERIOD_180_DAYS),
(360, texts.PERIOD_360_DAYS)
]
available_periods = settings.get_available_subscription_periods()
for days, text in periods:
keyboard.append([
InlineKeyboardButton(text=text, callback_data=f"period_{days}")
])
period_texts = {
14: texts.PERIOD_14_DAYS,
30: texts.PERIOD_30_DAYS,
60: texts.PERIOD_60_DAYS,
90: texts.PERIOD_90_DAYS,
180: texts.PERIOD_180_DAYS,
360: texts.PERIOD_360_DAYS
}
for days in available_periods:
if days in period_texts:
keyboard.append([
InlineKeyboardButton(
text=period_texts[days],
callback_data=f"period_{days}"
)
])
keyboard.append([
InlineKeyboardButton(text=texts.BACK, callback_data="back_to_menu")
@@ -869,29 +875,31 @@ def get_specific_app_keyboard(
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_extend_subscription_keyboard_with_prices(language: str, prices: dict) -> InlineKeyboardMarkup:
from app.localization.texts import get_texts
texts = get_texts(language)
keyboard = []
return InlineKeyboardMarkup(inline_keyboard=[
[
InlineKeyboardButton(
text=f"📅 30 дней - {texts.format_price(prices[30])}",
callback_data="extend_period_30"
)
],
[
InlineKeyboardButton(
text=f"📅 90 дней - {texts.format_price(prices[90])}",
callback_data="extend_period_90"
)
],
[
InlineKeyboardButton(
text=f"📅 180 дней - {texts.format_price(prices[180])}",
callback_data="extend_period_180"
)
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
]
available_periods = settings.get_available_renewal_periods()
period_display = {
14: "14 дней",
30: "30 дней",
60: "60 дней",
90: "90 дней",
180: "180 дней",
360: "360 дней"
}
for days in available_periods:
if days in prices and days in period_display:
keyboard.append([
InlineKeyboardButton(
text=f"📅 {period_display[days]} - {texts.format_price(prices[days])}",
callback_data=f"extend_period_{days}"
)
])
keyboard.append([
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
])
return InlineKeyboardMarkup(inline_keyboard=keyboard)
+21
View File
@@ -280,6 +280,27 @@ class RussianTexts(Texts):
• Поддержка до 3 устройств
⚡️ Успейте оформить до окончания тестового периода!
"""
MAINTENANCE_MODE_ACTIVE = """
🔧 Технические работы!
Сервис временно недоступен. Ведутся технические работы по улучшению качества обслуживания.
⏰ Ориентировочное время завершения: неизвестно
🔄 Попробуйте позже
Приносим извинения за временные неудобства.
"""
MAINTENANCE_MODE_API_ERROR = """
🔧 Технические работы!
Сервис временно недоступен из-за проблем с подключением к серверам.
⏰ Мы работаем над восстановлением. Попробуйте через несколько минут.
🔄 Последняя проверка: {last_check}
"""
SUBSCRIPTION_EXPIRING_PAID = """
+45
View File
@@ -0,0 +1,45 @@
import logging
from typing import Callable, Dict, Any, Awaitable
from aiogram import BaseMiddleware
from aiogram.types import Message, CallbackQuery, TelegramObject, User as TgUser
from app.config import settings
from app.services.maintenance_service import maintenance_service
logger = logging.getLogger(__name__)
class MaintenanceMiddleware(BaseMiddleware):
async def __call__(
self,
handler: Callable[[TelegramObject, Dict[str, Any]], Awaitable[Any]],
event: TelegramObject,
data: Dict[str, Any]
) -> Any:
user: TgUser = None
if isinstance(event, (Message, CallbackQuery)):
user = event.from_user
if not user or user.is_bot:
return await handler(event, data)
if not maintenance_service.is_maintenance_active():
return await handler(event, data)
if settings.is_admin(user.id):
return await handler(event, data)
maintenance_message = maintenance_service.get_maintenance_message()
try:
if isinstance(event, Message):
await event.answer(maintenance_message, parse_mode="HTML")
elif isinstance(event, CallbackQuery):
await event.answer(maintenance_message, show_alert=True)
except Exception as e:
logger.error(f"Ошибка отправки сообщения о техработах пользователю {user.id}: {e}")
logger.info(f"🔧 Пользователь {user.id} заблокирован во время техработ")
return
+268
View File
@@ -0,0 +1,268 @@
import asyncio
import logging
from datetime import datetime, timedelta
from typing import Optional, Dict, Any
from dataclasses import dataclass
from app.config import settings
from app.external.remnawave_api import RemnaWaveAPI, test_api_connection
from app.utils.cache import cache
logger = logging.getLogger(__name__)
@dataclass
class MaintenanceStatus:
is_active: bool
enabled_at: Optional[datetime] = None
last_check: Optional[datetime] = None
reason: Optional[str] = None
auto_enabled: bool = False
api_status: bool = True
consecutive_failures: int = 0
class MaintenanceService:
def __init__(self):
self._status = MaintenanceStatus(is_active=False)
self._check_task: Optional[asyncio.Task] = None
self._is_checking = False
self._max_consecutive_failures = 3
@property
def status(self) -> MaintenanceStatus:
return self._status
def is_maintenance_active(self) -> bool:
return self._status.is_active
def get_maintenance_message(self) -> str:
if self._status.auto_enabled:
return f"""
🔧 Технические работы!
Сервис временно недоступен из-за проблем с подключением к серверам.
⏰ Мы работаем над восстановлением. Попробуйте через несколько минут.
🔄 Последняя проверка: {self._status.last_check.strftime('%H:%M:%S') if self._status.last_check else 'неизвестно'}
"""
else:
return settings.get_maintenance_message()
async def enable_maintenance(self, reason: Optional[str] = None, auto: bool = False) -> bool:
try:
if self._status.is_active:
logger.warning("Режим техработ уже включен")
return True
self._status.is_active = True
self._status.enabled_at = datetime.utcnow()
self._status.reason = reason or ("Автоматическое включение" if auto else "Включено администратором")
self._status.auto_enabled = auto
await self._save_status_to_cache()
logger.warning(f"🔧 Режим техработ ВКЛЮЧЕН. Причина: {self._status.reason}")
return True
except Exception as e:
logger.error(f"Ошибка включения режима техработ: {e}")
return False
async def disable_maintenance(self) -> bool:
try:
if not self._status.is_active:
logger.info("Режим техработ уже выключен")
return True
self._status.is_active = False
self._status.enabled_at = None
self._status.reason = None
self._status.auto_enabled = False
self._status.consecutive_failures = 0
await self._save_status_to_cache()
logger.info("✅ Режим техработ ВЫКЛЮЧЕН")
return True
except Exception as e:
logger.error(f"Ошибка выключения режима техработ: {e}")
return False
async def start_monitoring(self) -> bool:
try:
if self._check_task and not self._check_task.done():
logger.warning("Мониторинг уже запущен")
return True
await self._load_status_from_cache()
self._check_task = asyncio.create_task(self._monitoring_loop())
logger.info(f"🔄 Запущен мониторинг API RemnaWave (интервал: {settings.get_maintenance_check_interval()}с)")
return True
except Exception as e:
logger.error(f"Ошибка запуска мониторинга: {e}")
return False
async def stop_monitoring(self) -> bool:
try:
if self._check_task and not self._check_task.done():
self._check_task.cancel()
try:
await self._check_task
except asyncio.CancelledError:
pass
logger.info("⏹️ Мониторинг API остановлен")
return True
except Exception as e:
logger.error(f"Ошибка остановки мониторинга: {e}")
return False
async def check_api_status(self) -> bool:
try:
if self._is_checking:
return self._status.api_status
self._is_checking = True
self._status.last_check = datetime.utcnow()
api = RemnaWaveAPI(settings.REMNAWAVE_API_URL, settings.REMNAWAVE_API_KEY)
async with api:
is_connected = await test_api_connection(api)
if is_connected:
self._status.api_status = True
self._status.consecutive_failures = 0
if self._status.is_active and self._status.auto_enabled:
await self.disable_maintenance()
logger.info("✅ API восстановился, режим техработ автоматически отключен")
return True
else:
self._status.api_status = False
self._status.consecutive_failures += 1
if (self._status.consecutive_failures >= self._max_consecutive_failures and
not self._status.is_active and
settings.is_maintenance_auto_enable()):
await self.enable_maintenance(
reason=f"Автоматическое включение после {self._status.consecutive_failures} неудачных проверок API",
auto=True
)
return False
except Exception as e:
logger.error(f"Ошибка проверки API: {e}")
self._status.api_status = False
self._status.consecutive_failures += 1
return False
finally:
self._is_checking = False
await self._save_status_to_cache()
async def _monitoring_loop(self):
while True:
try:
await self.check_api_status()
await asyncio.sleep(settings.get_maintenance_check_interval())
except asyncio.CancelledError:
logger.info("Мониторинг отменен")
break
except Exception as e:
logger.error(f"Ошибка в цикле мониторинга: {e}")
await asyncio.sleep(30)
async def _save_status_to_cache(self):
try:
status_data = {
"is_active": self._status.is_active,
"enabled_at": self._status.enabled_at.isoformat() if self._status.enabled_at else None,
"reason": self._status.reason,
"auto_enabled": self._status.auto_enabled,
"consecutive_failures": self._status.consecutive_failures,
"last_check": self._status.last_check.isoformat() if self._status.last_check else None
}
await cache.set("maintenance_status", status_data, expire=3600)
except Exception as e:
logger.error(f"Ошибка сохранения состояния в кеш: {e}")
async def _load_status_from_cache(self):
try:
status_data = await cache.get("maintenance_status")
if not status_data:
return
self._status.is_active = status_data.get("is_active", False)
self._status.reason = status_data.get("reason")
self._status.auto_enabled = status_data.get("auto_enabled", False)
self._status.consecutive_failures = status_data.get("consecutive_failures", 0)
if status_data.get("enabled_at"):
self._status.enabled_at = datetime.fromisoformat(status_data["enabled_at"])
if status_data.get("last_check"):
self._status.last_check = datetime.fromisoformat(status_data["last_check"])
logger.info(f"📥 Состояние техработ загружено из кеша: активен={self._status.is_active}")
except Exception as e:
logger.error(f"Ошибка загрузки состояния из кеша: {e}")
def get_status_info(self) -> Dict[str, Any]:
return {
"is_active": self._status.is_active,
"enabled_at": self._status.enabled_at,
"last_check": self._status.last_check,
"reason": self._status.reason,
"auto_enabled": self._status.auto_enabled,
"api_status": self._status.api_status,
"consecutive_failures": self._status.consecutive_failures,
"monitoring_active": self._check_task is not None and not self._check_task.done(),
"auto_enable_configured": settings.is_maintenance_auto_enable(),
"check_interval": settings.get_maintenance_check_interval()
}
async def force_api_check(self) -> Dict[str, Any]:
start_time = datetime.utcnow()
try:
api_status = await self.check_api_status()
end_time = datetime.utcnow()
response_time = (end_time - start_time).total_seconds()
return {
"success": True,
"api_available": api_status,
"response_time": round(response_time, 2),
"checked_at": end_time,
"consecutive_failures": self._status.consecutive_failures
}
except Exception as e:
end_time = datetime.utcnow()
response_time = (end_time - start_time).total_seconds()
return {
"success": False,
"api_available": False,
"error": str(e),
"response_time": round(response_time, 2),
"checked_at": end_time,
"consecutive_failures": self._status.consecutive_failures
}
maintenance_service = MaintenanceService()
+35 -26
View File
@@ -11,7 +11,7 @@ from app.config import settings
from app.services.yookassa_service import YooKassaService
from app.database.crud.yookassa import create_yookassa_payment, link_yookassa_payment_to_transaction
from app.database.crud.transaction import create_transaction
from app.database.crud.user import add_user_balance
from app.database.crud.user import add_user_balance, get_user_by_id
from app.database.models import TransactionType, PaymentMethod
logger = logging.getLogger(__name__)
@@ -93,11 +93,13 @@ class PaymentService:
yookassa_created_at = None
if yookassa_response.get("created_at"):
try:
yookassa_created_at = datetime.fromisoformat(
dt_with_tz = datetime.fromisoformat(
yookassa_response["created_at"].replace('Z', '+00:00')
)
except:
pass
yookassa_created_at = dt_with_tz.replace(tzinfo=None)
except Exception as e:
logger.warning(f"Не удалось парсить created_at: {e}")
yookassa_created_at = None
local_payment = await create_yookassa_payment(
db=db,
@@ -110,7 +112,7 @@ class PaymentService:
confirmation_url=yookassa_response.get("confirmation_url"),
metadata_json=payment_metadata,
payment_method_type=None,
yookassa_created_at=yookassa_created_at,
yookassa_created_at=yookassa_created_at,
test_mode=yookassa_response.get("test_mode", False)
)
@@ -156,7 +158,7 @@ class PaymentService:
captured_at = None
if status == "succeeded":
captured_at = datetime.utcnow()
captured_at = datetime.utcnow()
updated_payment = await update_yookassa_payment_status(
db,
@@ -184,22 +186,27 @@ class PaymentService:
db, yookassa_payment_id, transaction.id
)
from app.database.crud.user import add_user_balance
await add_user_balance(db, updated_payment.user_id, updated_payment.amount_kopeks)
if self.bot:
try:
await self.bot.send_message(
updated_payment.user.telegram_id,
f"✅ <b>Пополнение успешно!</b>\n\n"
f"💰 Сумма: {settings.format_price(updated_payment.amount_kopeks)}\n"
f" Способ: Банковская карта\n"
f"🆔 Транзакция: {yookassa_payment_id[:8]}...\n\n"
f"Баланс пополнен автоматически!",
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Ошибка отправки уведомления о пополнении: {e}")
user = await get_user_by_id(db, updated_payment.user_id)
if user:
await add_user_balance(db, user, updated_payment.amount_kopeks, f"Пополнение YooKassa: {updated_payment.amount_kopeks/100:.2f}")
if self.bot:
try:
await self.bot.send_message(
user.telegram_id,
f"✅ <b>Пополнение успешно!</b>\n\n"
f"💰 Сумма: {settings.format_price(updated_payment.amount_kopeks)}\n"
f"🏦 Способ: Банковская карта\n"
f"🆔 Транзакция: {yookassa_payment_id[:8]}...\n\n"
f"Баланс пополнен автоматически!",
parse_mode="HTML"
)
logger.info(f"Отправлено уведомление пользователю {user.telegram_id} о пополнении на {updated_payment.amount_kopeks/100:.2f}")
except Exception as e:
logger.error(f"Ошибка отправки уведомления о пополнении: {e}")
else:
logger.error(f"Пользователь с ID {updated_payment.user_id} не найден при пополнении баланса")
return False
return True
@@ -231,15 +238,17 @@ class PaymentService:
transaction_id=transaction.id
)
await add_user_balance(db, payment.user_id, payment.amount_kopeks)
user = await get_user_by_id(db, payment.user_id)
if user:
await add_user_balance(db, user, payment.amount_kopeks, f"Пополнение YooKassa: {payment.amount_kopeks/100:.2f}")
logger.info(f"Успешно обработан платеж YooKassa {payment.yookassa_payment_id}: "
f"пользователь {payment.user_id} получил {payment.amount_kopeks/100}")
if self.bot:
if self.bot and user:
try:
await self._send_payment_success_notification(
payment.user.telegram_id,
user.telegram_id,
payment.amount_kopeks
)
except Exception as e:
@@ -339,4 +348,4 @@ class PaymentService:
except Exception as e:
logger.error(f"Ошибка обработки платежа: {e}")
return False
return False
+21 -9
View File
@@ -20,7 +20,7 @@ async def process_referral_registration(
referrer = await get_user_by_id(db, referrer_id)
if not new_user or not referrer:
logger.error(f"Пользователи не найдены: {new_user_id}, {referrer_id}")
logger.error(f"Пользователи не найдены: new_user_id={new_user_id}, referrer_id={referrer_id}")
return False
if new_user.referred_by_id != referrer_id:
@@ -81,9 +81,13 @@ async def process_referral_purchase(
logger.info(f"🔍 Покупка реферала {user_id}: первая = {is_first_purchase}, сумма = {purchase_amount_kopeks/100}")
if is_first_purchase:
if is_first_purchase and settings.REFERRAL_REGISTRATION_REWARD > 0:
reward_amount = settings.REFERRAL_REGISTRATION_REWARD
if reward_amount > 1000000:
logger.error(f"❌ КРИТИЧЕСКАЯ ОШИБКА: reward_amount = {reward_amount} слишком большой! Проверьте настройки REFERRAL_REGISTRATION_REWARD")
reward_amount = 10000
await add_user_balance(
db, referrer, reward_amount,
f"Реферальная награда за первую покупку {user.full_name}"
@@ -100,12 +104,18 @@ async def process_referral_purchase(
logger.info(f"🎉 Первая покупка реферала: {referrer.telegram_id} получил {reward_amount/100}")
commission_amount = int(purchase_amount_kopeks * settings.REFERRAL_COMMISSION_PERCENT / 100)
if not (0 <= settings.REFERRAL_COMMISSION_PERCENT <= 100):
logger.error(f"❌ КРИТИЧЕСКАЯ ОШИБКА: REFERRAL_COMMISSION_PERCENT = {settings.REFERRAL_COMMISSION_PERCENT} некорректный! Должен быть от 0 до 100")
commission_percent = 10
else:
commission_percent = settings.REFERRAL_COMMISSION_PERCENT
commission_amount = int(purchase_amount_kopeks * commission_percent / 100)
if commission_amount > 0:
await add_user_balance(
db, referrer, commission_amount,
f"Комиссия {settings.REFERRAL_COMMISSION_PERCENT}% с покупки {user.full_name}"
f"Комиссия {commission_percent}% с покупки {user.full_name}"
)
await create_referral_earning(
@@ -128,6 +138,8 @@ async def process_referral_purchase(
except Exception as e:
logger.error(f"Ошибка обработки покупки реферала: {e}")
import traceback
logger.error(f"Полный traceback: {traceback.format_exc()}")
return False
@@ -141,7 +153,7 @@ async def get_referral_stats_for_user(db: AsyncSession, user_id: int) -> dict:
invited_count_result = await db.execute(
select(func.count(User.id)).where(User.referred_by_id == user_id)
)
invited_count = invited_count_result.scalar()
invited_count = invited_count_result.scalar() or 0
paid_referrals_result = await db.execute(
select(func.count(User.id)).where(
@@ -149,13 +161,13 @@ async def get_referral_stats_for_user(db: AsyncSession, user_id: int) -> dict:
User.has_had_paid_subscription == True
)
)
paid_referrals_count = paid_referrals_result.scalar()
paid_referrals_count = paid_referrals_result.scalar() or 0
total_earned = await get_referral_earnings_sum(db, user_id)
total_earned = await get_referral_earnings_sum(db, user_id) or 0
from datetime import datetime, timedelta
month_ago = datetime.utcnow() - timedelta(days=30)
month_earned = await get_referral_earnings_sum(db, user_id, start_date=month_ago)
month_earned = await get_referral_earnings_sum(db, user_id, start_date=month_ago) or 0
return {
"invited_count": invited_count,
@@ -171,4 +183,4 @@ async def get_referral_stats_for_user(db: AsyncSession, user_id: int) -> dict:
"paid_referrals_count": 0,
"total_earned_kopeks": 0,
"month_earned_kopeks": 0
}
}
+187 -198
View File
@@ -2,6 +2,8 @@ import logging
from typing import Dict, List, Any, Optional
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import delete
from datetime import datetime, timedelta
import re
from app.config import settings
from app.external.remnawave_api import (
@@ -25,6 +27,34 @@ class RemnaWaveService:
base_url=settings.REMNAWAVE_API_URL,
api_key=settings.REMNAWAVE_API_KEY
)
def _parse_remnawave_date(self, date_str: str) -> datetime:
if not date_str:
return datetime.utcnow() + timedelta(days=30)
try:
cleaned_date = date_str.strip()
if cleaned_date.endswith('Z'):
cleaned_date = cleaned_date[:-1] + '+00:00'
if '+00:00+00:00' in cleaned_date:
cleaned_date = cleaned_date.replace('+00:00+00:00', '+00:00')
cleaned_date = re.sub(r'(\+\d{2}:\d{2})\+\d{2}:\d{2}$', r'\1', cleaned_date)
parsed_date = datetime.fromisoformat(cleaned_date)
if parsed_date.tzinfo is not None:
parsed_date = parsed_date.replace(tzinfo=None)
logger.debug(f"Успешно распарсена дата: {date_str} -> {parsed_date}")
return parsed_date
except Exception as e:
logger.warning(f"⚠️ Не удалось распарсить дату '{date_str}': {e}. Используем дефолтную дату.")
return datetime.utcnow() + timedelta(days=30)
async def get_system_statistics(self) -> Dict[str, Any]:
try:
@@ -58,7 +88,6 @@ class RemnaWaveService:
logger.error(f"Ошибка получения статистики нод: {e}")
nodes_stats = {}
from datetime import datetime
total_download = sum(node.get('downloadBytes', 0) for node in realtime_usage)
total_upload = sum(node.get('uploadBytes', 0) for node in realtime_usage)
@@ -199,7 +228,6 @@ class RemnaWaveService:
'GBPS': 1024 ** 3
}
import re
match = re.match(r'([0-9.,]+)([A-Z]+)', bandwidth_str)
if match:
value_str = match.group(1).replace(',', '.')
@@ -397,212 +425,185 @@ class RemnaWaveService:
async def sync_users_from_panel(self, db: AsyncSession, sync_type: str = "all") -> Dict[str, int]:
try:
stats = {"created": 0, "updated": 0, "errors": 0, "deleted": 0}
logger.info(f"🔄 Начинаем синхронизацию типа: {sync_type}")
async with self.api as api:
panel_users_data = await api._make_request('GET', '/api/users')
panel_users = panel_users_data['response']['users']
logger.info(f"👥 Найдено пользователей в панели: {len(panel_users)}")
logger.info(f"🔄 Начинаем синхронизацию типа: {sync_type}")
async with self.api as api:
panel_users = []
start = 0
size = 100
bot_users = await get_users_list(db, offset=0, limit=10000)
bot_users_by_telegram_id = {user.telegram_id: user for user in bot_users}
panel_telegram_ids = set()
for i, panel_user in enumerate(panel_users):
try:
telegram_id = panel_user.get('telegramId')
if not telegram_id:
logger.debug(f"➡️ Пропускаем пользователя без telegram_id")
continue
panel_telegram_ids.add(telegram_id)
logger.info(f"🔄 Обрабатываем пользователя {i+1}/{len(panel_users)}: {telegram_id}")
while True:
logger.info(f"📥 Загружаем пользователей: start={start}, size={size}")
db_user = bot_users_by_telegram_id.get(telegram_id)
response = await api.get_all_users(start=start, size=size)
users_batch = response['users']
total_users = response['total']
if not db_user:
if sync_type in ["new_only", "all"]:
logger.info(f"🆕 Создание пользователя для telegram_id {telegram_id}")
from app.database.crud.user import create_user
logger.info(f"📊 Получено {len(users_batch)} пользователей из {total_users}")
for user_obj in users_batch:
user_dict = {
'uuid': user_obj.uuid,
'shortUuid': user_obj.short_uuid,
'username': user_obj.username,
'status': user_obj.status.value,
'telegramId': user_obj.telegram_id,
'expireAt': user_obj.expire_at.isoformat() + 'Z',
'trafficLimitBytes': user_obj.traffic_limit_bytes,
'usedTrafficBytes': user_obj.used_traffic_bytes,
'hwidDeviceLimit': user_obj.hwid_device_limit,
'subscriptionUrl': user_obj.subscription_url,
'activeInternalSquads': user_obj.active_internal_squads
}
panel_users.append(user_dict)
if len(users_batch) < size:
break
start += size
if start > total_users:
break
logger.info(f"✅ Всего загружено пользователей из панели: {len(panel_users)}")
bot_users = await get_users_list(db, offset=0, limit=10000)
bot_users_by_telegram_id = {user.telegram_id: user for user in bot_users}
logger.info(f"📊 Пользователей в боте: {len(bot_users)}")
panel_users_with_tg = [
user for user in panel_users
if user.get('telegramId') is not None
]
logger.info(f"📊 Пользователей в панели с Telegram ID: {len(panel_users_with_tg)}")
panel_telegram_ids = set()
for i, panel_user in enumerate(panel_users_with_tg):
try:
telegram_id = panel_user.get('telegramId')
if not telegram_id:
continue
panel_telegram_ids.add(telegram_id)
if (i + 1) % 10 == 0:
logger.info(f"🔄 Обрабатываем пользователя {i+1}/{len(panel_users_with_tg)}: {telegram_id}")
db_user = bot_users_by_telegram_id.get(telegram_id)
if not db_user:
if sync_type in ["new_only", "all"]:
logger.info(f"🆕 Создание пользователя для telegram_id {telegram_id}")
db_user = await create_user(
db=db,
telegram_id=telegram_id,
username=panel_user.get('username') or f"user_{telegram_id}",
first_name=f"Panel User {telegram_id}",
language="ru"
)
from app.database.crud.user import create_user
db_user = await create_user(
db=db,
telegram_id=telegram_id,
username=panel_user.get('username') or f"user_{telegram_id}",
first_name=f"Panel User {telegram_id}",
language="ru"
)
await update_user(db, db_user, remnawave_uuid=panel_user.get('uuid'))
await self._create_subscription_from_panel_data(db, db_user, panel_user)
stats["created"] += 1
logger.info(f"✅ Создан пользователь {telegram_id} с подпиской")
else:
if sync_type in ["update_only", "all"]:
logger.debug(f"🔄 Обновление пользователя {telegram_id}")
if not db_user.remnawave_uuid:
await update_user(db, db_user, remnawave_uuid=panel_user.get('uuid'))
await self._create_subscription_from_panel_data(db, db_user, panel_user)
await self._update_subscription_from_panel_data(db, db_user, panel_user)
stats["created"] += 1
logger.info(f"Создан пользователь {telegram_id} с подпиской")
stats["updated"] += 1
logger.debug(f"Обновлён пользователь {telegram_id}")
else:
if sync_type in ["update_only", "all"]:
logger.debug(f"🔄 Обновление пользователя {telegram_id}")
if not db_user.remnawave_uuid:
await update_user(db, db_user, remnawave_uuid=panel_user.get('uuid'))
await self._update_subscription_from_panel_data(db, db_user, panel_user)
stats["updated"] += 1
logger.debug(f"✅ Обновлён пользователь {telegram_id}")
except Exception as user_error:
logger.error(f"❌ Ошибка обработки пользователя {telegram_id}: {user_error}")
stats["errors"] += 1
continue
except Exception as user_error:
logger.error(f"❌ Ошибка обработки пользователя {telegram_id}: {user_error}")
stats["errors"] += 1
continue
if sync_type == "all":
logger.info("🗑️ Деактивация подписок пользователей, отсутствующих в панели...")
if sync_type == "all":
logger.info("🗑️ ПОЛНАЯ очистка подписок пользователей, отсутствующих в панели...")
for telegram_id, db_user in bot_users_by_telegram_id.items():
if telegram_id not in panel_telegram_ids and db_user.subscription:
try:
logger.info(f"🗑️ ПОЛНАЯ очистка данных подписки пользователя {telegram_id} (нет в панели)")
subscription = db_user.subscription
if db_user.remnawave_uuid:
try:
for telegram_id, db_user in bot_users_by_telegram_id.items():
if telegram_id not in panel_telegram_ids and db_user.subscription:
try:
logger.info(f"🗑️ Деактивация подписки пользователя {telegram_id} (нет в панели)")
subscription = db_user.subscription
if db_user.remnawave_uuid:
try:
async with self.api as api:
devices_reset = await api.reset_user_devices(db_user.remnawave_uuid)
if devices_reset:
logger.info(f"🔧 Сброшены HWID устройства для пользователя {telegram_id}")
else:
logger.warning(f"⚠️ Не удалось сбросить HWID устройства для пользователя {telegram_id}")
except Exception as hwid_error:
logger.error(f"❌ Ошибка сброса HWID устройств для {telegram_id}: {hwid_error}")
except Exception as hwid_error:
logger.error(f"❌ Ошибка сброса HWID устройств для {telegram_id}: {hwid_error}")
try:
from sqlalchemy import delete
from app.database.models import SubscriptionServer
try:
from app.database.crud.subscription import get_subscription_server_ids, remove_subscription_servers
from sqlalchemy import delete
from app.database.models import SubscriptionServer
await db.execute(
delete(SubscriptionServer).where(
SubscriptionServer.subscription_id == subscription.id
)
await db.execute(
delete(SubscriptionServer).where(
SubscriptionServer.subscription_id == subscription.id
)
logger.info(f"🗑️ УДАЛЕНЫ ВСЕ записи SubscriptionServer для подписки {subscription.id}")
except Exception as servers_error:
logger.warning(f"⚠️ Не удалось удалить серверы подписки: {servers_error}")
try:
from sqlalchemy import delete
from app.database.models import Transaction
await db.execute(
delete(Transaction).where(Transaction.user_id == db_user.id)
)
logger.info(f"🗑️ УДАЛЕНЫ ВСЕ транзакции пользователя {telegram_id}")
except Exception as transactions_error:
logger.warning(f"⚠️ Не удалось удалить транзакции: {transactions_error}")
try:
from sqlalchemy import delete
from app.database.models import ReferralEarning
await db.execute(
delete(ReferralEarning).where(ReferralEarning.user_id == db_user.id)
)
await db.execute(
delete(ReferralEarning).where(ReferralEarning.referral_id == db_user.id)
)
logger.info(f"🗑️ УДАЛЕНЫ ВСЕ реферальные доходы пользователя {telegram_id}")
except Exception as referral_error:
logger.warning(f"⚠️ Не удалось удалить реферальные доходы: {referral_error}")
try:
from sqlalchemy import delete
from app.database.models import PromoCodeUse
await db.execute(
delete(PromoCodeUse).where(PromoCodeUse.user_id == db_user.id)
)
logger.info(f"🗑️ УДАЛЕНЫ ВСЕ использования промокодов пользователя {telegram_id}")
except Exception as promo_error:
logger.warning(f"⚠️ Не удалось удалить использования промокодов: {promo_error}")
try:
db_user.balance_kopeks = 0
logger.info(f"💰 Сброшен баланс пользователя {telegram_id}")
except Exception as balance_error:
logger.warning(f"⚠️ Не удалось сбросить баланс: {balance_error}")
from app.database.models import SubscriptionStatus
from datetime import datetime
subscription.status = SubscriptionStatus.DISABLED.value
subscription.is_trial = True
subscription.end_date = datetime.utcnow()
subscription.traffic_limit_gb = 0
subscription.traffic_used_gb = 0.0
subscription.device_limit = 1
subscription.connected_squads = []
subscription.autopay_enabled = False
subscription.autopay_days_before = 3
subscription.remnawave_short_uuid = None
subscription.subscription_url = ""
db_user.remnawave_uuid = None
db_user.has_had_paid_subscription = False
db_user.used_promocodes = 0
await db.commit()
stats["deleted"] += 1
logger.info(f"✅ ПОЛНОСТЬЮ очищены ВСЕ данные подписки пользователя {telegram_id}")
except Exception as delete_error:
logger.error(f"❌ Ошибка полной очистки данных подписки {telegram_id}: {delete_error}")
stats["errors"] += 1
await db.rollback()
logger.info(f"🎯 Синхронизация завершена: создано {stats['created']}, обновлено {stats['updated']}, удалено {stats['deleted']}, ошибок {stats['errors']}")
)
logger.info(f"🗑️ Удалены серверы подписки для {telegram_id}")
except Exception as servers_error:
logger.warning(f"⚠️ Не удалось удалить серверы подписки: {servers_error}")
from app.database.models import SubscriptionStatus
subscription.status = SubscriptionStatus.DISABLED.value
subscription.is_trial = True
subscription.end_date = datetime.utcnow()
subscription.traffic_limit_gb = 0
subscription.traffic_used_gb = 0.0
subscription.device_limit = 1
subscription.connected_squads = []
subscription.autopay_enabled = False
subscription.remnawave_short_uuid = None
subscription.subscription_url = ""
db_user.remnawave_uuid = None
await db.commit()
stats["deleted"] += 1
logger.info(f"✅ Деактивирована подписка пользователя {telegram_id} (сохранен баланс)")
except Exception as delete_error:
logger.error(f"❌ Ошибка деактивации подписки {telegram_id}: {delete_error}")
stats["errors"] += 1
await db.rollback()
logger.info(f"🎯 Синхронизация завершена: создано {stats['created']}, обновлено {stats['updated']}, деактивировано {stats['deleted']}, ошибок {stats['errors']}")
return stats
except Exception as e:
logger.error(f"❌ Критическая ошибка синхронизации пользователей: {e}")
return {"created": 0, "updated": 0, "errors": 1, "deleted": 0}
async def _create_subscription_from_panel_data(self, db: AsyncSession, user, panel_user):
try:
from app.database.crud.subscription import create_subscription
from app.database.models import SubscriptionStatus
from datetime import datetime, timedelta
import pytz
expire_at_str = panel_user.get('expireAt', '')
try:
if expire_at_str:
if expire_at_str.endswith('Z'):
expire_at_str = expire_at_str[:-1] + '+00:00'
expire_at = datetime.fromisoformat(expire_at_str)
if expire_at.tzinfo is not None:
expire_at = expire_at.replace(tzinfo=None)
else:
expire_at = datetime.utcnow() + timedelta(days=30)
except Exception as date_error:
logger.warning(f"⚠️ Ошибка парсинга даты {expire_at_str}: {date_error}")
expire_at = datetime.utcnow() + timedelta(days=30)
expire_at = self._parse_remnawave_date(expire_at_str)
panel_status = panel_user.get('status', 'ACTIVE')
current_time = datetime.utcnow()
@@ -672,7 +673,6 @@ class RemnaWaveService:
try:
from app.database.crud.subscription import get_subscription_by_user_id
from app.database.models import SubscriptionStatus
from datetime import datetime, timedelta
subscription = await get_subscription_by_user_id(db, user.id)
@@ -683,19 +683,12 @@ class RemnaWaveService:
panel_status = panel_user.get('status', 'ACTIVE')
expire_at_str = panel_user.get('expireAt', '')
try:
if expire_at_str:
if expire_at_str.endswith('Z'):
expire_at_str = expire_at_str[:-1] + '+00:00'
expire_at = datetime.fromisoformat(expire_at_str)
if expire_at.tzinfo is not None:
expire_at = expire_at.replace(tzinfo=None)
if abs((subscription.end_date - expire_at).total_seconds()) > 60:
subscription.end_date = expire_at
logger.debug(f"Обновлена дата окончания подписки до {expire_at}")
except Exception as date_error:
logger.warning(f"⚠️ Ошибка парсинга даты при обновлении {expire_at_str}: {date_error}")
if expire_at_str:
expire_at = self._parse_remnawave_date(expire_at_str)
if abs((subscription.end_date - expire_at).total_seconds()) > 60:
subscription.end_date = expire_at
logger.debug(f"Обновлена дата окончания подписки до {expire_at}")
current_time = datetime.utcnow()
if panel_status == 'ACTIVE' and subscription.end_date > current_time:
@@ -1001,7 +994,6 @@ class RemnaWaveService:
node_realtime = stats
break
from datetime import datetime, timedelta
end_date = datetime.now()
start_date = end_date - timedelta(days=7)
@@ -1089,7 +1081,6 @@ class RemnaWaveService:
logger.error(f"❌ Ошибка удаления связанных записей: {records_error}")
try:
from datetime import datetime
user.balance_kopeks = 0
user.remnawave_uuid = None
@@ -1213,7 +1204,6 @@ class RemnaWaveService:
from app.database.crud.subscription import get_all_subscriptions
from app.database.models import SubscriptionStatus
from datetime import datetime
page = 1
limit = 100
@@ -1266,7 +1256,6 @@ class RemnaWaveService:
from app.database.crud.subscription import get_all_subscriptions
from app.database.models import SubscriptionStatus
from datetime import datetime
page = 1
limit = 100
+80 -15
View File
@@ -48,7 +48,6 @@ class TributeService:
if not payment_url:
return None
return payment_url
except Exception as e:
@@ -91,19 +90,28 @@ class TributeService:
return {"status": "ok", "event": event_type}
async def _handle_successful_payment(self, payment_data: Dict[str, Any]):
"""Обработка успешного платежа - ПОЛНОСТЬЮ ПЕРЕРАБОТАННАЯ ВЕРСИЯ"""
try:
user_id = payment_data["user_id"]
amount_kopeks = payment_data["amount_kopeks"]
payment_id = payment_data["payment_id"]
logger.info(f"Обрабатываем успешный Tribute платеж: user_id={user_id}, amount={amount_kopeks}, payment_id={payment_id}")
async for session in get_db():
existing_transaction = await get_transaction_by_external_id(
session, f"donation_{payment_id}", PaymentMethod.TRIBUTE
from app.database.crud.transaction import check_tribute_payment_duplicate, create_unique_tribute_transaction
duplicate_transaction = await check_tribute_payment_duplicate(
session, payment_id, amount_kopeks, user_id
)
if existing_transaction:
logger.warning(f"Транзакция с donation_request_id {payment_id} уже существует")
if duplicate_transaction:
logger.warning(f"Найден дубликат платежа:")
logger.warning(f" Transaction ID: {duplicate_transaction.id}")
logger.warning(f" Amount: {duplicate_transaction.amount_kopeks} коп")
logger.warning(f" Created: {duplicate_transaction.created_at}")
logger.warning(f"Платеж игнорирован")
return
user = await get_user_by_telegram_id(session, user_id)
@@ -111,29 +119,34 @@ class TributeService:
logger.error(f"Пользователь {user_id} не найден")
return
transaction = await create_transaction(
logger.info(f"Найден пользователь {user.telegram_id}, текущий баланс: {user.balance_kopeks} коп")
transaction = await create_unique_tribute_transaction(
db=session,
user_id=user.id,
type=TransactionType.DEPOSIT,
user_id=user.id,
payment_id=payment_id,
amount_kopeks=amount_kopeks,
description=f"Пополнение через Tribute: {amount_kopeks/100}",
payment_method=PaymentMethod.TRIBUTE,
external_id=f"donation_{payment_id}",
is_completed=True
description=f"Пополнение через Tribute: {amount_kopeks/100} (ID: {payment_id})"
)
old_balance = user.balance_kopeks
user.balance_kopeks += amount_kopeks
user.updated_at = datetime.utcnow()
await session.commit()
logger.info(f"Баланс пользователя {user_id} обновлен: {old_balance} -> {user.balance_kopeks} коп (+{amount_kopeks})")
await self._send_success_notification(user_id, amount_kopeks)
logger.info(f"Успешно обработан Tribute платеж: {amount_kopeks/100}₽ для пользователя {user_id}")
break
except Exception as e:
logger.error(f"Ошибка обработки успешного Tribute платежа: {e}")
logger.error(f"Ошибка обработки успешного Tribute платежа: {e}", exc_info=True)
async def _handle_failed_payment(self, payment_data: Dict[str, Any]):
"""Обработка неудачного платежа"""
try:
user_id = payment_data["user_id"]
@@ -141,7 +154,7 @@ class TributeService:
async for session in get_db():
transaction = await get_transaction_by_external_id(
session, payment_id, PaymentMethod.TRIBUTE
session, f"donation_{payment_id}", PaymentMethod.TRIBUTE
)
if transaction:
@@ -157,6 +170,7 @@ class TributeService:
logger.error(f"Ошибка обработки неудачного Tribute платежа: {e}")
async def _handle_refund(self, refund_data: Dict[str, Any]):
"""Обработка возврата"""
try:
user_id = refund_data["user_id"]
@@ -189,6 +203,7 @@ class TributeService:
logger.error(f"Ошибка обработки возврата Tribute: {e}")
async def _send_success_notification(self, user_id: int, amount_kopeks: int):
"""Отправка уведомления об успешном платеже"""
try:
amount_rubles = amount_kopeks / 100
@@ -217,6 +232,7 @@ class TributeService:
logger.error(f"Ошибка отправки уведомления об успешном платеже: {e}")
async def _send_failure_notification(self, user_id: int):
"""Отправка уведомления о неудачном платеже"""
try:
text = (
@@ -245,6 +261,7 @@ class TributeService:
logger.error(f"Ошибка отправки уведомления о неудачном платеже: {e}")
async def _send_refund_notification(self, user_id: int, amount_kopeks: int):
"""Отправка уведомления о возврате"""
try:
amount_rubles = amount_kopeks / 100
@@ -272,6 +289,54 @@ class TributeService:
except Exception as e:
logger.error(f"Ошибка отправки уведомления о возврате: {e}")
async def force_process_payment(
self,
payment_id: str,
user_id: int,
amount_kopeks: int,
description: str = "Принудительная обработка Tribute платежа"
) -> bool:
"""Принудительная обработка платежа (для отладки)"""
try:
logger.info(f"🔧 ПРИНУДИТЕЛЬНАЯ ОБРАБОТКА: payment_id={payment_id}, user_id={user_id}, amount={amount_kopeks}")
async for session in get_db():
user = await get_user_by_telegram_id(session, user_id)
if not user:
logger.error(f"❌ Пользователь {user_id} не найден")
return False
external_id = f"force_donation_{payment_id}_{int(datetime.utcnow().timestamp())}"
transaction = await create_transaction(
db=session,
user_id=user.id,
type=TransactionType.DEPOSIT,
amount_kopeks=amount_kopeks,
description=description,
payment_method=PaymentMethod.TRIBUTE,
external_id=external_id,
is_completed=True
)
old_balance = user.balance_kopeks
user.balance_kopeks += amount_kopeks
user.updated_at = datetime.utcnow()
await session.commit()
logger.info(f"💰 ПРИНУДИТЕЛЬНО обновлен баланс: {old_balance} -> {user.balance_kopeks} коп")
await self._send_success_notification(user_id, amount_kopeks)
logger.info(f"✅ Принудительно обработан платеж {payment_id}")
return True
except Exception as e:
logger.error(f"❌ Ошибка принудительной обработки: {e}", exc_info=True)
return False
async def get_payment_status(self, payment_id: str) -> Optional[Dict[str, Any]]:
return await self.tribute_api.get_payment_status(payment_id)
@@ -281,4 +346,4 @@ class TributeService:
amount_kopeks: Optional[int] = None,
reason: str = "Возврат по запросу"
) -> Optional[Dict[str, Any]]:
return await self.tribute_api.refund_payment(payment_id, amount_kopeks, reason)
return await self.tribute_api.refund_payment(payment_id, amount_kopeks, reason)
+2 -3
View File
@@ -135,16 +135,15 @@ class UserService:
return False
if amount_kopeks > 0:
await add_user_balance(db, user, amount_kopeks, description)
await add_user_balance(db, user, amount_kopeks, description=description)
logger.info(f"Админ {admin_id} пополнил баланс пользователя {user_id} на {amount_kopeks/100}")
return True
else:
success = await subtract_user_balance(db, user, abs(amount_kopeks), description)
if success:
logger.info(f"Админ {admin_id} списал с баланса пользователя {user_id} {abs(amount_kopeks)/100}")
return success
return True
except Exception as e:
logger.error(f"Ошибка изменения баланса пользователя: {e}")
return False
+123 -12
View File
@@ -2,6 +2,7 @@ import asyncio
import logging
import sys
import os
import signal
from pathlib import Path
sys.path.append(str(Path(__file__).parent))
@@ -10,10 +11,23 @@ from app.bot import setup_bot
from app.config import settings
from app.database.database import init_db
from app.services.monitoring_service import monitoring_service
from app.services.maintenance_service import maintenance_service
from app.services.payment_service import PaymentService
from app.external.webhook_server import WebhookServer
from app.external.yookassa_webhook import start_yookassa_webhook_server
from app.database.universal_migration import run_universal_migration
class GracefulExit:
def __init__(self):
self.exit = False
def exit_gracefully(self, signum, frame):
logging.getLogger(__name__).info(f"Получен сигнал {signum}. Корректное завершение работы...")
self.exit = True
async def main():
logging.basicConfig(
level=getattr(logging, settings.LOG_LEVEL),
@@ -27,7 +41,15 @@ async def main():
logger = logging.getLogger(__name__)
logger.info("🚀 Запуск Bedolaga Remnawave Bot...")
killer = GracefulExit()
signal.signal(signal.SIGINT, killer.exit_gracefully)
signal.signal(signal.SIGTERM, killer.exit_gracefully)
webhook_server = None
yookassa_server_task = None
monitoring_task = None
maintenance_task = None
polling_task = None
try:
logger.info("📊 Инициализация базы данных...")
@@ -56,38 +78,127 @@ async def main():
monitoring_service.bot = bot
payment_service = PaymentService(bot)
if settings.TRIBUTE_ENABLED:
logger.info("🌐 Запуск webhook сервера для Tribute...")
logger.info("🌐 Запуск Tribute webhook сервера...")
webhook_server = WebhookServer(bot)
await webhook_server.start()
else:
logger.info("ℹ️ Tribute отключен, webhook сервер не запускается")
if settings.is_yookassa_enabled():
logger.info("💳 Запуск YooKassa webhook сервера...")
yookassa_server_task = asyncio.create_task(
start_yookassa_webhook_server(payment_service)
)
else:
logger.info("ℹ️ YooKassa отключена, webhook сервер не запускается")
logger.info("🔍 Запуск службы мониторинга...")
monitoring_task = asyncio.create_task(monitoring_service.start_monitoring())
logger.info("🔧 Запуск службы техработ...")
maintenance_task = asyncio.create_task(maintenance_service.start_monitoring())
logger.info("🔄 Запуск polling...")
polling_task = asyncio.create_task(dp.start_polling(bot, skip_updates=True))
logger.info("=" * 50)
logger.info("🎯 Активные webhook endpoints:")
if settings.TRIBUTE_ENABLED:
logger.info(f" Tribute: {settings.WEBHOOK_URL}:{settings.TRIBUTE_WEBHOOK_PORT}{settings.TRIBUTE_WEBHOOK_PATH}")
if settings.is_yookassa_enabled():
logger.info(f" YooKassa: {settings.WEBHOOK_URL}:{settings.YOOKASSA_WEBHOOK_PORT}{settings.YOOKASSA_WEBHOOK_PATH}")
logger.info("=" * 50)
try:
await asyncio.gather(
dp.start_polling(bot),
monitoring_task
)
while not killer.exit:
await asyncio.sleep(1)
if yookassa_server_task and yookassa_server_task.done():
exception = yookassa_server_task.exception()
if exception:
logger.error(f"YooKassa webhook сервер завершился с ошибкой: {exception}")
logger.info("🔄 Перезапуск YooKassa webhook сервера...")
yookassa_server_task = asyncio.create_task(
start_yookassa_webhook_server(payment_service)
)
if monitoring_task.done():
exception = monitoring_task.exception()
if exception:
logger.error(f"Служба мониторинга завершилась с ошибкой: {exception}")
monitoring_task = asyncio.create_task(monitoring_service.start_monitoring())
if maintenance_task.done():
exception = maintenance_task.exception()
if exception:
logger.error(f"Служба техработ завершилась с ошибкой: {exception}")
maintenance_task = asyncio.create_task(maintenance_service.start_monitoring())
if polling_task.done():
exception = polling_task.exception()
if exception:
logger.error(f"Polling завершился с ошибкой: {exception}")
break
except Exception as e:
logger.error(f"Ошибка в основном цикле: {e}")
monitoring_service.stop_monitoring()
if webhook_server:
await webhook_server.stop()
raise
except Exception as e:
logger.error(f"❌ Критическая ошибка при запуске: {e}")
raise
finally:
logger.info("🛑 Завершение работы бота")
monitoring_service.stop_monitoring()
logger.info("🛑 Начинается корректное завершение работы...")
if yookassa_server_task and not yookassa_server_task.done():
logger.info("⏹️ Остановка YooKassa webhook сервера...")
yookassa_server_task.cancel()
try:
await yookassa_server_task
except asyncio.CancelledError:
pass
if monitoring_task and not monitoring_task.done():
logger.info("⏹️ Остановка службы мониторинга...")
monitoring_service.stop_monitoring()
monitoring_task.cancel()
try:
await monitoring_task
except asyncio.CancelledError:
pass
if maintenance_task and not maintenance_task.done():
logger.info("⏹️ Остановка службы техработ...")
await maintenance_service.stop_monitoring()
maintenance_task.cancel()
try:
await maintenance_task
except asyncio.CancelledError:
pass
if polling_task and not polling_task.done():
logger.info("⏹️ Остановка polling...")
polling_task.cancel()
try:
await polling_task
except asyncio.CancelledError:
pass
if webhook_server:
logger.info("⏹️ Остановка Tribute webhook сервера...")
await webhook_server.stop()
if 'bot' in locals():
try:
await bot.session.close()
logger.info("✅ Сессия бота закрыта")
except Exception as e:
logger.error(f"Ошибка закрытия сессии бота: {e}")
logger.info("✅ Завершение работы бота завершено")
if __name__ == "__main__":
+3
View File
@@ -18,6 +18,9 @@ yookassa==3.0.0
# Логирование и мониторинг
structlog==23.2.0
# Планировщик задач для техработ
APScheduler==3.10.4
# Утилиты
python-dateutil==2.8.2
pytz==2023.4