chore: переезд на serWebhook

This commit is contained in:
protokey committed 2026-09-26 18:20:35 +04:00
1 parent 45e7a1d22f
commit f1414e5ace
42 files changed
+1391 -224

No files matched your search

+10
View File
@@ -0,0 +1,10 @@
.git
.idea
.venv
.env
extra
docs
**/__pycache__
src/db.sqlite3
src/design/static_root
docker-services-data
+44
View File
@@ -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", "-"]
+30 -4
View File
@@ -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 <pkg> [<pkg> ...] - установить пакет(ы) в $(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=<dump> - восстановить базу из бэкапа (сайт и воркер остановлены)"
@echo "make docker_build / docker_up - собрать образ и поднять все в docker (профиль app)"
@echo "make serverus_xui ip=<ip> password=<root_password> [domain=<domain>] [ssl=no] [ca=yes] - настроить свежий сервер и поставить 3x-ui, доступы в $(SERVERUS_ACCESS_DIR)/<ip>/"
@echo "make serverus_key ip=<ip> action=add|del|stats name=<name> - ключ на панели, настроенной через serverus_xui"
@echo "make serverus_cert ip=<ip> [domain=<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=<ip> [cmd='uptime']
serve_ssh:
@test -n "$(ip)" || { echo "Укажите сервер: make serve_ssh ip=<ip> [cmd='<команда>']"; exit 1; }
+20 -29
View File
@@ -299,10 +299,10 @@ https://1.2.3.4:12839/9g72PmSGAQ3V/<subId>
### Запуск
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/<subId>
Сервис - На серверах у нас живут сервисы. Поднимается или настраивается какой-то сервис, чтобы решить какую-то полезную задачу. И реализуют 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; статистика - только у ключей; ботом может пользоваться любой.
+57 -4
View File
@@ -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:
+14 -1
View File
@@ -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=
+2 -1
View File
@@ -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
+4
View File
@@ -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=<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` - кнопки и тексты.
+28 -2
View File
@@ -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
+57
View File
@@ -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:]
+4 -7
View File
@@ -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)
+17 -4
View File
@@ -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'Не удалось применить настройку <code>{escape(point.name)}</code>:\n{escape(str(error))}\n'
'Точка работает с прежними настройками.'
)
def task_interrupted(point: Point) -> str:
return f'<code>{escape(point.name)}</code>: {escape(point.error)}.'
def health_problem(point: Point, problem: str) -> str:
return f'Проверка <code>{escape(point.name)}</code>: {escape(problem)}.'
def health_restored(point: Point) -> str:
return f'Проверка <code>{escape(point.name)}</code>: снова все в порядке.'
+47
View File
@@ -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 в бэкап не входит: храните его отдельно.')
@@ -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)}')
@@ -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 отдельно от бэкапов базы: без него секреты не расшифровать.')
+35
View File
@@ -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.')
@@ -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),
),
]
+68 -1
View File
@@ -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:
+37
View File
@@ -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':
+157 -30
View File
@@ -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()
+186 -13
View File
@@ -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()
+136
View File
@@ -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()
+33
View File
@@ -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)
+1
View File
@@ -12,3 +12,4 @@ from .project import *
```
- `test.py` - настройки для тестов: фоновые задачи идут в stub брокер, а не в rabbitmq. `make dj_test` использует их.
- `docker.py` - настройки для docker: все из окружения, `DEBUG=False`, `SECRET_KEY` обязателен, статика через whitenoise, логи в консоль.
+47
View File
@@ -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')},
}
+8
View File
@@ -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=''),
+2 -1
View File
@@ -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')),
]
+6 -3
View File
@@ -19,7 +19,7 @@
- `context.py` - `Context`, который получает обработчик: текст, команда, собеседник, состояние, `reply()`.
- `messages.py` - `send()` для отправки сообщения в чат вне обработчика, например уведомления из фоновой задачи.
- `dispatcher.py` - разбор обновления: сохраняет пользователя, чат, сообщение и передает его интерфейсу бота.
- `runner.py` - long polling: по потоку на каждый опубликованный бот, супервизор сверяется с базой каждые 5 секунд.
- `views.py` - webhook `/telegram/webhook/<id бота>/`: проверяет секрет из заголовка `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, чтобы он не повторял то же обновление.
+22 -9
View File
@@ -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):
+1
View File
@@ -13,6 +13,7 @@ def configure_api_server():
CLOUD_API_URL = 'https://api.telegram.org'
ALLOWED_UPDATES = ['message', 'my_chat_member']
class ApiError(Exception):
-1
View File
@@ -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')
View File
Whitespace-only changes.
Whitespace-only changes.
@@ -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()
@@ -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),
),
]
+27 -5
View File
@@ -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):
-91
View File
@@ -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()
+7 -2
View File
@@ -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:
+71 -1
View File
@@ -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):
+9
View File
@@ -0,0 +1,9 @@
from django.urls import path
from . import views
app_name = 'telegram'
urlpatterns = [
path('webhook/<int:bot_id>/', views.webhook, name='webhook'),
]
+41
View File
@@ -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()
+4
View File
@@ -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={