Compare commits
66 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 92ff4f661f | |||
| dfebc5e608 | |||
| 8e6121cd0b | |||
| 6c0d5714bd | |||
| ed29697434 | |||
| 07aebef82f | |||
| 7ce390fd50 | |||
| cebe392d1c | |||
| fd160f585a | |||
| 5c0516fbc0 | |||
| ba075326c1 | |||
| c32425341e | |||
| cc8df1eff3 | |||
| 49cda0791c | |||
| 60ea363a62 | |||
| f66619b225 | |||
| 5999402e1a | |||
| 3aebccc090 | |||
| ed577dd51a | |||
| e517487596 | |||
| f9441909a4 | |||
| f13a015bfa | |||
| 18bf946464 | |||
| e4c46b2488 | |||
| e0539b116c | |||
| acdf7af1d9 | |||
| d5b7828079 | |||
| 03be071edc | |||
| e8742f30a6 | |||
| 3d7727dbf2 | |||
| 5bff67860f | |||
| 429e967da2 | |||
| c3df3beeae | |||
| 4d664c15d7 | |||
| 53be365e38 | |||
| 3c3ffed32e | |||
| 015726b3a6 | |||
| 4a3267712b | |||
| 9b45aae410 | |||
| f43b391d0e | |||
| b09b7a0c84 | |||
| 541b44b933 | |||
| b21fce0eff | |||
| 017aa27c0f | |||
| 496978e8a2 | |||
| b3711a590d | |||
| 631a0f2c4c | |||
| 503f7ca455 | |||
| 80c771b7bc | |||
| 6e344d350b | |||
| 7738bc83af | |||
| 6f7a1eecf6 | |||
| bf147a736a | |||
| ea3decf1ac | |||
| b1c0315a6e | |||
| f6fa7e8f41 | |||
| 33d4d7994c | |||
| 0ce4ce2b08 | |||
| 3a9163e919 | |||
| 376cff748b | |||
| d5cef080da | |||
| 33e31213b3 | |||
| 8077ba0c35 | |||
| 67ad3847f8 | |||
| 1eda042b34 | |||
| 961ff46f43 |
@@ -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
@@ -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
@@ -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"]
|
||||
|
||||
@@ -44,6 +44,7 @@
|
||||
- 🎁 **Промо-система** - коды на деньги, дни подписки, триал-периоды
|
||||
- 3 режима показа ссылки подписки: 1) С гайдом по подключению прямо в боте(тянущий данные приложений и ссылок на скачку из app-config.json) 2) Обычное открытие ссылки подписки в миниапе 3) Интеграция сабпейджа maposia - кастомно прописать ссылку можно
|
||||
- Возможность переключаться между пакетной продажей трафика и фиксированной(Пропуская шаг выбора пакета трафика при оформлении/настройки подписки юзера)
|
||||
- Возможность задать доступные дни для покупки первой подписки и при продлении
|
||||
|
||||
### 💪 **Enterprise готовность**
|
||||
- 🏗️ **Современная архитектура** - AsyncIO, PostgreSQL, Redis
|
||||
@@ -62,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
|
||||
@@ -121,6 +123,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
|
||||
|
||||
@@ -422,7 +430,6 @@ bedolaga_bot/
|
||||
```
|
||||
project/
|
||||
├── docker-compose.yml # 🚀 Продакшн
|
||||
├── docker-compose.local.yml # 🏠 Разработка
|
||||
├── .env # ⚙️ Конфиг
|
||||
└── .env.example # 📝 Пример
|
||||
```
|
||||
@@ -433,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:
|
||||
@@ -502,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>
|
||||
@@ -614,8 +527,6 @@ networks:
|
||||
<summary>📄 Показать dev конфигурацию</summary>
|
||||
|
||||
```yaml
|
||||
version: '3.8'
|
||||
|
||||
services:
|
||||
# 🗄️ PostgreSQL для разработки
|
||||
postgres-dev:
|
||||
|
||||
+41
-1
@@ -33,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
|
||||
@@ -216,6 +218,44 @@ class Settings(BaseSettings):
|
||||
|
||||
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",
|
||||
|
||||
Vendored
+75
-21
@@ -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_kopeks = 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_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)}")
|
||||
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
|
||||
|
||||
Vendored
+69
-13
@@ -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
|
||||
})
|
||||
|
||||
Vendored
+150
-59
@@ -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
|
||||
|
||||
+106
-56
@@ -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(
|
||||
|
||||
@@ -500,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:
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -599,3 +599,11 @@ def get_maintenance_keyboard(language: str = "ru", is_active: bool = False, moni
|
||||
]
|
||||
|
||||
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
@@ -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)
|
||||
|
||||
@@ -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__)
|
||||
@@ -186,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
|
||||
|
||||
@@ -233,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:
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -12,7 +12,9 @@ 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
|
||||
|
||||
|
||||
@@ -44,6 +46,7 @@ async def main():
|
||||
signal.signal(signal.SIGTERM, killer.exit_gracefully)
|
||||
|
||||
webhook_server = None
|
||||
yookassa_server_task = None
|
||||
monitoring_task = None
|
||||
maintenance_task = None
|
||||
polling_task = None
|
||||
@@ -75,13 +78,23 @@ 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())
|
||||
|
||||
@@ -91,10 +104,27 @@ async def main():
|
||||
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:
|
||||
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:
|
||||
@@ -123,6 +153,14 @@ async def main():
|
||||
finally:
|
||||
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()
|
||||
@@ -150,7 +188,7 @@ async def main():
|
||||
pass
|
||||
|
||||
if webhook_server:
|
||||
logger.info("⏹️ Остановка webhook сервера...")
|
||||
logger.info("⏹️ Остановка Tribute webhook сервера...")
|
||||
await webhook_server.stop()
|
||||
|
||||
if 'bot' in locals():
|
||||
|
||||
Reference in New Issue
Block a user