diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..290e278 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,10 @@ +.git +.idea +.venv +.env +extra +docs +**/__pycache__ +src/db.sqlite3 +src/design/static_root +docker-services-data diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..dad96cd --- /dev/null +++ b/Dockerfile @@ -0,0 +1,44 @@ +# Сайт (он же webhook ботов) и воркер dramatiq: сервисы web и worker профиля app в docker-compose.yaml +ARG PYTHON_IMAGE=docker.local/python:3.12.2-slim +FROM ${PYTHON_IMAGE} + +# pg_dump и pg_restore для manage.py backup/restore - той же major-версии, что и сервис postgres +ARG PG_MAJOR=16 +# uid владельца смонтированных папок extra/ на хосте +ARG UID=1000 + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 \ + DJANGO_SETTINGS_MODULE=MyPointVPN.settings.docker + +# openssh-client - serverus ходит на серверы по ssh; postgresql-client - из репозитория PGDG, в Debian другая версия +RUN apt-get update \ + && apt-get install -y --no-install-recommends ca-certificates curl openssh-client \ + && install -d /usr/share/postgresql-common/pgdg \ + && curl -fsSL https://www.postgresql.org/media/keys/ACCC4CF8.asc -o /usr/share/postgresql-common/pgdg/apt.postgresql.org.asc \ + && . /etc/os-release \ + && echo "deb [signed-by=/usr/share/postgresql-common/pgdg/apt.postgresql.org.asc] https://apt.postgresql.org/pub/repos/apt ${VERSION_CODENAME}-pgdg main" > /etc/apt/sources.list.d/pgdg.list \ + && apt-get update \ + && apt-get install -y --no-install-recommends postgresql-client-${PG_MAJOR} \ + && apt-get purge -y curl \ + && apt-get autoremove -y \ + && rm -rf /var/lib/apt/lists/* + +RUN useradd --create-home --uid ${UID} app + +WORKDIR /app +COPY requirements.txt . +# wheels/ - то, что собрал make install: без индекса, если все пакеты там есть +COPY wheels/ wheels/ +RUN pip install --no-cache-dir --find-links=wheels -r requirements.txt && rm -rf wheels + +COPY src/ src/ +WORKDIR /app/src +# SECRET_KEY нужен только чтобы загрузить настройки, в образе он не остается +RUN SECRET_KEY=collectstatic python manage.py collectstatic --noinput \ + && mkdir -p /app/extra/serverus/ca /app/extra/backups \ + && chown -R app:app /app/extra + +USER app +EXPOSE 8000 +CMD ["gunicorn", "MyPointVPN.wsgi", "--bind", "0.0.0.0:8000", "--workers", "2", "--threads", "4", "--timeout", "120", "--access-logfile", "-"] diff --git a/Makefile b/Makefile index 5bf25db..10e4428 100644 --- a/Makefile +++ b/Makefile @@ -19,7 +19,7 @@ SERVERUS_CA = $(if $(filter yes,$(ca)),--ca $(SERVERUS_CA_DIR)) MAKECMDGOALS_TARGETS := install freeze wheels sync sync-offline clean help PKG := $(filter-out $(MAKECMDGOALS_TARGETS),$(MAKECMDGOALS)) -.PHONY: install freeze wheels sync sync-offline clean help venv dj_makemigrations dj_migrate dj_superuser dj_run dj_test dj_startapp dj_worker tg_run serve_ssh serverus_xui serverus_key serverus_cert serverus_outbound serverus_release +.PHONY: install freeze wheels sync sync-offline clean help venv dj_makemigrations dj_migrate dj_superuser dj_run dj_test dj_startapp dj_worker dj_secrets dj_backup dj_restore dj_check_points docker_build docker_up serve_ssh serverus_xui serverus_key serverus_cert serverus_outbound serverus_release help: @echo "make install [ ...] - установить пакет(ы) в $(VENV), обновить $(REQUIREMENTS), собрать wheel в $(WHEELS_DIR)/" @@ -28,6 +28,10 @@ help: @echo "make freeze - перезаписать $(REQUIREMENTS) текущим состоянием окружения" @echo "make wheels - собрать wheel для всех пакетов из $(REQUIREMENTS)" @echo "make clean - удалить $(WHEELS_DIR)/" + @echo "make dj_secrets - сгенерировать недостающие секреты в .env" + @echo "make dj_backup [keep=14] - бэкап базы в extra/backups, keep - сколько последних оставить" + @echo "make dj_restore file= - восстановить базу из бэкапа (сайт и воркер остановлены)" + @echo "make docker_build / docker_up - собрать образ и поднять все в docker (профиль app)" @echo "make serverus_xui ip= password= [domain=] [ssl=no] [ca=yes] - настроить свежий сервер и поставить 3x-ui, доступы в $(SERVERUS_ACCESS_DIR)//" @echo "make serverus_key ip= action=add|del|stats name= - ключ на панели, настроенной через serverus_xui" @echo "make serverus_cert ip= [domain=] [ca=yes] - сертификат панели на домен или на ip" @@ -86,12 +90,34 @@ dj_test: dj_worker: $(MANAGE) rundramatiq +# секреты в .env: дописывает только отсутствующие и пустые, FIELD_ENCRYPTION_KEYS не перезаписывает +dj_secrets: + $(MANAGE) generate_secrets --env-file ../.env + +# бэкап базы в extra/backups: make dj_backup [keep=14] +dj_backup: + $(MANAGE) backup $(if $(keep),--keep $(keep)) + +# восстановление: make dj_restore file=extra/backups/<файл>.dump. Сайт и воркер должны быть остановлены +dj_restore: + @test -n "$(file)" || { echo "Укажите бэкап: make dj_restore file=extra/backups/<файл>.dump"; exit 1; } + $(MANAGE) restore $(abspath $(file)) + +# перезапустить ежедневные проверки точек (нужен воркер) +dj_check_points: + $(MANAGE) check_points + +# образ приложения: wheels/ в .gitignore, поэтому собираются перед сборкой +docker_build: wheels + docker compose --profile app build + +# инфраструктура, миграции, сайт и воркер в docker +docker_up: + docker compose --profile app up -d + dj_startapp: $(MANAGE) startapp $(app) -tg_run: - $(MANAGE) telegram_run - # ручной вход на сервер точки с ключом из базы: make serve_ssh ip= [cmd='uptime'] serve_ssh: @test -n "$(ip)" || { echo "Укажите сервер: make serve_ssh ip= [cmd='<команда>']"; exit 1; } diff --git a/README.md b/README.md index 70d3b41..2958bf5 100644 --- a/README.md +++ b/README.md @@ -299,10 +299,10 @@ https://1.2.3.4:12839/9g72PmSGAQ3V/ ### Запуск -1. `cp env_example .env` и заполнить: база, RabbitMQ, `FIELD_ENCRYPTION_KEYS`, `SERVERUS_CERTIFICATES`. +1. `cp env_example .env` и заполнить: база, RabbitMQ, `FIELD_ENCRYPTION_KEYS`, `SERVERUS_CERTIFICATES`, `TELEGRAM_WEBHOOK_URL`, `ALLOWED_HOSTS`. 2. `docker compose up -d` - Postgres, RabbitMQ, свой Bot API сервер. 3. `make sync`, `make dj_migrate`, `make dj_superuser`. -4. Три процесса: `make dj_run` (сайт и админка), `make dj_worker` (установка и возврат серверов), `make tg_run` (бот, запускать строго один экземпляр). +4. Два процесса: `make dj_run` (сайт, админка и webhook ботов), `make dj_worker` (установка и возврат серверов). 5. В админке добавить бота с токеном, выбрать интерфейс MyPoint и опубликовать. Подробности: `src/App/README.md` (сценарии, секреты, ручной вход `make serve_ssh`), `src/serverus/README.md` (плейбуки, ручная проверка на сервере), `src/Telegram/README.md` (боты). @@ -329,39 +329,30 @@ https://1.2.3.4:12839/9g72PmSGAQ3V/ Сервис - На серверах у нас живут сервисы. Поднимается или настраивается какой-то сервис, чтобы решить какую-то полезную задачу. И реализуют api для управления сервисом на сервере. Сейчас эту функцию выполняет Роль сервера, потому что подразумевается что сервис может быть 1. -### Что осталось +### Запуск в docker -Проверено на локальном стенде (Debian 13, свой CA): установка точки из бота с маскировкой, DoT и запасной DoH, ключи и подписка, socks5 и WARP на выходе, отказ нерабочего прокси, переустановка с восстановлением ключей, возврат серверов, `serve_ssh`, шифрование секретов в Postgres. +Образ (`Dockerfile`) - сайт с webhook ботов на gunicorn и воркер dramatiq, настройки `MyPointVPN.settings.docker`: все из окружения, `DEBUG=False`, статика админки через whitenoise. nginx не нужен: webhook принимает сайт внутри сети compose, Bot API сервер ходит на него по имени сервиса (`TELEGRAM_WEBHOOK_URL=http://web:8000`), домен не нужен. -**Функции** +1. `make dj_secrets` - дописывает в `.env` недостающие секреты (`SECRET_KEY`, `FIELD_ENCRYPTION_KEYS`, пароли базы и RabbitMQ), заданные не трогает. +2. `make docker_build` - собирает wheels и образ. +3. `make docker_up` - инфраструктура, миграции (сервис `migrate`), `web` и `worker`. Адреса сервисов внутри сети compose задаются в `docker-compose.yaml` поверх `.env`, так что один `.env` годится и для разработки с хоста. +4. Админка только на `127.0.0.1:8000` хоста, с другой машины - через `ssh -L 8000:127.0.0.1:8000`. `manage.py createsuperuser`: `docker compose --profile app exec web python manage.py createsuperuser`. -- [x] Кнопки "Домен" и "Прокси" у готовой точки, смена на горячую. -- [ ] Скрипт маскировки (сейчас маскировка - перенаправление портов nftables). -- [ ] Статистика точки в главном меню и карточке: трафик и число ключей онлайн (сейчас статистика только у ключей). -- [ ] Кто может пользоваться ботом: сейчас любой, кто написал боту, может добавить сервер. Решить: белый список, инвайты или открыто. +Бэкап: `make dj_backup keep=14` (или `docker compose --profile app exec web python manage.py backup --keep 14` по cron) - `pg_dump` в `extra/backups`, секреты в нем зашифрованы. Восстановление: остановить `web` и `worker`, `make dj_restore file=extra/backups/<файл>.dump`, команда проверит, что `FIELD_ENCRYPTION_KEYS` подходит. `pg_dump` нужен той же версии, что сервер (16): в образе он есть, на хосте - `apt install postgresql-client-16`. `FIELD_ENCRYPTION_KEYS` храните отдельно от бэкапов: без него база бесполезна, с ним бэкап - доступ ко всем серверам. -**Надежность** +### Надежность -- [ ] Точка, зависшая в "устанавливается" или "возвращается", если воркер упал посреди задачи: в боте с ней ничего не сделать, а в админке статус только для чтения. Нужно действие в админке (сбросить в "ошибка") или проверка по времени. -- [ ] Периодическая проверка точек: панель отвечает, сертификат не истекает (короткий сертификат на ip продлевается каждые несколько дней). +- Долгие задачи (установка, переустановка, возврат, смена домена и прокси) идут с супервизором: он запускается вместе с задачей, просыпается раз в минуту, пока задача работает, и снимает ее, если она перестала отмечаться (упал или перезапущен воркер, вышел лимит времени). Точка получает понятный статус, пользователь - уведомление. Повторно доставленное после падения воркера сообщение задачу не повторяет. +- Раз в сутки каждая готовая точка проверяется: панель отвечает, Xray работает, сертификат подписки не истекает. О проблеме и о восстановлении владелец получает сообщение, текущая проблема видна в карточке ("Проверка: ..."). -**Проверка на настоящих серверах** +### TODO -- [ ] Let's Encrypt на домен и на ip, продление по cron. -- [ ] Возврат серверов и Debian 12. +Проверено на локальном стенде (Debian 13, свой CA): установка точки из бота с маскировкой, DoT и запасной DoH, ключи и подписка, socks5 и WARP на выходе, отказ нерабочего прокси, переустановка с восстановлением ключей, возврат серверов, `serve_ssh`, шифрование секретов в Postgres, бэкап и восстановление через образ. + +- [ ] Let's Encrypt на домен и на ip на сервере с публичным адресом, продление короткого сертификата по cron. +- [ ] Возврат серверов на настоящем сервере, Debian 12. - [ ] Точка с маскировкой и настоящим доменом. - -**Продакшен** - -- [ ] Настройки: `SECRET_KEY`, `DEBUG`, `ALLOWED_HOSTS` из `.env` (сейчас в `base.py` небезопасные значения для разработки). -- [ ] Запуск на сервере: systemd юниты или контейнеры для сайта (gunicorn за nginx), воркера и `telegram_run` (строго один экземпляр), логи. -- [ ] Бэкапы Postgres. `FIELD_ENCRYPTION_KEYS` хранить отдельно от бэкапов: без него база бесполезна. - -**Уборка** - -- [ ] Удалить `django-fernet-fields` (не работает на Django 6.1, заменен своим полем) и обновить `requirements.txt`. -- [ ] Удалить `src/db.sqlite3` и `extra/db-before-encryption-*.sqlite3`: в бэкапе секреты открытым текстом. -- [ ] Контейнер Postgres называется `shop-postgres`: переименовать. -- [ ] `Key` в админке (только для чтения). -- [x] Карточка точки показывает прокси вместе с логином и паролем: маскировать пароль. +- [ ] Прогнать установку в docker целиком: `web` + `telegram` + webhook по имени сервиса. - [ ] Пароль root при переустановке уходит в задачу dramatiq открытым текстом (сообщения RabbitMQ лежат на диске): шифровать, как прокси. + +Решено оставить как есть: маскировка - перенаправление портов nftables; статистика - только у ключей; ботом может пользоваться любой. diff --git a/docker-compose.yaml b/docker-compose.yaml index a8216b1..0aa9b29 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -16,11 +16,14 @@ services: environment: TELEGRAM_API_ID: $TG_APP_ID TELEGRAM_API_HASH: $TG_APP_HASH + # --local: webhook на http и любой порт, в том числе на сайт на хосте + TELEGRAM_LOCAL: "1" volumes: - telegram-bot-api-data:/var/lib/telegram-bot-api - ports: - - "127.0.0.1:8081:8081" - - "127.0.0.1:8082:8082" + network_mode: host +# ports: +# - "127.0.0.1:8081:8081" +# - "127.0.0.1:8082:8082" redis: image: docker.local/redis:8.8.3-alpine @@ -49,7 +52,6 @@ services: postgres: image: docker.local/postgres:16.15-alpine3.23 - container_name: shop-postgres environment: POSTGRES_DB: $DB_NAME POSTGRES_USER: $DB_USER @@ -59,6 +61,57 @@ services: volumes: - postgres_data:/var/lib/postgresql/data + # приложение: docker compose --profile app up -d --build. Без профиля поднимается только инфраструктура для разработки + migrate: + <<: &app + build: . + image: mypointvpn:latest + profiles: [ app ] + env_file: .env + # адреса сервисов внутри сети compose поверх .env, который настроен на разработку с хоста + environment: + DJANGO_SETTINGS_MODULE: MyPointVPN.settings.docker + DB_HOST: postgres + DB_PORT: "5432" + RABBITMQ_ADDRESS: rabbitmq:5672 + TELEGRAM_API_URL: http://telegram:8081 + # Bot API сервер ходит на webhook по имени сервиса сайта, домен не нужен + TELEGRAM_WEBHOOK_URL: http://web:8000 + ALLOWED_HOSTS: web,127.0.0.1,localhost + volumes: + # свой CA (SERVERUS_CERTIFICATES=local) и бэкапы - на хосте + - ./extra/serverus/ca:/app/extra/serverus/ca + - ./extra/backups:/app/extra/backups + command: python manage.py migrate --noinput + depends_on: [ postgres ] + + web: + <<: *app + # админка только с хоста: ssh -L 8000:127.0.0.1:8000 сервер + ports: + - "127.0.0.1:8000:8000" + restart: unless-stopped + depends_on: + migrate: + condition: service_completed_successfully + rabbitmq: + condition: service_started + telegram: + condition: service_started + + worker: + <<: *app + # check_points заново запускает ежедневные проверки точек: цепочки переживают потерю очереди + command: sh -c "python manage.py check_points && exec python manage.py rundramatiq --processes 1 --threads 8" + restart: unless-stopped + # установка может идти до часа: даем задачам закончиться, иначе супервизор снимет их как прерванные + stop_grace_period: 5m + depends_on: + migrate: + condition: service_completed_successfully + rabbitmq: + condition: service_started + volumes: telegram-bot-api-data: redis_data: diff --git a/env_example b/env_example index cf3b009..e3e2904 100644 --- a/env_example +++ b/env_example @@ -2,7 +2,11 @@ TG_APP_ID= TG_APP_HASH= # Адрес своего Bot API сервера для Django. Пусто - api.telegram.org -TELEGRAM_API_URL=http://127.0.0.1:8081 +TELEGRAM_API_URL=http://localhost:8081 +# Адрес сайта для webhook ботов, как его видит Bot API сервер (он в сети хоста, network_mode: host) +TELEGRAM_WEBHOOK_URL=http://localhost:8000 +# Хосты сайта через запятую, в том числе хост из TELEGRAM_WEBHOOK_URL +ALLOWED_HOSTS=127.0.0.1,localhost # PostgreSQL DB_NAME= @@ -31,3 +35,12 @@ FIELD_ENCRYPTION_KEYS= SERVERUS_CERTIFICATES=letsencrypt # Папка своего CA, по умолчанию extra/serverus/ca SERVERUS_CA_DIR= + +# Django в docker (settings/docker.py). Сгенерировать секреты: make dj_secrets +SECRET_KEY= +# Для админки через другой адрес, чем 127.0.0.1:8000, например https://admin.example.com +CSRF_TRUSTED_ORIGINS= +LOG_LEVEL=INFO + +# Бэкапы базы (make dj_backup), по умолчанию extra/backups +BACKUP_DIR= diff --git a/requirements.txt b/requirements.txt index 8c7a0a6..8cca288 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,10 +10,10 @@ cffi==2.1.1 charset-normalizer==3.5.1 cryptography==50.0.1 Django==6.1.1 -django-fernet-fields==0.6 django_dramatiq==0.15.0 dramatiq==2.2.1 frozenlist==1.8.0 +gunicorn==26.2.0 idna==3.20 Jinja2==3.1.6 lockfile==0.12.2 @@ -36,4 +36,5 @@ sqlparse==0.6.0 typing_extensions==4.16.0 urllib3==2.8.0 watchdog==6.0.0 +whitenoise==6.12.0 yarl==1.25.1 diff --git a/src/App/README.md b/src/App/README.md index 1ebe0a1..fc11396 100644 --- a/src/App/README.md +++ b/src/App/README.md @@ -6,6 +6,7 @@ Django приложение связывающее все части сайта -[x] Модель данных для поддержки интерфейса и функционала -[ ] Настройки домена и прокси у готовой точки -[x] Клиенты и статистика через API панели 3x-ui +-[x] Супервизор долгих задач и ежедневная проверка точек Зависимости: `Telegram` (интерфейс бота, пользователи, уведомления) и `serverus` (настройка серверов). @@ -33,6 +34,9 @@ Django приложение связывающее все части сайта - `fields.py` - `EncryptedTextField`: секреты в базе зашифрованы Fernet ключами `FIELD_ENCRYPTION_KEYS` из `.env`. - `management/commands/serve_ssh.py` - ручной вход на сервер: `make serve_ssh ip=`. - `tasks.py` - задачи dramatiq: запускают сценарий и присылают результат в чат. Воркер: `make dj_worker`. + - Долгая задача запускается через `tasks.start(actor, point, status, chat, *args)`: точка получает статус и token задачи, вместе с задачей стартует `supervise`. Задача раз в 30 секунд отмечается в точке (`task_heartbeat`), супервизор раз в минуту проверяет отметку и, если задача молчит 3 минуты (или не начиналась час), снимает ее: статус по `Point.INTERRUPTED`, уведомление в чат. Кончилась задача - кончился и супервизор. Повторная доставка сообщения после падения воркера задачу не повторяет, а снимает. + - `check_point` - ежедневная проверка точки (`services.check_point`): своя цепочка у каждой точки, запускается после установки и командой `check_points` (ее выполняет воркер в docker при старте). Сутки проходятся перескоками по 15 минут: RabbitMQ не держит неподтвержденное отложенное сообщение дольше `consumer_timeout`. +- `backups.py` и команды `backup`, `restore` - бэкап базы `pg_dump`/`pg_restore`, `generate_secrets` - секреты для `.env`. - `bot/interface.py` - интерфейс бота `MyPointInterface` (выбирается у бота в админке). - `bot/texts.py` - кнопки и тексты. diff --git a/src/App/admin.py b/src/App/admin.py index 62506d3..eef87e3 100644 --- a/src/App/admin.py +++ b/src/App/admin.py @@ -3,7 +3,7 @@ from django.contrib.auth.admin import UserAdmin as BaseUserAdmin from django.contrib.auth.forms import AdminUserCreationForm, UserChangeForm from django.utils.translation import gettext_lazy as _ -from .models import Point, Server, User +from .models import Key, Point, Server, User class UserCreationForm(AdminUserCreationForm): @@ -48,7 +48,10 @@ class PointAdmin(admin.ModelAdmin): list_display = ('name', 'owner', 'status', 'domain', 'created_at') list_filter = ('status',) search_fields = ('name', 'servers__address') - readonly_fields = ('owner', 'name', 'status', 'error', 'created_at') + readonly_fields = ( + 'owner', 'name', 'status', 'error', 'health_error', 'checked_at', 'task_started_at', 'task_heartbeat', + 'created_at', + ) inlines = [ServerInline] @@ -70,3 +73,26 @@ class ServerAdmin(admin.ModelAdmin): for server in queryset: server.reset_fingerprint() self.message_user(request, f'Сброшено отпечатков: {len(queryset)}') + + +@admin.register(Key) +class KeyAdmin(admin.ModelAdmin): + """Только просмотр: ключи выпускаются и отзываются через панель точки (services), иначе база и панель разойдутся. + + uuid и ссылка подписки не показываются: это доступ в сеть клиента. + """ + + list_display = ('name', 'point', 'created_at') + list_filter = ('point',) + search_fields = ('name', 'point__name') + fields = ('name', 'point', 'created_at') + readonly_fields = fields + + def has_add_permission(self, request): + return False + + def has_change_permission(self, request, obj=None): + return False + + def has_delete_permission(self, request, obj=None): + return False diff --git a/src/App/backups.py b/src/App/backups.py new file mode 100644 index 0000000..816f60c --- /dev/null +++ b/src/App/backups.py @@ -0,0 +1,57 @@ +"""Бэкап базы для команд backup и restore: pg_dump и pg_restore той же major-версии, что и сервер Postgres. + +Дамп берется из базы как есть, поэтому секреты (App.fields) в нем зашифрованы. Восстановленная база читается +только с тем же FIELD_ENCRYPTION_KEYS, check_keys это проверяет. +""" +import os +import subprocess +from pathlib import Path + +from django.conf import settings +from django.core.exceptions import ImproperlyConfigured +from django.core.management.base import CommandError + +SUFFIX = '.dump' + + +def database() -> dict: + database = settings.DATABASES['default'] + if 'postgresql' not in database['ENGINE']: + raise CommandError('бэкап умеет только PostgreSQL: задайте DB_NAME и остальные DB_* в .env') + return database + + +def run(args: list[str], database: dict, **kwargs): + """Запускает утилиту postgres с доступом к базе из переменных окружения: пароль не виден в списке процессов.""" + env = { + **os.environ, + 'PGHOST': database.get('HOST') or '127.0.0.1', + 'PGPORT': str(database.get('PORT') or 5432), + 'PGUSER': database.get('USER', ''), + 'PGPASSWORD': database.get('PASSWORD', ''), + 'PGDATABASE': database['NAME'], + } + try: + result = subprocess.run(args, env=env, stderr=subprocess.PIPE, text=True, **kwargs) + except FileNotFoundError: + raise CommandError(f'{args[0]} не найден: поставьте postgresql-client той же версии, что и сервер') + if result.returncode: + raise CommandError(f'{args[0]}: {result.stderr.strip()}') + + +def check_keys(): + """Секреты в базе расшифровываются текущим FIELD_ENCRYPTION_KEYS: иначе с базы нет толку.""" + from .models import Server + try: + # поля расшифровываются при чтении из базы + list(Server.objects.all()) + except ImproperlyConfigured: + raise CommandError('FIELD_ENCRYPTION_KEYS не подходит к секретам в базе: нужен ключ, с которым они шифровались') + + +def rotate(directory: Path, keep: int) -> list[Path]: + """Удаляет из directory бэкапы, кроме keep последних. Возвращает удаленные.""" + dumps = sorted(directory.glob(f'*{SUFFIX}'), key=lambda path: path.stat().st_mtime, reverse=True) + for old in dumps[keep:]: + old.unlink() + return dumps[keep:] diff --git a/src/App/bot/interface.py b/src/App/bot/interface.py index 526305a..b4e2c5d 100644 --- a/src/App/bot/interface.py +++ b/src/App/bot/interface.py @@ -186,7 +186,7 @@ class MyPointInterface(Interface): domain=draft.get('domain', ''), outbound_proxy=draft.get('proxy', ''), ) - tasks.install_point.send(point.pk, ctx.chat.pk) + tasks.start(tasks.install_point, point, Point.Status.INSTALLING, ctx.chat) ctx.reset() ctx.reply(texts.INSTALL_STARTED, keyboard=texts.MAIN_KEYBOARD_WITH_POINTS) @@ -243,8 +243,7 @@ class MyPointInterface(Interface): point = self.idle_point(ctx) if point is None: return - point.set_status(Point.Status.RELEASING) - tasks.release_point.send(point.pk, ctx.chat.pk) + tasks.start(tasks.release_point, point, Point.Status.RELEASING, ctx.chat) ctx.reset() ctx.reply(texts.RELEASE_STARTED, keyboard=texts.MAIN_KEYBOARD_WITH_POINTS) @@ -273,8 +272,7 @@ class MyPointInterface(Interface): self.start_reinstall(ctx, point) def start_reinstall(self, ctx, point: Point, address: str | None = None, password: str | None = None): - point.set_status(Point.Status.INSTALLING) - tasks.reinstall_point.send(point.pk, ctx.chat.pk, address, password) + tasks.start(tasks.reinstall_point, point, Point.Status.INSTALLING, ctx.chat, address, password) ctx.reset() ctx.reply(texts.REINSTALL_STARTED, keyboard=texts.MAIN_KEYBOARD_WITH_POINTS) @@ -347,8 +345,7 @@ class MyPointInterface(Interface): self.start_change(ctx, point, tasks.set_outbound_proxy, fields.encrypt(proxy) if proxy else '') def start_change(self, ctx, point: Point, actor, value: str): - point.set_status(Point.Status.UPDATING) - actor.send(point.pk, ctx.chat.pk, value) + tasks.start(actor, point, Point.Status.UPDATING, ctx.chat, value) ctx.reset() ctx.reply(texts.CHANGE_STARTED, keyboard=texts.MAIN_KEYBOARD_WITH_POINTS) diff --git a/src/App/bot/texts.py b/src/App/bot/texts.py index 66bea20..f8cf1f8 100644 --- a/src/App/bot/texts.py +++ b/src/App/bot/texts.py @@ -151,6 +151,8 @@ def point_card(point: Point) -> str: ] if point.error: lines.append(f'Ошибка: {escape(point.error)}') + if point.health_error: + lines.append(f'Проверка: {escape(point.health_error)}') return '\n'.join(lines) @@ -214,12 +216,11 @@ def key_card(key: Key, traffic: dict | None) -> str: def proxy_label(proxy: str) -> str: - """Прокси для показа в чате: без пароля.""" + """Прокси для показа в чате: без логина и пароля.""" parts = urlsplit(proxy) - if not parts.password: + if '@' not in parts.netloc: return proxy - host = parts.netloc.rpartition('@')[2] - return urlunsplit(parts._replace(netloc=f'{parts.username}:***@{host}')) + return urlunsplit(parts._replace(netloc='***@' + parts.netloc.rpartition('@')[2])) def domain_prompt(point: Point) -> str: @@ -264,3 +265,15 @@ def change_failed(point: Point, error: Exception) -> str: f'Не удалось применить настройку {escape(point.name)}:\n{escape(str(error))}\n' 'Точка работает с прежними настройками.' ) + + +def task_interrupted(point: Point) -> str: + return f'{escape(point.name)}: {escape(point.error)}.' + + +def health_problem(point: Point, problem: str) -> str: + return f'Проверка {escape(point.name)}: {escape(problem)}.' + + +def health_restored(point: Point) -> str: + return f'Проверка {escape(point.name)}: снова все в порядке.' diff --git a/src/App/management/commands/backup.py b/src/App/management/commands/backup.py new file mode 100644 index 0000000..dbe8e01 --- /dev/null +++ b/src/App/management/commands/backup.py @@ -0,0 +1,47 @@ +import datetime +import os +import subprocess +import tempfile +from pathlib import Path + +from django.conf import settings +from django.core.management.base import BaseCommand + +from App import backups + + +class Command(BaseCommand): + help = ( + 'Бэкап базы: pg_dump в формате custom, согласованный снимок без остановки сайта. ' + 'Секреты в бэкапе остаются зашифрованными, ключ FIELD_ENCRYPTION_KEYS в него не входит' + ) + + def add_arguments(self, parser): + parser.add_argument('--output', type=Path, help='файл бэкапа, по умолчанию в BACKUP_DIR') + parser.add_argument('--keep', type=int, default=0, help='оставить в BACKUP_DIR столько последних бэкапов, 0 - все') + + def handle(self, output, keep, **options): + database = backups.database() + backups.check_keys() + directory = Path(settings.BACKUP_DIR) + if output is None: + stamp = datetime.datetime.now().strftime('%Y%m%d-%H%M%S') + output = directory / f"{database['NAME']}-{stamp}{backups.SUFFIX}" + output.parent.mkdir(parents=True, exist_ok=True) + # пишем во временный файл рядом: оборванный дамп не должен выглядеть как бэкап + descriptor, temporary = tempfile.mkstemp(dir=output.parent, prefix='.backup-') + os.close(descriptor) + temporary = Path(temporary) + try: + backups.run(['pg_dump', '--format=custom', f'--file={temporary}'], database) + # архив читается - значит, дописан до конца + backups.run(['pg_restore', '--list', str(temporary)], database, stdout=subprocess.DEVNULL) + temporary.chmod(0o600) + temporary.replace(output) + finally: + temporary.unlink(missing_ok=True) + self.stdout.write(f'бэкап: {output} ({output.stat().st_size // 1024} КБ)') + if keep: + for old in backups.rotate(directory, keep): + self.stdout.write(f'удален старый бэкап: {old}') + self.stdout.write('Ключ FIELD_ENCRYPTION_KEYS в бэкап не входит: храните его отдельно.') diff --git a/src/App/management/commands/check_points.py b/src/App/management/commands/check_points.py new file mode 100644 index 0000000..a2b5a41 --- /dev/null +++ b/src/App/management/commands/check_points.py @@ -0,0 +1,17 @@ +from django.core.management.base import BaseCommand + +from App import tasks +from App.models import Point + + +class Command(BaseCommand): + help = ( + 'Запускает ежедневные проверки всех точек заново: прежние цепочки проверок останавливаются. ' + 'Нужен воркер. Запускается при старте воркера в docker, так цепочки переживают потерю очереди' + ) + + def handle(self, **options): + points = list(Point.objects.all()) + for point in points: + tasks.schedule_checks(point) + self.stdout.write(f'проверки запущены: {len(points)}') diff --git a/src/App/management/commands/generate_secrets.py b/src/App/management/commands/generate_secrets.py new file mode 100644 index 0000000..e67e8ce --- /dev/null +++ b/src/App/management/commands/generate_secrets.py @@ -0,0 +1,66 @@ +import secrets +from pathlib import Path + +from cryptography.fernet import Fernet +from django.core.management.base import BaseCommand + +# FIELD_ENCRYPTION_KEYS не перезаписывается никогда: с новым ключом зашифрованные секреты в базе не прочитать. +# Только [A-Za-z0-9_=-]: .env читают make (# - комментарий, $ - переменная) и docker compose ($ - подстановка) +GENERATORS = { + 'SECRET_KEY': lambda: secrets.token_urlsafe(50), + 'FIELD_ENCRYPTION_KEYS': lambda: Fernet.generate_key().decode(), + 'DB_PASSWD': lambda: secrets.token_urlsafe(24), + 'RABBITMQ_PASSWORD': lambda: secrets.token_urlsafe(24), + 'RABBITMQ_ERLANG_COOKIE': lambda: secrets.token_hex(32), +} + + +def fill_env(path: Path) -> list[str]: + """Дописывает в env-файл секреты, которых там нет или которые пустые. Заданные не трогает. Возвращает имена.""" + lines = path.read_text().splitlines() if path.exists() else [] + filled = [] + present = set() + for index, line in enumerate(lines): + key, sep, value = line.partition('=') + key = key.strip() + if not sep or key not in GENERATORS: + continue + present.add(key) + if not value.strip(): + lines[index] = f'{key}={GENERATORS[key]()}' + filled.append(key) + for key, generate in GENERATORS.items(): + if key not in present: + lines.append(f'{key}={generate()}') + filled.append(key) + if filled: + path.write_text('\n'.join(lines) + '\n') + path.chmod(0o600) + return filled + + +class Command(BaseCommand): + help = ( + 'Секреты для .env: SECRET_KEY, FIELD_ENCRYPTION_KEYS, пароли базы и RabbitMQ. ' + 'Без --env-file печатает их, с ним - дописывает только отсутствующие и пустые' + ) + + def add_arguments(self, parser): + parser.add_argument('--env-file', type=Path, help='env-файл, например ../.env') + + def handle(self, env_file, **options): + if env_file is None: + for key, generate in GENERATORS.items(): + self.stdout.write(f'{key}={generate()}') + return + filled = fill_env(env_file) + if not filled: + self.stdout.write('все секреты уже заданы, ничего не менял') + return + self.stdout.write(f"{env_file}: сгенерированы {', '.join(filled)}") + if {'DB_PASSWD', 'RABBITMQ_PASSWORD', 'RABBITMQ_ERLANG_COOKIE'} & set(filled): + self.stdout.write( + 'Пароли базы и RabbitMQ применяются только при первом запуске контейнеров с пустыми томами.' + ) + if 'FIELD_ENCRYPTION_KEYS' in filled: + self.stdout.write('Сохраните FIELD_ENCRYPTION_KEYS отдельно от бэкапов базы: без него секреты не расшифровать.') diff --git a/src/App/management/commands/restore.py b/src/App/management/commands/restore.py new file mode 100644 index 0000000..b3de75b --- /dev/null +++ b/src/App/management/commands/restore.py @@ -0,0 +1,35 @@ +from pathlib import Path + +from django.core.management.base import BaseCommand, CommandError +from django.db import connections + +from App import backups + + +class Command(BaseCommand): + help = ( + 'Восстановление базы из бэкапа manage.py backup: текущие данные заменяются целиком, в одной транзакции. ' + 'Перед запуском остановите сайт и воркер' + ) + + def add_arguments(self, parser): + parser.add_argument('path', type=Path, help='файл бэкапа') + parser.add_argument('--yes', action='store_true', help='не спрашивать подтверждение') + + def handle(self, path, yes, **options): + database = backups.database() + if not path.is_file(): + raise CommandError(f'{path}: нет такого файла') + if not yes: + answer = input(f"База {database['NAME']} будет заменена данными из {path}. Продолжить? [yes/no] ") + if answer.strip().lower() != 'yes': + raise CommandError('отменено') + # свое соединение с базой держало бы блокировки, мешая --clean + connections.close_all() + backups.run([ + 'pg_restore', '--clean', '--if-exists', '--no-owner', '--no-privileges', + '--single-transaction', '--exit-on-error', f"--dbname={database['NAME']}", str(path), + ], database) + self.stdout.write(f'восстановлено из {path}') + backups.check_keys() + self.stdout.write('Если бэкап старше кода, выполните manage.py migrate.') diff --git a/src/App/migrations/0009_point_task_supervisor.py b/src/App/migrations/0009_point_task_supervisor.py new file mode 100644 index 0000000..b31b14a --- /dev/null +++ b/src/App/migrations/0009_point_task_supervisor.py @@ -0,0 +1,50 @@ +# Generated by Django 6.1.1 on 2026-09-26 13:45 + +import django.db.models.deletion +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('App', '0008_point_status_updating'), + ('Telegram', '0003_bot_webhook'), + ] + + operations = [ + migrations.AddField( + model_name='point', + name='check_token', + field=models.CharField(blank=True, editable=False, max_length=32), + ), + migrations.AddField( + model_name='point', + name='checked_at', + field=models.DateTimeField(blank=True, editable=False, null=True), + ), + migrations.AddField( + model_name='point', + name='health_error', + field=models.TextField(blank=True, editable=False, help_text='Что не так по последней проверке'), + ), + migrations.AddField( + model_name='point', + name='task_chat', + field=models.ForeignKey(blank=True, editable=False, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='+', to='Telegram.chat'), + ), + migrations.AddField( + model_name='point', + name='task_heartbeat', + field=models.DateTimeField(blank=True, editable=False, null=True), + ), + migrations.AddField( + model_name='point', + name='task_started_at', + field=models.DateTimeField(blank=True, editable=False, null=True), + ), + migrations.AddField( + model_name='point', + name='task_token', + field=models.CharField(blank=True, editable=False, max_length=32), + ), + ] diff --git a/src/App/models.py b/src/App/models.py index effd6b5..24663c2 100644 --- a/src/App/models.py +++ b/src/App/models.py @@ -5,6 +5,7 @@ from dataclasses import dataclass from django.contrib.auth.base_user import BaseUserManager from django.contrib.auth.models import AbstractUser from django.db import models, transaction +from django.utils import timezone from serverus import CF_WARP, Host, generate_password, generate_port, generate_username from .fields import EncryptedTextField @@ -85,6 +86,9 @@ class PointQuerySet(models.QuerySet): def is_name_taken(self, name: str) -> bool: return self.filter(name=name).exists() + def busy(self) -> 'PointQuerySet': + return self.filter(status__in=Point.BUSY_STATUSES) + class PointManager(models.Manager.from_queryset(PointQuerySet)): def create_point( @@ -113,6 +117,46 @@ class PointManager(models.Manager.from_queryset(PointQuerySet)): Server.objects.create_server(point, Server.Role.MASK, mask) return point + # долгая задача точки: token связывает точку, задачу и ее супервизора (App.tasks) + + def start_task(self, point: 'Point', status: str, token: str, chat=None): + """Точка занята задачей token. chat - куда сообщить, если задача прервется.""" + point.status = status + point.error = '' + point.task_token = token + point.task_chat = chat + point.task_started_at = None + point.task_heartbeat = timezone.now() + point.save(update_fields=['status', 'error', 'task_token', 'task_chat', 'task_started_at', 'task_heartbeat']) + + def claim_task(self, pk, token: str) -> bool: + """Задача token начала работу. False - точка уже не ее или это повтор после падения воркера.""" + now = timezone.now() + return bool(self.filter(pk=pk, task_token=token, task_started_at__isnull=True) + .update(task_started_at=now, task_heartbeat=now)) + + def beat(self, pk, token: str): + """Задача token жива.""" + self.filter(pk=pk, task_token=token).update(task_heartbeat=timezone.now()) + + def finish_task(self, pk, token: str): + """Задача token закончилась, статус она выставила сама.""" + self.filter(pk=pk, task_token=token).update(**Point.NO_TASK) + + def abort_task(self, pk, token: str) -> bool: + """Задача token прервалась: статус по INTERRUPTED. False - задача уже снята или успела выставить итог.""" + with transaction.atomic(): + point = self.select_for_update().filter(pk=pk, task_token=token).first() + if point is None: + return False + if point.status not in Point.INTERRUPTED: + # сценарий закончился, упало то, что после него (например уведомление): итог уже в статусе + self.filter(pk=pk).update(**Point.NO_TASK) + return False + status, error = Point.INTERRUPTED[point.status] + self.filter(pk=pk).update(status=status, error=error, **Point.NO_TASK) + return True + def _unique_name(self) -> str: while True: name = f'point-{secrets.token_hex(3)}' @@ -137,6 +181,16 @@ class Point(models.Model): ERROR = 'error', 'ошибка' RELEASING = 'releasing', 'возвращается' + BUSY_STATUSES = (Status.INSTALLING, Status.UPDATING, Status.RELEASING) + # статус и ошибка, если задача прервалась: упал воркер, вышел лимит времени + INTERRUPTED = { + Status.INSTALLING: (Status.ERROR, 'установка прервалась, нажмите "Переустановить"'), + # домен и прокси serverus применяет в самом конце, прежние настройки почти наверняка на месте + Status.UPDATING: (Status.READY, 'смена настройки прервалась, настройки прежние'), + Status.RELEASING: (Status.ERROR, 'возврат прервался, нажмите "Вернуть" еще раз'), + } + NO_TASK = {'task_token': '', 'task_chat': None, 'task_started_at': None, 'task_heartbeat': None} + owner = models.ForeignKey('Telegram.User', on_delete=models.PROTECT, related_name='points') name = models.CharField(max_length=64, unique=True) status = models.CharField(max_length=16, choices=Status.choices, default=Status.INSTALLING) @@ -154,6 +208,19 @@ class Point(models.Model): sub_port = models.PositiveIntegerField(help_text='Порт подписки панели') sub_path = models.CharField(max_length=64, help_text='Путь подписки без слешей') + # долгая задача (установка, смена настройки, возврат) и ее супервизор, см. App.tasks + task_token = models.CharField(max_length=32, blank=True, editable=False) + task_chat = models.ForeignKey( + 'Telegram.Chat', on_delete=models.SET_NULL, null=True, blank=True, editable=False, related_name='+', + ) + task_started_at = models.DateTimeField(null=True, blank=True, editable=False) + task_heartbeat = models.DateTimeField(null=True, blank=True, editable=False) + + # ежедневная проверка, см. App.tasks.check_point + check_token = models.CharField(max_length=32, blank=True, editable=False) + checked_at = models.DateTimeField(null=True, blank=True, editable=False) + health_error = models.TextField(blank=True, editable=False, help_text='Что не так по последней проверке') + created_at = models.DateTimeField(auto_now_add=True) objects = PointManager() @@ -167,7 +234,7 @@ class Point(models.Model): @property def is_busy(self) -> bool: - return self.status in (self.Status.INSTALLING, self.Status.UPDATING, self.Status.RELEASING) + return self.status in self.BUSY_STATUSES @property def panel_host(self) -> str: diff --git a/src/App/services.py b/src/App/services.py index ccf7dab..8b1405d 100644 --- a/src/App/services.py +++ b/src/App/services.py @@ -1,4 +1,7 @@ """Сценарии точки доступа поверх serverus: установка, переустановка и возврат серверов.""" +import socket +import ssl +import time from contextlib import contextmanager import serverus @@ -122,6 +125,40 @@ def restore_keys(point: Point): api.add_client(point.panel_inbound_id, key.name, uuid=key.uuid, sub_id=key.sub_id) +def check_point(point: Point) -> str: + """Проверка готовой точки: панель отвечает, Xray работает, сертификат подписки не истекает. + + Ходит так же, как клиенты: по домену через маскировку, без домена - на ip панели. Возвращает, что не так, или ''. + """ + try: + status = xui_api(point).status() + except serverus.XuiApiError as error: + return f'панель не отвечает: {error}' + xray = status.get('xray', {}).get('state') + if xray and xray != 'running': + return f'Xray не работает: {xray}' + try: + days = _certificate_days_left(point.panel_host, point.sub_port) + except (OSError, ssl.SSLError) as error: + return f'подписка не отвечает: {error}' + if days < CERTIFICATE_WARN_DAYS: + return f'сертификат истекает через {days:.1f} дн., не продлился' + return '' + + +# сертификат на ip живет ~6 дней, acme.sh продлевает его заранее +CERTIFICATE_WARN_DAYS = 2 + + +def _certificate_days_left(host: str, port: int) -> float: + ca = local_ca() + context = ssl.create_default_context(cafile=str(ca.cert_path) if ca else None) + with socket.create_connection((host, port), timeout=15) as sock: + with context.wrap_socket(sock, server_hostname=host) as tls: + not_after = ssl.cert_time_to_seconds(tls.getpeercert()['notAfter']) + return (not_after - time.time()) / 86400 + + def local_ca() -> serverus.LocalCA | None: """Свой CA, если панели работают не с Let's Encrypt (SERVERUS_CERTIFICATES=local).""" if settings.SERVERUS_CERTIFICATES == 'local': diff --git a/src/App/tasks.py b/src/App/tasks.py index ea53156..2ad1306 100644 --- a/src/App/tasks.py +++ b/src/App/tasks.py @@ -1,7 +1,21 @@ -"""Фоновые задачи: долгие операции с серверами и уведомление пользователя о результате.""" +"""Фоновые задачи: долгие операции с серверами, их супервизор и ежедневная проверка точек. + +Долгая задача запускается через start(): точка получает статус и token задачи, вместе с задачей запускается +супервизор. Пока задача работает, она раз в HEARTBEAT_SECONDS отмечается в точке. Супервизор просыпается +раз в минуту, пока задача не закончится, и снимает ее, если отметки прекратились: воркер упал, его убили, +сообщение потерялось. Точка получает статус по Point.INTERRUPTED, пользователь - уведомление. + +Повторная доставка сообщения (RabbitMQ возвращает в очередь сообщение упавшего воркера) задачу не повторяет, +а снимает: повторять половину работы с серверами опасно, пользователь сам решит, переустанавливать ли. +""" +import datetime import logging +import secrets +import threading import dramatiq +from django.db import connection +from django.utils import timezone from Telegram.messages import send from Telegram.models import Chat @@ -13,13 +27,87 @@ logger = logging.getLogger(__name__) # установка обновляет систему и ставит панель, это долго TIME_LIMIT_MS = 60 * 60 * 1000 +HEARTBEAT_SECONDS = 30 +SUPERVISE_DELAY_MS = 60 * 1000 +# без отметок дольше - задача мертва +TASK_STALE = datetime.timedelta(minutes=3) +# задача ждет свободный поток воркера; дольше - сообщение потерялось +QUEUE_TIMEOUT = datetime.timedelta(hours=1) + +CHECK_INTERVAL = datetime.timedelta(days=1) +# сутки проходятся перескоками: отложенное сообщение воркер держит неподтвержденным, +# а RabbitMQ закрывает канал, если сообщение не подтверждено дольше consumer_timeout (30 минут) +CHECK_HOP_MS = 15 * 60 * 1000 -def _notify(chat_id: int, text: str, keyboard=None): +def _notify(chat_id: int | None, text: str, keyboard=None): + if chat_id is None: + return send(Chat.objects.get(pk=chat_id), text, keyboard=keyboard, parse_mode='HTML') -def _run_installation(point: Point, chat_id: int, operation, *args): +# долгие задачи и супервизор + +def start(actor, point: Point, status: str, chat, *args): + """Запускает долгую задачу точки вместе с супервизором. chat - куда сообщить о результате.""" + token = secrets.token_hex(16) + Point.objects.start_task(point, status, token, chat) + actor.send(point.pk, token, *args) + supervise.send_with_options(args=(point.pk, token), delay=SUPERVISE_DELAY_MS) + + +@dramatiq.actor(max_retries=3) +def supervise(point_id: int, token: str): + point = Point.objects.filter(pk=point_id, task_token=token).first() + if point is None: + # задача закончилась, супервизор больше не нужен + return + limit = TASK_STALE if point.task_started_at else QUEUE_TIMEOUT + if timezone.now() - point.task_heartbeat > limit: + logger.error('Point %s task is stuck: no heartbeat since %s', point, point.task_heartbeat) + _abort(point_id, token) + return + supervise.send_with_options(args=(point_id, token), delay=SUPERVISE_DELAY_MS) + + +def _run(point_id: int, token: str, work): + """Выполняет work(point) как задачу token: с отметками для супервизора и снятием, если она прервется.""" + if not Point.objects.claim_task(point_id, token): + # точка удалена, задачу уже снял супервизор или это повторная доставка после падения воркера + _abort(point_id, token) + return + stop = threading.Event() + heartbeat = threading.Thread(target=_heartbeat, args=(point_id, token, stop), daemon=True) + heartbeat.start() + try: + work(Point.objects.get(pk=point_id)) + except BaseException: + # лимит времени dramatiq и остановка воркера - не Exception, сценарий их не ловит + logger.exception('Point %s task interrupted', point_id) + _abort(point_id, token) + raise + finally: + stop.set() + heartbeat.join() + Point.objects.finish_task(point_id, token) + + +def _heartbeat(point_id: int, token: str, stop: threading.Event): + try: + while not stop.wait(HEARTBEAT_SECONDS): + Point.objects.beat(point_id, token) + finally: + connection.close() + + +def _abort(point_id: int, token: str): + chat_id = Point.objects.filter(pk=point_id).values_list('task_chat', flat=True).first() + if Point.objects.abort_task(point_id, token): + _notify(chat_id, texts.task_interrupted(Point.objects.get(pk=point_id)), keyboard=[[texts.SERVER_MENU]]) + + +def _install(point: Point, operation, *args): + chat_id = point.task_chat_id try: operation(point, *args) except Exception as error: @@ -27,46 +115,85 @@ def _run_installation(point: Point, chat_id: int, operation, *args): _notify(chat_id, texts.install_failed(point, error), keyboard=[[texts.SERVER_MENU]]) return _notify(chat_id, texts.installed(point), keyboard=[[texts.SERVER_MENU]]) + schedule_checks(point) -# повтор не нужен: пользователь сам решает, переустанавливать ли после ошибки -@dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) -def install_point(point_id: int, chat_id: int): - _run_installation(Point.objects.get(pk=point_id), chat_id, services.install_point) - - -@dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) -def reinstall_point(point_id: int, chat_id: int, address: str | None = None, password: str | None = None): - _run_installation(Point.objects.get(pk=point_id), chat_id, services.reinstall, address, password) - - -@dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) -def release_point(point_id: int, chat_id: int): - point = Point.objects.get(pk=point_id) - name = point.name - passwords, errors = services.release_point(point) - _notify(chat_id, texts.released(name, passwords, errors)) - - -def _run_change(point: Point, chat_id: int, operation, value: str, done: str): +def _change(point: Point, operation, value: str, done): + chat_id = point.task_chat_id try: operation(point, value) except Exception as error: logger.exception('Point %s change failed', point) _notify(chat_id, texts.change_failed(point, error), keyboard=[[texts.SERVER_MENU]]) return - _notify(chat_id, done, keyboard=[[texts.SERVER_MENU]]) + _notify(chat_id, done(point, value), keyboard=[[texts.SERVER_MENU]]) + + +def _release(point: Point): + chat_id, name = point.task_chat_id, point.name + passwords, errors = services.release_point(point) + _notify(chat_id, texts.released(name, passwords, errors)) + + +# повтор не нужен: пользователь сам решает, переустанавливать ли после ошибки +@dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) +def install_point(point_id: int, token: str): + _run(point_id, token, lambda point: _install(point, services.install_point)) @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) -def set_domain(point_id: int, chat_id: int, domain: str): - point = Point.objects.get(pk=point_id) - _run_change(point, chat_id, services.set_domain, domain, texts.domain_changed(point, domain)) +def reinstall_point(point_id: int, token: str, address: str | None = None, password: str | None = None): + _run(point_id, token, lambda point: _install(point, services.reinstall, address, password)) @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) -def set_outbound_proxy(point_id: int, chat_id: int, encrypted_proxy: str): +def release_point(point_id: int, token: str): + _run(point_id, token, _release) + + +@dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) +def set_domain(point_id: int, token: str, domain: str): + _run(point_id, token, lambda point: _change(point, services.set_domain, domain, texts.domain_changed)) + + +@dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) +def set_outbound_proxy(point_id: int, token: str, encrypted_proxy: str): """encrypted_proxy - fields.encrypt(proxy): в прокси логин и пароль, а сообщения брокера лежат на диске.""" - point = Point.objects.get(pk=point_id) proxy = fields.decrypt(encrypted_proxy) if encrypted_proxy else '' - _run_change(point, chat_id, services.set_outbound_proxy, proxy, texts.proxy_changed(point, proxy)) + _run(point_id, token, lambda point: _change(point, services.set_outbound_proxy, proxy, texts.proxy_changed)) + + +# ежедневная проверка: своя цепочка у каждой точки, живет, пока жива точка + +def schedule_checks(point: Point): + """(Пере)запускает цепочку проверок точки. Прежняя цепочка остановится: у нее другой token.""" + token = secrets.token_hex(16) + Point.objects.filter(pk=point.pk).update(check_token=token) + check_point.send(point.pk, token) + + +@dramatiq.actor(max_retries=3) +def check_point(point_id: int, token: str): + point = Point.objects.filter(pk=point_id, check_token=token).first() + if point is None: + # точка удалена или цепочку сменила новая + return + # следующий перескок сначала: ошибка проверки не должна обрывать цепочку + check_point.send_with_options(args=(point_id, token), delay=CHECK_HOP_MS) + due = point.checked_at is None or timezone.now() - point.checked_at >= CHECK_INTERVAL + if point.status == Point.Status.READY and due: + _check(point) + + +def _check(point: Point): + problem = services.check_point(point) + previous = point.health_error + Point.objects.filter(pk=point.pk).update(checked_at=timezone.now(), health_error=problem) + if problem == previous: + return + text = texts.health_problem(point, problem) if problem else texts.health_restored(point) + _notify(_owner_chat(point), text, keyboard=[[texts.SERVER_MENU]]) + + +def _owner_chat(point: Point) -> int | None: + return Chat.objects.filter(user=point.owner, deleted_at__isnull=True).order_by('-pk').values_list('pk', flat=True).first() diff --git a/src/App/tests_bot.py b/src/App/tests_bot.py index e7277d0..de504d8 100644 --- a/src/App/tests_bot.py +++ b/src/App/tests_bot.py @@ -1,11 +1,14 @@ +import datetime from unittest import mock from django.test import TestCase +from django.utils import timezone +from dramatiq.middleware.time_limit import TimeLimitExceeded from Telegram import interfaces from Telegram.models import Bot, Chat, User as TelegramUser from Telegram.testing import Conversation, FakeApi -from . import fields, tasks +from . import fields, services, tasks from .bot import texts from .bot.interface import MyPointInterface, is_proxy_url, parse_credentials from .models import Credentials, Point, Server @@ -83,7 +86,9 @@ class AddServerTests(BotTestCase): self.assertEqual((point.mask.address, point.mask.root_password), ('10.0.0.2', 'mask-pw')) self.assertEqual((point.domain, point.outbound_proxy), ('example.com', 'cf_warp')) self.assertEqual(point.owner, TelegramUser.objects.get()) - send.assert_called_once_with(point.pk, Chat.objects.get().pk) + send.assert_called_once_with(point.pk, mock.ANY) + point.refresh_from_db() + self.assertEqual((point.task_chat, point.task_token), (Chat.objects.get(), send.call_args.args[1])) # сообщения с паролями удалены из чата self.assertEqual(len(self.api.deleted), 2) self.assertEqual(Chat.objects.get().state, '') @@ -158,7 +163,7 @@ class PointMenuTests(BotTestCase): self.assertIn(self.point.name, self.say(texts.NO)) self.assertIn('Продолжить?', self.say(texts.RELEASE)) self.assertEqual(self.say(texts.YES), texts.RELEASE_STARTED) - send.assert_called_once_with(self.point.pk, Chat.objects.get(user__tg_id=7).pk) + send.assert_called_once_with(self.point.pk, mock.ANY) self.point.refresh_from_db() self.assertEqual(self.point.status, Point.Status.RELEASING) @@ -311,8 +316,9 @@ class PointSettingsTests(BotTestCase): self.point.outbound_proxy = 'socks5://user:secret@proxy.example:1080' self.point.save() card = self.say(texts.SERVER_MENU) - self.assertIn('socks5://user:***@proxy.example:1080', card) + self.assertIn('socks5://***@proxy.example:1080', card) self.assertNotIn('secret', card) + self.assertNotIn('user', card) class TasksTests(TestCase): @@ -323,13 +329,21 @@ class TasksTests(TestCase): self.point = create_point(owner) self.api = FakeApi() - def run_task(self, actor, *args, runner=None): - with patch_runner(runner or FakeRunner()), mock.patch.object(Bot, 'api', return_value=self.api): - actor.fn(*args) + def run_task(self, actor, *args, status=Point.Status.INSTALLING, runner=None): + """Запускает задачу, как бот, и выполняет ее сообщение.""" + self.point.refresh_from_db() + with mock.patch.object(actor, 'send') as send: + tasks.start(actor, self.point, status, self.chat, *args) + with ( + patch_runner(runner or FakeRunner()), + mock.patch.object(Bot, 'api', return_value=self.api), + mock.patch.object(tasks.check_point, 'send'), + ): + actor.fn(*send.call_args.args) return self.api.last_text def test_install_notifies_success(self): - text = self.run_task(tasks.install_point, self.point.pk, self.chat.pk) + text = self.run_task(tasks.install_point) self.assertIn(self.point.name, text) self.assertIn('10.0.0.2', text) self.assertIn('готов к работе', text) @@ -337,12 +351,12 @@ class TasksTests(TestCase): def test_install_notifies_failure(self): with self.assertLogs('App.tasks', 'ERROR'): - text = self.run_task(tasks.install_point, self.point.pk, self.chat.pk, runner=FakeRunner(fail='xui')) + text = self.run_task(tasks.install_point, runner=FakeRunner(fail='xui')) self.assertIn('Не удалось', text) self.assertIn('boom', text) def test_release_notifies_passwords(self): - text = self.run_task(tasks.release_point, self.point.pk, self.chat.pk) + text = self.run_task(tasks.release_point, status=Point.Status.RELEASING) self.assertIn('10.0.0.1 - пароль root', text) self.assertFalse(Point.objects.exists()) @@ -351,10 +365,169 @@ class TasksTests(TestCase): self.point.status = Point.Status.READY self.point.save() proxy = fields.encrypt('socks5://user:secret@proxy.example:1080') - text = self.run_task(tasks.set_outbound_proxy, self.point.pk, self.chat.pk, proxy) - self.assertIn('user:***@proxy.example', text) + text = self.run_task(tasks.set_outbound_proxy, proxy, status=Point.Status.UPDATING) + self.assertIn('socks5://***@proxy.example', text) with self.assertLogs('App.tasks', 'ERROR'): - text = self.run_task(tasks.set_domain, self.point.pk, self.chat.pk, 'vpn.example.com', runner=FakeRunner(fail='xui_cert')) + text = self.run_task( + tasks.set_domain, 'vpn.example.com', status=Point.Status.UPDATING, runner=FakeRunner(fail='xui_cert'), + ) self.assertIn('прежними настройками', text) self.point.refresh_from_db() self.assertEqual((self.point.status, self.point.domain), (Point.Status.READY, '')) + + +class SupervisorTests(TestCase): + """Задача, которая прервалась, не оставляет точку занятой навсегда.""" + + def setUp(self): + bot = Bot.objects.create(token='1:test') + owner = TelegramUser.objects.create(tg_id=1) + self.chat = Chat.objects.create(bot=bot, chat_id=1, type='private', user=owner) + self.point = create_point(owner) + self.api = FakeApi() + patcher = mock.patch.object(Bot, 'api', return_value=self.api) + patcher.start() + self.addCleanup(patcher.stop) + + def start(self, actor=tasks.install_point, status=Point.Status.INSTALLING, *args): + with mock.patch.object(actor, 'send') as send, mock.patch.object(tasks.supervise, 'send_with_options') as sup: + tasks.start(actor, self.point, status, self.chat, *args) + sup.assert_called_once_with(args=(self.point.pk, send.call_args.args[1]), delay=tasks.SUPERVISE_DELAY_MS) + return send.call_args.args + + def supervise(self, token): + with mock.patch.object(tasks.supervise, 'send_with_options') as again: + tasks.supervise.fn(self.point.pk, token) + self.point.refresh_from_db() + return again.called + + def age_heartbeat(self, delta): + Point.objects.filter(pk=self.point.pk).update(task_heartbeat=timezone.now() - delta) + + def test_stuck_task_is_aborted(self): + _, token = self.start() + Point.objects.claim_task(self.point.pk, token) + self.assertTrue(self.supervise(token)) + self.age_heartbeat(tasks.TASK_STALE + datetime.timedelta(seconds=1)) + self.assertFalse(self.supervise(token)) + self.assertEqual(self.point.status, Point.Status.ERROR) + self.assertIn('Переустановить', self.point.error) + self.assertEqual(self.point.task_token, '') + self.assertIn('установка прервалась', self.api.last_text) + + def test_queued_task_waits_longer(self): + _, token = self.start() + self.age_heartbeat(tasks.TASK_STALE * 2) + self.assertTrue(self.supervise(token)) + self.age_heartbeat(tasks.QUEUE_TIMEOUT + datetime.timedelta(seconds=1)) + self.assertFalse(self.supervise(token)) + self.assertEqual(self.point.status, Point.Status.ERROR) + # сообщение все же дошло: задача снята, работы нет + runner = FakeRunner() + with patch_runner(runner): + tasks.install_point.fn(self.point.pk, token) + self.assertEqual(runner.calls, []) + + def test_interrupted_update_keeps_point_ready(self): + token = self.start(tasks.set_domain, Point.Status.UPDATING, 'vpn.example.com')[1] + Point.objects.claim_task(self.point.pk, token) + self.age_heartbeat(tasks.TASK_STALE * 2) + self.supervise(token) + self.assertEqual((self.point.status, self.point.domain), (Point.Status.READY, '')) + self.assertIn('прервалась', self.point.error) + + def test_supervisor_stops_after_task(self): + args = self.start() + with patch_runner(FakeRunner()), mock.patch.object(tasks.check_point, 'send'): + tasks.install_point.fn(*args) + self.assertFalse(self.supervise(args[1])) + self.assertEqual((self.point.status, self.point.task_token), (Point.Status.READY, '')) + + def test_redelivered_message_is_not_repeated(self): + # воркер упал посреди установки, RabbitMQ доставил сообщение снова + args = self.start() + Point.objects.claim_task(*args) + runner = FakeRunner() + with patch_runner(runner): + tasks.install_point.fn(*args) + self.assertEqual(runner.calls, []) + self.point.refresh_from_db() + self.assertEqual(self.point.status, Point.Status.ERROR) + + def test_time_limit_aborts_task(self): + args = self.start() + with ( + mock.patch.object(services, 'install_point', side_effect=TimeLimitExceeded), + self.assertRaises(TimeLimitExceeded), + self.assertLogs('App.tasks', 'ERROR'), + ): + tasks.install_point.fn(*args) + self.point.refresh_from_db() + self.assertEqual((self.point.status, self.point.task_token), (Point.Status.ERROR, '')) + + def test_failed_notification_keeps_result(self): + args = self.start() + with ( + patch_runner(FakeRunner()), + mock.patch.object(tasks, 'schedule_checks', side_effect=RuntimeError('boom')), + self.assertRaises(RuntimeError), + self.assertLogs('App.tasks', 'ERROR'), + ): + tasks.install_point.fn(*args) + self.point.refresh_from_db() + self.assertEqual((self.point.status, self.point.task_token), (Point.Status.READY, '')) + + +class CheckTests(TestCase): + def setUp(self): + bot = Bot.objects.create(token='1:test') + owner = TelegramUser.objects.create(tg_id=1) + Chat.objects.create(bot=bot, chat_id=1, type='private', user=owner) + self.point = create_point(owner) + self.point.status = Point.Status.READY + self.point.save() + self.api = FakeApi() + patcher = mock.patch.object(Bot, 'api', return_value=self.api) + patcher.start() + self.addCleanup(patcher.stop) + with mock.patch.object(tasks.check_point, 'send') as send: + tasks.schedule_checks(self.point) + self.token = send.call_args.args[1] + + def check(self, problem='', token=None): + with ( + mock.patch.object(services, 'check_point', return_value=problem) as check, + mock.patch.object(tasks.check_point, 'send_with_options') as again, + ): + tasks.check_point.fn(self.point.pk, token or self.token) + self.point.refresh_from_db() + return check.called, again.called + + def test_checks_once_a_day_and_reports_changes(self): + self.assertEqual(self.check('панель не отвечает'), (True, True)) + self.assertEqual(self.point.health_error, 'панель не отвечает') + self.assertIn('панель не отвечает', self.api.last_text) + self.assertIn('Проверка: панель не отвечает', texts.point_card(self.point)) + # до следующих суток не проверяет, но цепочка живет + self.assertEqual(self.check(), (False, True)) + Point.objects.filter(pk=self.point.pk).update(checked_at=timezone.now() - tasks.CHECK_INTERVAL) + sent = len(self.api.sent) + self.check('панель не отвечает') + self.assertEqual(len(self.api.sent), sent) + Point.objects.filter(pk=self.point.pk).update(checked_at=timezone.now() - tasks.CHECK_INTERVAL) + self.check() + self.assertIn('снова все в порядке', self.api.last_text) + + def test_busy_point_is_not_checked(self): + self.point.set_status(Point.Status.INSTALLING) + self.assertEqual(self.check(), (False, True)) + + def test_old_chain_stops(self): + with mock.patch.object(tasks.check_point, 'send'): + tasks.schedule_checks(self.point) + self.assertEqual(self.check(token=self.token), (False, False)) + point_id = self.point.pk + self.point.delete() + with mock.patch.object(tasks.check_point, 'send_with_options') as again: + tasks.check_point.fn(point_id, self.token) + again.assert_not_called() diff --git a/src/App/tests_commands.py b/src/App/tests_commands.py new file mode 100644 index 0000000..91d8bce --- /dev/null +++ b/src/App/tests_commands.py @@ -0,0 +1,136 @@ +import subprocess +import tempfile +from io import StringIO +from pathlib import Path +from unittest import mock + +from cryptography.fernet import Fernet +from django.core.management import CommandError, call_command +from django.test import TestCase, override_settings + +from . import backups, tasks +from .management.commands.generate_secrets import GENERATORS, fill_env +from .tests_services import create_point + +DATABASE = {'ENGINE': 'django.db.backends.postgresql', 'NAME': 'vpn', 'USER': 'vpn', 'PASSWORD': 'db-secret', 'HOST': 'db'} + + +class GenerateSecretsTests(TestCase): + def test_prints_all_secrets(self): + out = StringIO() + call_command('generate_secrets', stdout=out) + lines = dict(line.split('=', 1) for line in out.getvalue().splitlines()) + self.assertEqual(set(lines), set(GENERATORS)) + Fernet(lines['FIELD_ENCRYPTION_KEYS']) + for value in lines.values(): + self.assertRegex(value, r'^[A-Za-z0-9_=-]+$') + + def test_fills_only_missing_and_empty(self): + with tempfile.TemporaryDirectory() as directory: + env = Path(directory) / '.env' + env.write_text('# база\nDB_NAME=vpn\nDB_PASSWD=\nFIELD_ENCRYPTION_KEYS=keep-me\n') + filled = fill_env(env) + text = env.read_text() + self.assertNotIn('FIELD_ENCRYPTION_KEYS', filled) + self.assertIn('FIELD_ENCRYPTION_KEYS=keep-me\n', text) + self.assertIn('# база\nDB_NAME=vpn\n', text) + self.assertNotIn('DB_PASSWD=\n', text) + self.assertEqual(set(filled), set(GENERATORS) - {'FIELD_ENCRYPTION_KEYS'}) + self.assertEqual(env.stat().st_mode & 0o777, 0o600) + self.assertEqual(fill_env(env), []) + + +class FakePostgres: + """Подменяет pg_dump и pg_restore.""" + + def __init__(self, fail: str | None = None): + self.calls = [] + self.fail = fail + + def __call__(self, args, env, **kwargs): + self.calls.append((args, env)) + if args[0] == self.fail: + return subprocess.CompletedProcess(args, 1, stderr='connection refused') + if args[0] == 'pg_dump': + Path(args[-1].removeprefix('--file=')).write_bytes(b'PGDMP') + return subprocess.CompletedProcess(args, 0, stderr='') + + +class BackupTests(TestCase): + def setUp(self): + create_point(mask=None) + directory = tempfile.TemporaryDirectory() + self.addCleanup(directory.cleanup) + self.directory = Path(directory.name) + patcher = mock.patch.object(backups, 'database', return_value=DATABASE) + patcher.start() + self.addCleanup(patcher.stop) + settings = override_settings(BACKUP_DIR=str(self.directory)) + settings.enable() + self.addCleanup(settings.disable) + + def backup(self, postgres, *args): + with mock.patch('subprocess.run', postgres): + call_command('backup', *args, stdout=StringIO()) + + def test_backup(self): + postgres = FakePostgres() + self.backup(postgres) + [dump] = self.directory.glob('vpn-*.dump') + self.assertEqual(dump.stat().st_mode & 0o777, 0o600) + (dump_args, env), (list_args, _) = postgres.calls + self.assertEqual(dump_args[:2], ['pg_dump', '--format=custom']) + self.assertEqual(list_args[:2], ['pg_restore', '--list']) + # пароль базы только в окружении, не в аргументах + self.assertEqual((env['PGPASSWORD'], env['PGHOST'], env['PGDATABASE']), ('db-secret', 'db', 'vpn')) + self.assertNotIn('db-secret', ' '.join(dump_args)) + + def test_failed_dump_leaves_nothing(self): + with self.assertRaisesMessage(CommandError, 'connection refused'): + self.backup(FakePostgres(fail='pg_dump')) + self.assertEqual(list(self.directory.iterdir()), []) + + def test_rotation(self): + for name in ('vpn-1.dump', 'vpn-2.dump', 'other.txt'): + (self.directory / name).write_text('old') + self.backup(FakePostgres(), '--keep', '2') + self.assertEqual(len(list(self.directory.glob('*.dump'))), 2) + self.assertTrue((self.directory / 'other.txt').exists()) + + def test_wrong_key_is_reported(self): + with override_settings(FIELD_ENCRYPTION_KEYS=[Fernet.generate_key().decode()]): + with self.assertRaisesMessage(CommandError, 'FIELD_ENCRYPTION_KEYS'): + self.backup(FakePostgres()) + + def test_restore(self): + dump = self.directory / 'vpn-1.dump' + dump.write_bytes(b'PGDMP') + postgres = FakePostgres() + with mock.patch('subprocess.run', postgres): + call_command('restore', str(dump), '--yes', stdout=StringIO()) + [(args, _)] = postgres.calls + self.assertIn('--single-transaction', args) + self.assertIn('--clean', args) + self.assertEqual(args[-1], str(dump)) + + def test_restore_asks_confirmation(self): + dump = self.directory / 'vpn-1.dump' + dump.write_bytes(b'PGDMP') + with mock.patch('builtins.input', return_value='no'), self.assertRaisesMessage(CommandError, 'отменено'): + call_command('restore', str(dump)) + + + +class BackupDatabaseTests(TestCase): + def test_postgres_only(self): + # тесты идут на SQLite + with self.assertRaisesMessage(CommandError, 'PostgreSQL'): + call_command('backup') + + +class CheckPointsCommandTests(TestCase): + def test_schedules_every_point(self): + create_point(mask=None) + with mock.patch.object(tasks.check_point, 'send') as send: + call_command('check_points', stdout=StringIO()) + send.assert_called_once() diff --git a/src/App/tests_services.py b/src/App/tests_services.py index e08723d..db6a1ee 100644 --- a/src/App/tests_services.py +++ b/src/App/tests_services.py @@ -48,6 +48,7 @@ class FakeXuiApi: self.clients = {} self.down = False self.added = [] + self.xray_state = 'running' def _check(self): if self.down: @@ -76,6 +77,11 @@ class FakeXuiApi: self.get_client(email) return {'up': 1024, 'down': 3 * 1024 * 1024, 'lastOnline': 0} + def status(self): + self._check() + return {'xray': {'state': self.xray_state}} + + def patch_xui_api(api): return mock.patch('App.services.xui_api', lambda point: api) @@ -289,6 +295,33 @@ class ReadyPointTests(TestCase): self.assertEqual(runner.calls, []) +class CheckPointTests(TestCase): + def setUp(self): + self.point = create_point(mask=None) + self.panel_api = FakeXuiApi() + patcher = patch_xui_api(self.panel_api) + patcher.start() + self.addCleanup(patcher.stop) + + def check(self, days=60): + with mock.patch.object(services, '_certificate_days_left', return_value=days): + return services.check_point(self.point) + + def test_healthy(self): + self.assertEqual(self.check(), '') + + def test_problems(self): + self.assertIn('истекает', self.check(days=1)) + self.panel_api.xray_state = 'stop' + self.assertIn('Xray не работает', self.check()) + self.panel_api.down = True + self.assertIn('панель не отвечает', self.check()) + + def test_subscription_unreachable(self): + with mock.patch.object(services, '_certificate_days_left', side_effect=ConnectionRefusedError('refused')): + self.assertIn('подписка не отвечает', services.check_point(self.point)) + + class KeyTests(TestCase): def setUp(self): self.point = create_point(mask=None) diff --git a/src/MyPointVPN/settings/README.md b/src/MyPointVPN/settings/README.md index 12fefb6..7ecffcf 100644 --- a/src/MyPointVPN/settings/README.md +++ b/src/MyPointVPN/settings/README.md @@ -12,3 +12,4 @@ from .project import * ``` - `test.py` - настройки для тестов: фоновые задачи идут в stub брокер, а не в rabbitmq. `make dj_test` использует их. +- `docker.py` - настройки для docker: все из окружения, `DEBUG=False`, `SECRET_KEY` обязателен, статика через whitenoise, логи в консоль. diff --git a/src/MyPointVPN/settings/docker.py b/src/MyPointVPN/settings/docker.py new file mode 100644 index 0000000..80717e5 --- /dev/null +++ b/src/MyPointVPN/settings/docker.py @@ -0,0 +1,47 @@ +"""Настройки для запуска в docker (Dockerfile, сервисы профиля app в docker-compose.yaml). + +Все параметры - из окружения: env_file .env плюс адреса сервисов внутри сети compose. Секреты: manage.py generate_secrets. +""" +import os + +from django.core.exceptions import ImproperlyConfigured + +from .project import * + + +def _env_list(name: str) -> list[str]: + return [item.strip() for item in os.environ.get(name, '').split(',') if item.strip()] + + +DEBUG = False + +SECRET_KEY = os.environ.get('SECRET_KEY', '') +if not SECRET_KEY: + raise ImproperlyConfigured('SECRET_KEY is not set: python manage.py generate_secrets --env-file .env') + +# сайт открыт только внутри сети compose (webhook от Bot API сервера) и на 127.0.0.1 хоста (админка) +ALLOWED_HOSTS = _env_list('ALLOWED_HOSTS') +CSRF_TRUSTED_ORIGINS = _env_list('CSRF_TRUSTED_ORIGINS') + +# статика админки без nginx +MIDDLEWARE = [ + MIDDLEWARE[0], + 'whitenoise.middleware.WhiteNoiseMiddleware', + *MIDDLEWARE[1:], +] +STORAGES = { + 'default': {'BACKEND': 'django.core.files.storage.FileSystemStorage'}, + 'staticfiles': {'BACKEND': 'whitenoise.storage.CompressedManifestStaticFilesStorage'}, +} + +LOGGING = { + 'version': 1, + 'disable_existing_loggers': False, + 'formatters': { + 'plain': {'format': '{asctime} {levelname} {name}: {message}', 'style': '{'}, + }, + 'handlers': { + 'console': {'class': 'logging.StreamHandler', 'formatter': 'plain'}, + }, + 'root': {'handlers': ['console'], 'level': os.environ.get('LOG_LEVEL', 'INFO')}, +} diff --git a/src/MyPointVPN/settings/project.py b/src/MyPointVPN/settings/project.py index 60c8c1e..0295432 100644 --- a/src/MyPointVPN/settings/project.py +++ b/src/MyPointVPN/settings/project.py @@ -38,6 +38,11 @@ INSTALLED_APPS += [ # Свой Bot API сервер (сервис telegram в docker-compose.yaml). Пусто - api.telegram.org TELEGRAM_API_URL = os.environ.get('TELEGRAM_API_URL', '') +# Адрес этого сайта, как его видит Bot API сервер: на него ставится webhook ботов +TELEGRAM_WEBHOOK_URL = os.environ.get('TELEGRAM_WEBHOOK_URL', '') + +# Через запятую. Пусто при DEBUG - только localhost +ALLOWED_HOSTS = [host for host in os.environ.get('ALLOWED_HOSTS', '').split(',') if host] # PostgreSQL из docker-compose.yaml. Без DB_NAME - SQLite из base.py (разработка без docker) if os.environ.get('DB_NAME'): @@ -64,6 +69,9 @@ FIELD_ENCRYPTION_KEYS = [key for key in os.environ.get('FIELD_ENCRYPTION_KEYS', SERVERUS_CERTIFICATES = os.environ.get('SERVERUS_CERTIFICATES', 'letsencrypt') SERVERUS_CA_DIR = os.environ.get('SERVERUS_CA_DIR', str(BASE_DIR.parent / 'extra/serverus/ca')) +# Бэкапы базы: manage.py backup и restore +BACKUP_DIR = os.environ.get('BACKUP_DIR') or str(BASE_DIR.parent / 'extra/backups') + # Фоновые задачи: dramatiq поверх rabbitmq. Воркер: make dj_worker RABBITMQ_URL = 'amqp://{user}:{password}@{address}/{vhost}'.format( user=quote(os.environ.get('RABBITMQ_USER', 'guest'), safe=''), diff --git a/src/MyPointVPN/urls.py b/src/MyPointVPN/urls.py index b2ba8d1..c879c75 100644 --- a/src/MyPointVPN/urls.py +++ b/src/MyPointVPN/urls.py @@ -15,8 +15,9 @@ Including another URLconf 2. Add a URL to urlpatterns: path('blog/', include('blog.urls')) """ from django.contrib import admin -from django.urls import path +from django.urls import include, path urlpatterns = [ path('admin/', admin.site.urls), + path('telegram/', include('Telegram.urls')), ] diff --git a/src/Telegram/README.md b/src/Telegram/README.md index b4f0590..71c4040 100644 --- a/src/Telegram/README.md +++ b/src/Telegram/README.md @@ -19,7 +19,7 @@ - `context.py` - `Context`, который получает обработчик: текст, команда, собеседник, состояние, `reply()`. - `messages.py` - `send()` для отправки сообщения в чат вне обработчика, например уведомления из фоновой задачи. - `dispatcher.py` - разбор обновления: сохраняет пользователя, чат, сообщение и передает его интерфейсу бота. -- `runner.py` - long polling: по потоку на каждый опубликованный бот, супервизор сверяется с базой каждые 5 секунд. +- `views.py` - webhook `/telegram/webhook//`: проверяет секрет из заголовка `X-Telegram-Bot-Api-Secret-Token` и передает обновление диспетчеру. Бот читается из базы на каждое обновление. ## Интерфейс бота @@ -64,8 +64,11 @@ class Shop(Interface): ## Запуск 1. В админке добавить бота с токеном и выбрать интерфейс. -2. Действие "Опубликовать" проверяет токен, подтягивает имя бота и включает бота. "Отозвать" выключает. +2. Действие "Опубликовать" проверяет токен, подтягивает имя бота и ставит webhook с новым секретом. "Отозвать" выключает бота и снимает webhook. Если задан `TELEGRAM_API_URL`, при первой публикации бот отключается от облачного API (`logOut`). Вернуть его в облако можно не раньше чем через 10 минут. -3. `make tg_run` запускает раннер. Смена интерфейса, публикация и отзыв применяются без перезапуска. +3. Отдельный процесс не нужен: обновления принимает сайт. Смена интерфейса применяется без перезапуска. `TELEGRAM_API_URL` направляет бота на свой Bot API сервер (сервис `telegram` в `docker-compose.yaml`). +`TELEGRAM_WEBHOOK_URL` - адрес сайта, как его видит Bot API сервер, к нему дописывается путь webhook. Сервер в docker запущен с `--local` (`TELEGRAM_LOCAL`), поэтому webhook может быть на http и на любом порту. Хост из этого адреса нужно добавить в `ALLOWED_HOSTS`. + +Состояние webhook (адрес, очередь, последняя ошибка доставки) видно на странице бота в админке. Ошибка обработчика пишется в `last_error`, а телеграму отвечаем 200, чтобы он не повторял то же обновление. diff --git a/src/Telegram/admin.py b/src/Telegram/admin.py index c00095f..d7a43af 100644 --- a/src/Telegram/admin.py +++ b/src/Telegram/admin.py @@ -1,12 +1,8 @@ -from datetime import timedelta - from django import forms from django.contrib import admin, messages -from django.utils import timezone from . import interfaces from .models import Bot, Chat, Message, User -from .runner import HEARTBEAT_TIMEOUT class BotForm(forms.ModelForm): @@ -27,7 +23,7 @@ class BotAdmin(admin.ModelAdmin): form = BotForm list_display = ('__str__', 'name', 'interface_title', 'is_active', 'status') readonly_fields = ( - 'username', 'name', 'is_active', 'status', 'cloud_logged_out', 'heartbeat_at', 'last_error', 'created_at', + 'username', 'name', 'is_active', 'status', 'cloud_logged_out', 'webhook', 'last_error', 'created_at', ) actions = ('publish', 'revoke') @@ -42,8 +38,21 @@ class BotAdmin(admin.ModelAdmin): return 'выключен' if bot.last_error: return f'ошибка: {bot.last_error}' - online = bot.heartbeat_at and timezone.now() - bot.heartbeat_at < timedelta(seconds=HEARTBEAT_TIMEOUT) - return 'работает' if online else 'ожидает раннер' + return 'работает' + + @admin.display(description='webhook') + def webhook(self, bot): + """Состояние webhook со стороны Bot API сервера: только на странице бота, это запрос в API.""" + if not bot.pk or not bot.is_active: + return '-' + try: + info = bot.api().get_webhook_info() + except Exception as error: + return f'не удалось получить: {error}' + lines = [info.url or 'не установлен', f'в очереди: {info.pending_update_count}'] + if info.last_error_message: + lines.append(f'последняя ошибка доставки: {info.last_error_message}') + return '\n'.join(lines) @admin.action(description='Опубликовать') def publish(self, request, queryset): @@ -58,8 +67,12 @@ class BotAdmin(admin.ModelAdmin): @admin.action(description='Отозвать') def revoke(self, request, queryset): for bot in queryset: - bot.revoke() - self.message_user(request, f'Отозвано ботов: {len(queryset)}') + try: + bot.revoke() + except Exception as error: + self.message_user(request, f'{bot} выключен, но webhook не снят: {error}', messages.WARNING) + else: + self.message_user(request, f'{bot} отозван') class ReadOnlyAdmin(admin.ModelAdmin): diff --git a/src/Telegram/api.py b/src/Telegram/api.py index 461ce3a..5dea1c3 100644 --- a/src/Telegram/api.py +++ b/src/Telegram/api.py @@ -13,6 +13,7 @@ def configure_api_server(): CLOUD_API_URL = 'https://api.telegram.org' +ALLOWED_UPDATES = ['message', 'my_chat_member'] class ApiError(Exception): diff --git a/src/Telegram/dispatcher.py b/src/Telegram/dispatcher.py index b5e7c24..cbb2657 100644 --- a/src/Telegram/dispatcher.py +++ b/src/Telegram/dispatcher.py @@ -8,7 +8,6 @@ from .models import Bot, Chat, Message, User logger = logging.getLogger(__name__) -ALLOWED_UPDATES = ['message', 'my_chat_member'] PRIVATE = 'private' LEFT_STATUSES = ('left', 'kicked') diff --git a/src/Telegram/management/__init__.py b/src/Telegram/management/__init__.py deleted file mode 100644 index e69de29..0000000 diff --git a/src/Telegram/management/commands/__init__.py b/src/Telegram/management/commands/__init__.py deleted file mode 100644 index e69de29..0000000 diff --git a/src/Telegram/management/commands/telegram_run.py b/src/Telegram/management/commands/telegram_run.py deleted file mode 100644 index c82c723..0000000 --- a/src/Telegram/management/commands/telegram_run.py +++ /dev/null @@ -1,15 +0,0 @@ -from django.core.management.base import BaseCommand - -from Telegram.runner import Supervisor - - -class Command(BaseCommand): - help = 'Запускает long polling для всех опубликованных ботов' - - def handle(self, *args, **options): - supervisor = Supervisor() - self.stdout.write('Telegram runner started. Ctrl+C to stop.') - try: - supervisor.run() - except KeyboardInterrupt: - supervisor.stop() diff --git a/src/Telegram/migrations/0003_bot_webhook.py b/src/Telegram/migrations/0003_bot_webhook.py new file mode 100644 index 0000000..c4e2bef --- /dev/null +++ b/src/Telegram/migrations/0003_bot_webhook.py @@ -0,0 +1,26 @@ +# Generated by Django 6.1.1 on 2026-09-26 13:44 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('Telegram', '0002_bot_cloud_logged_out'), + ] + + operations = [ + migrations.RemoveField( + model_name='bot', + name='heartbeat_at', + ), + migrations.RemoveField( + model_name='bot', + name='update_offset', + ), + migrations.AddField( + model_name='bot', + name='webhook_secret', + field=models.CharField(blank=True, editable=False, help_text='Telegram присылает его в X-Telegram-Bot-Api-Secret-Token, меняется при каждой публикации', max_length=64), + ), + ] diff --git a/src/Telegram/models.py b/src/Telegram/models.py index a7ad6e5..952464a 100644 --- a/src/Telegram/models.py +++ b/src/Telegram/models.py @@ -1,7 +1,16 @@ +import secrets + +from django.conf import settings from django.db import models +from django.urls import reverse from . import interfaces -from .api import log_out_from_cloud, make_api, uses_own_server +from .api import ALLOWED_UPDATES, ApiError, log_out_from_cloud, make_api, uses_own_server + + +class BotQuerySet(models.QuerySet): + def find_published(self, pk: int) -> 'Bot | None': + return self.filter(pk=pk, is_active=True).first() class Bot(models.Model): @@ -22,11 +31,15 @@ class Bot(models.Model): 'отключен от облачного API', default=False, editable=False, help_text='logOut в api.telegram.org выполнен, бот обслуживается своим Bot API сервером', ) - update_offset = models.BigIntegerField(default=0, editable=False) - heartbeat_at = models.DateTimeField(null=True, blank=True, editable=False) + webhook_secret = models.CharField( + max_length=64, blank=True, editable=False, + help_text='Telegram присылает его в X-Telegram-Bot-Api-Secret-Token, меняется при каждой публикации', + ) last_error = models.TextField(blank=True, editable=False) created_at = models.DateTimeField(auto_now_add=True) + objects = BotQuerySet.as_manager() + class Meta: verbose_name = 'бот' verbose_name_plural = 'боты' @@ -41,8 +54,14 @@ class Bot(models.Model): cls = interfaces.get(self.interface) return cls() if cls else None + def webhook_url(self) -> str: + base = getattr(settings, 'TELEGRAM_WEBHOOK_URL', '') + if not base: + raise ApiError('TELEGRAM_WEBHOOK_URL не задан') + return base.rstrip('/') + reverse('telegram:webhook', args=[self.pk]) + def publish(self): - """Проверяет токен, подтягивает имя бота и включает его обслуживание раннером. + """Проверяет токен, подтягивает имя бота и ставит webhook на этот сайт. Если настроен свой Bot API сервер, один раз отключает бота от облачного. """ @@ -53,7 +72,8 @@ class Bot(models.Model): self.save(update_fields=['cloud_logged_out']) api = self.api() me = api.get_me() - api.delete_webhook() + self.webhook_secret = secrets.token_urlsafe(32) + api.set_webhook(url=self.webhook_url(), secret_token=self.webhook_secret, allowed_updates=ALLOWED_UPDATES) self.username = me.username or '' self.name = me.first_name or '' self.is_active = True @@ -61,8 +81,10 @@ class Bot(models.Model): self.save() def revoke(self): + """Выключает бота и снимает webhook. Выключенному боту вьюха не отвечает, даже если снять webhook не вышло.""" self.is_active = False self.save(update_fields=['is_active']) + self.api().delete_webhook() class User(models.Model): diff --git a/src/Telegram/runner.py b/src/Telegram/runner.py deleted file mode 100644 index 5423819..0000000 --- a/src/Telegram/runner.py +++ /dev/null @@ -1,91 +0,0 @@ -"""Long polling раннер: по потоку на каждый опубликованный бот. - -Супервизор периодически сверяется с базой, поэтому публикация, отзыв и смена интерфейса бота в админке -применяются без перезапуска. -""" -import logging -import threading - -from django.db import close_old_connections -from django.utils import timezone - -from .dispatcher import ALLOWED_UPDATES, handle_update -from .models import Bot - -logger = logging.getLogger(__name__) - -POLL_TIMEOUT = 25 -SYNC_INTERVAL = 5 -ERROR_DELAY = 5 -HEARTBEAT_TIMEOUT = POLL_TIMEOUT * 2 - - -class BotWorker(threading.Thread): - def __init__(self, bot: Bot): - super().__init__(name=f'telegram-{bot.pk}', daemon=True) - self.bot_id = bot.pk - self.token = bot.token - self._stop_event = threading.Event() - - def stop(self): - self._stop_event.set() - - def run(self): - api = Bot(token=self.token).api() - while not self._stop_event.is_set(): - try: - self._poll(api) - except Exception as error: - logger.exception('Bot %s polling failed', self.bot_id) - Bot.objects.filter(pk=self.bot_id).update(last_error=str(error)) - self._stop_event.wait(ERROR_DELAY) - finally: - close_old_connections() - - def _poll(self, api): - offset = Bot.objects.values_list('update_offset', flat=True).get(pk=self.bot_id) - # TODO: переписать на setWebhook - updates = api.get_updates( - offset=offset, timeout=POLL_TIMEOUT, long_polling_timeout=POLL_TIMEOUT, allowed_updates=ALLOWED_UPDATES, - ) - for update in updates: - if self._stop_event.is_set(): - break - # бот перечитывается на каждое обновление, чтобы смена интерфейса применялась сразу - bot = Bot.objects.get(pk=self.bot_id) - try: - handle_update(api, bot, update) - except Exception: - logger.exception('Bot %s failed to handle update %s', self.bot_id, update.update_id) - Bot.objects.filter(pk=self.bot_id).update(update_offset=update.update_id + 1) - Bot.objects.filter(pk=self.bot_id).update(heartbeat_at=timezone.now(), last_error='') - - -class Supervisor: - def __init__(self): - self.workers: dict[int, BotWorker] = {} - self._stop_event = threading.Event() - - def stop(self): - self._stop_event.set() - - def run(self): - while not self._stop_event.is_set(): - self.sync() - close_old_connections() - self._stop_event.wait(SYNC_INTERVAL) - for worker in self.workers.values(): - worker.stop() - - def sync(self): - active = {bot.pk: bot for bot in Bot.objects.filter(is_active=True)} - for bot_id, worker in list(self.workers.items()): - bot = active.get(bot_id) - if bot is None or bot.token != worker.token or not worker.is_alive(): - worker.stop() - del self.workers[bot_id] - for bot_id, bot in active.items(): - if bot_id not in self.workers: - logger.info('Starting %s', bot) - self.workers[bot_id] = worker = BotWorker(bot) - worker.start() diff --git a/src/Telegram/testing.py b/src/Telegram/testing.py index ddd0233..a190349 100644 --- a/src/Telegram/testing.py +++ b/src/Telegram/testing.py @@ -33,7 +33,12 @@ class FakeApi: def message_update(update_id: int, text: str, chat_type: str = 'private', user_id: int = USER_ID) -> types.Update: - return types.Update.de_json({ + return types.Update.de_json(message_update_json(update_id, text, chat_type, user_id)) + + +def message_update_json(update_id: int, text: str, chat_type: str = 'private', user_id: int = USER_ID) -> dict: + """Обновление с сообщением в том виде, в каком его присылает Bot API.""" + return { 'update_id': update_id, 'message': { 'message_id': update_id, @@ -43,7 +48,7 @@ def message_update(update_id: int, text: str, chat_type: str = 'private', user_i 'title': 'group' if chat_type != 'private' else None}, 'from': {'id': user_id, 'is_bot': False, 'first_name': 'Ann', 'username': f'user{user_id}'}, }, - }) + } class Conversation: diff --git a/src/Telegram/tests.py b/src/Telegram/tests.py index 2e70b58..d23d63f 100644 --- a/src/Telegram/tests.py +++ b/src/Telegram/tests.py @@ -1,3 +1,4 @@ +import json from types import SimpleNamespace from unittest import mock @@ -8,7 +9,7 @@ from . import interfaces from .dispatcher import handle_update from .interfaces import Interface, command, register, state, text from .models import Bot, Chat, Message, User -from .testing import FakeApi, message_update +from .testing import FakeApi, message_update, message_update_json @register @@ -127,6 +128,7 @@ class DispatcherTests(TestCase): self.assertEqual(Chat.objects.count(), 1) +@override_settings(TELEGRAM_WEBHOOK_URL='http://localhost:8000/') class PublishTests(TestCase): def setUp(self): self.bot = Bot.objects.create(token='1:test') @@ -153,6 +155,74 @@ class PublishTests(TestCase): self.assertEqual(self.publish().call_count, 0) self.assertFalse(self.bot.cloud_logged_out) + def test_sets_webhook_with_new_secret(self): + self.publish() + first_secret = self.bot.webhook_secret + self.api.set_webhook.assert_called_once_with( + url=f'http://localhost:8000/telegram/webhook/{self.bot.pk}/', + secret_token=first_secret, allowed_updates=['message', 'my_chat_member'], + ) + self.publish() + self.assertNotEqual(self.bot.webhook_secret, first_secret) + + @override_settings(TELEGRAM_WEBHOOK_URL='') + def test_requires_webhook_url(self): + with self.assertRaisesMessage(Exception, 'TELEGRAM_WEBHOOK_URL'): + self.publish() + self.bot.refresh_from_db() + self.assertFalse(self.bot.is_active) + + def test_revoke_deletes_webhook(self): + self.publish() + with mock.patch.object(Bot, 'api', return_value=self.api): + self.bot.revoke() + self.api.delete_webhook.assert_called_once() + self.bot.refresh_from_db() + self.assertFalse(self.bot.is_active) + + +class WebhookTests(TestCase): + def setUp(self): + self.bot = Bot.objects.create( + token='1:test', interface=interfaces.key(EchoInterface), is_active=True, webhook_secret='secret', + ) + self.api = FakeApi() + patcher = mock.patch.object(Bot, 'api', return_value=self.api) + patcher.start() + self.addCleanup(patcher.stop) + + def post(self, text='/start', secret='secret', bot_id=None, body=None): + if body is None: + body = json.dumps(message_update_json(1, text)) + return self.client.post( + f'/telegram/webhook/{bot_id or self.bot.pk}/', body, content_type='application/json', + headers={'X-Telegram-Bot-Api-Secret-Token': secret}, + ) + + def test_handles_update(self): + self.assertEqual(self.post().status_code, 200) + self.assertEqual(self.api.last_text, 'menu') + + def test_rejects_wrong_secret_and_inactive_bot(self): + self.assertEqual(self.post(secret='wrong').status_code, 403) + self.assertEqual(self.post(bot_id=self.bot.pk + 1).status_code, 403) + Bot.objects.filter(pk=self.bot.pk).update(is_active=False) + self.assertEqual(self.post().status_code, 403) + self.assertEqual(self.api.sent, []) + + def test_bad_body(self): + self.assertEqual(self.post(body='not json').status_code, 400) + + def test_handler_error_is_saved_and_acknowledged(self): + with mock.patch('Telegram.views.handle_update', side_effect=RuntimeError('boom')), \ + self.assertLogs('Telegram.views', 'ERROR'): + self.assertEqual(self.post().status_code, 200) + self.bot.refresh_from_db() + self.assertEqual(self.bot.last_error, 'boom') + self.post() + self.bot.refresh_from_db() + self.assertEqual(self.bot.last_error, '') + class InterfaceRegistryTests(TestCase): def test_registered_interfaces_are_choices(self): diff --git a/src/Telegram/urls.py b/src/Telegram/urls.py new file mode 100644 index 0000000..e021981 --- /dev/null +++ b/src/Telegram/urls.py @@ -0,0 +1,9 @@ +from django.urls import path + +from . import views + +app_name = 'telegram' + +urlpatterns = [ + path('webhook//', views.webhook, name='webhook'), +] diff --git a/src/Telegram/views.py b/src/Telegram/views.py new file mode 100644 index 0000000..069ac8a --- /dev/null +++ b/src/Telegram/views.py @@ -0,0 +1,41 @@ +import logging + +from django.http import HttpResponse, HttpResponseBadRequest, HttpResponseForbidden +from django.utils.crypto import constant_time_compare +from django.views.decorators.csrf import csrf_exempt +from django.views.decorators.http import require_POST +from telebot import types + +from .dispatcher import handle_update +from .models import Bot + +logger = logging.getLogger(__name__) + +SECRET_HEADER = 'X-Telegram-Bot-Api-Secret-Token' + + +@csrf_exempt +@require_POST +def webhook(request, bot_id: int): + """Принимает обновление от Bot API сервера. + + Бот читается из базы на каждое обновление, поэтому смена интерфейса в админке применяется сразу. + Ошибка обработчика не возвращается телеграму: иначе он будет повторять то же обновление. + """ + bot = Bot.objects.find_published(bot_id) + if bot is None or not constant_time_compare(request.headers.get(SECRET_HEADER, ''), bot.webhook_secret): + return HttpResponseForbidden() + try: + update = types.Update.de_json(request.body.decode()) + except ValueError: + return HttpResponseBadRequest() + + try: + handle_update(bot.api(), bot, update) + except Exception as error: + logger.exception('Bot %s failed to handle update %s', bot_id, update.update_id) + Bot.objects.filter(pk=bot_id).update(last_error=str(error)) + else: + if bot.last_error: + Bot.objects.filter(pk=bot_id).update(last_error='') + return HttpResponse() diff --git a/src/serverus/xui_api.py b/src/serverus/xui_api.py index dea6d27..f3dbb21 100644 --- a/src/serverus/xui_api.py +++ b/src/serverus/xui_api.py @@ -70,6 +70,10 @@ class XuiApi: def list_clients(self) -> list[dict]: return self._call('GET', 'clients/list') or [] + def status(self) -> dict: + """Состояние сервера панели: cpu, mem, xray (state, version), uptime и остальное.""" + return self._call('GET', 'server/status') or {} + def _call(self, method: str, path: str, body: dict | None = None): data = json.dumps(body).encode() if body is not None else None request = urllib.request.Request(f'{self.url}/{path}', data=data, method=method, headers={