chore: база
This commit is contained in:
1 parent
93790da617
commit
5f53cadf89
39 files changed
+1457
-18
No files matched your search
@@ -3,11 +3,12 @@ PYTHON := $(VENV)/bin/python
|
||||
PIP := $(VENV)/bin/pip
|
||||
WHEELS_DIR := wheels
|
||||
REQUIREMENTS := requirements.txt
|
||||
MANAGE := cd src/ && ../$(PYTHON) manage.py
|
||||
|
||||
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
|
||||
.PHONY: install freeze wheels sync sync-offline clean help venv dj_makemigrations dj_migrate dj_superuser dj_run dj_test dj_startapp tg_run
|
||||
|
||||
help:
|
||||
@echo "make install <pkg> [<pkg> ...] - установить пакет(ы) в $(VENV), обновить $(REQUIREMENTS), собрать wheel в $(WHEELS_DIR)/"
|
||||
@@ -52,19 +53,22 @@ clean:
|
||||
@:
|
||||
|
||||
dj_makemigrations:
|
||||
cd src/ && ./manage.py makemigrations
|
||||
$(MANAGE) makemigrations
|
||||
|
||||
dj_migrate:
|
||||
cd src/ && ./manage.py migrate
|
||||
$(MANAGE) migrate
|
||||
|
||||
dj_superuser:
|
||||
cd src/ && ./manage.py createsuperuser
|
||||
$(MANAGE) createsuperuser
|
||||
|
||||
dj_run:
|
||||
cd src/ && ./manage.py runserver
|
||||
$(MANAGE) runserver
|
||||
|
||||
dj_test:
|
||||
cd src/ && ./manage.py test
|
||||
$(MANAGE) test
|
||||
|
||||
dj_startapp:
|
||||
cd src/ && ./manage.py startapp $(app)
|
||||
$(MANAGE) startapp $(app)
|
||||
|
||||
tg_run:
|
||||
$(MANAGE) telegram_run
|
||||
@@ -1,6 +1,8 @@
|
||||
# Telegram Bot API server (https://my.telegram.org/apps)
|
||||
TG_APP_ID=
|
||||
TG_APP_HASH=
|
||||
# Адрес своего Bot API сервера для Django. Пусто - api.telegram.org
|
||||
TELEGRAM_API_URL=http://127.0.0.1:8081
|
||||
|
||||
# PostgreSQL
|
||||
DB_NAME=
|
||||
|
||||
@@ -1,21 +1,35 @@
|
||||
aiohappyeyeballs==2.7.1
|
||||
aiohttp==3.14.3
|
||||
aiosignal==1.4.0
|
||||
ansible-core==2.21.4
|
||||
ansible-runner==2.4.3
|
||||
asgiref==3.12.1
|
||||
attrs==26.1.0
|
||||
certifi==2026.7.22
|
||||
cffi==2.1.1
|
||||
charset-normalizer==3.5.1
|
||||
cryptography==50.0.1
|
||||
Django==6.1.1
|
||||
dramatiq==2.2.1
|
||||
frozenlist==1.8.0
|
||||
idna==3.20
|
||||
Jinja2==3.1.6
|
||||
lockfile==0.12.2
|
||||
MarkupSafe==3.0.3
|
||||
multidict==6.9.1
|
||||
packaging==26.3
|
||||
pexpect==4.9.0
|
||||
pika==1.4.4
|
||||
propcache==0.5.4
|
||||
psycopg2-binary==2.9.13
|
||||
ptyprocess==0.7.0
|
||||
pycparser==3.0
|
||||
pyTelegramBotAPI==4.37.0
|
||||
python-daemon==3.1.2
|
||||
PyYAML==6.0.3
|
||||
redis==8.1.0
|
||||
requests==2.34.2
|
||||
resolvelib==1.2.1
|
||||
sqlparse==0.6.0
|
||||
typing_extensions==4.16.0
|
||||
urllib3==2.8.0
|
||||
|
||||
@@ -3,3 +3,6 @@ from django.apps import AppConfig
|
||||
|
||||
class AppConfig(AppConfig):
|
||||
name = 'App'
|
||||
|
||||
def ready(self):
|
||||
from . import bot # noqa: F401 регистрирует интерфейс бота
|
||||
@@ -0,0 +1,39 @@
|
||||
from Telegram.interfaces import Interface, command, register, text
|
||||
|
||||
MAIN_MENU = 'главное меню'
|
||||
CANCEL = 'отмена'
|
||||
FAQ = 'ФАК'
|
||||
ADD_SERVER = 'добавить сервер'
|
||||
HELP = 'помощь'
|
||||
|
||||
MAIN_KEYBOARD = [[FAQ, ADD_SERVER, HELP]]
|
||||
|
||||
|
||||
@register
|
||||
class MyPointInterface(Interface):
|
||||
title = 'MyPoint VPN'
|
||||
|
||||
@command('start')
|
||||
@text(MAIN_MENU, CANCEL)
|
||||
def main_menu(self, ctx):
|
||||
ctx.reset()
|
||||
# TODO: перечислить серверы пользователя и добавить кнопку "управление сервером", когда появится модель данных
|
||||
ctx.reply('привет. выбери действие', keyboard=MAIN_KEYBOARD)
|
||||
|
||||
@text(FAQ)
|
||||
def faq(self, ctx):
|
||||
# TODO: текст FAQ
|
||||
ctx.reply('FAQ пока пуст')
|
||||
|
||||
@text(HELP)
|
||||
def help(self, ctx):
|
||||
# TODO: текст помощи
|
||||
ctx.reply('Помощь пока пуста')
|
||||
|
||||
@text(ADD_SERVER)
|
||||
def add_server(self, ctx):
|
||||
# TODO: меню добавления сервера, когда появится модель данных App
|
||||
ctx.reply('Добавление сервера пока не реализовано')
|
||||
|
||||
def fallback(self, ctx):
|
||||
ctx.reply('Не понял. Воспользуйся кнопками меню.', keyboard=MAIN_KEYBOARD)
|
||||
@@ -0,0 +1,18 @@
|
||||
from django.test import TestCase
|
||||
|
||||
from Telegram import interfaces
|
||||
from Telegram.dispatcher import handle_update
|
||||
from Telegram.models import Bot
|
||||
from Telegram.tests import FakeApi, message_update
|
||||
|
||||
from .bot import MAIN_KEYBOARD, MyPointInterface
|
||||
|
||||
|
||||
class MyPointInterfaceTests(TestCase):
|
||||
def test_start_and_main_menu_button_are_same(self):
|
||||
api = FakeApi()
|
||||
bot = Bot.objects.create(token='1:test', interface=interfaces.key(MyPointInterface))
|
||||
handle_update(api, bot, message_update(1, '/start'))
|
||||
handle_update(api, bot, message_update(2, 'главное меню'))
|
||||
self.assertEqual(api.sent[0][1], api.sent[1][1])
|
||||
self.assertEqual([button['text'] for button in api.sent[0][2].keyboard[0]], MAIN_KEYBOARD[0])
|
||||
@@ -1,3 +1,5 @@
|
||||
import os
|
||||
|
||||
from .base import *
|
||||
|
||||
DEBUG = True
|
||||
@@ -13,3 +15,10 @@ MEDIA_ROOT = BASE_DIR.parent / 'extra/media'
|
||||
MEDIA_URL = '/media/'
|
||||
|
||||
AUTH_USER_MODEL = 'App.User'
|
||||
|
||||
INSTALLED_APPS += [
|
||||
'Telegram.apps.TelegramConfig',
|
||||
]
|
||||
|
||||
# Свой Bot API сервер (сервис telegram в docker-compose.yaml). Пусто - api.telegram.org
|
||||
TELEGRAM_API_URL = os.environ.get('TELEGRAM_API_URL', '')
|
||||
@@ -4,6 +4,7 @@
|
||||
|
||||
- `design` - фронт
|
||||
- `App` - главное приложение
|
||||
- `Telegram` - Django приложение для чат-ботов
|
||||
- `MyPointVPN` - настройки django
|
||||
- `manage.py` - точка входа
|
||||
|
||||
|
||||
@@ -9,3 +9,55 @@
|
||||
Смена интерфейса бота в админке переключает поведение бота в рантайме.
|
||||
|
||||
В первой итерации нужно реализовать поддержку создания интерфейсов чат-ботов. Которые сообщениями отвечают на сообщения. И поддержку приватных сообщений только. Но мы подразумеваем что в будущем это приложение может быть расширено.
|
||||
|
||||
## Как устроено
|
||||
|
||||
Зависимости: Django и pyTelegramBotAPI. Приложение ничего не знает о других приложениях проекта: оно предоставляет базовый класс интерфейса, а приложения регистрируют свои интерфейсы сами.
|
||||
|
||||
- `models.py` - `Bot`, `User`, `Chat` (у каждого чата состояние диалога `state` и данные `data`), `Message`.
|
||||
- `interfaces.py` - базовый класс `Interface`, декораторы `command`, `text`, `state` и реестр `register`.
|
||||
- `context.py` - `Context`, который получает обработчик: текст, команда, собеседник, состояние, `reply()`.
|
||||
- `messages.py` - `send()` для отправки сообщения в чат вне обработчика, например уведомления из фоновой задачи.
|
||||
- `dispatcher.py` - разбор обновления: сохраняет пользователя, чат, сообщение и передает его интерфейсу бота.
|
||||
- `runner.py` - long polling: по потоку на каждый опубликованный бот, супервизор сверяется с базой каждые 5 секунд.
|
||||
|
||||
## Интерфейс бота
|
||||
|
||||
```python
|
||||
from Telegram.interfaces import Interface, command, register, state, text
|
||||
|
||||
|
||||
@register
|
||||
class Shop(Interface):
|
||||
title = 'Магазин'
|
||||
|
||||
@command('start')
|
||||
@text('главное меню')
|
||||
def menu(self, ctx):
|
||||
ctx.reset()
|
||||
ctx.reply('Привет', keyboard=[['адрес']])
|
||||
|
||||
@text('адрес')
|
||||
def ask_address(self, ctx):
|
||||
ctx.state = 'await_address'
|
||||
ctx.reply('Пришли адрес', keyboard=[['главное меню']])
|
||||
|
||||
@state('await_address')
|
||||
def address(self, ctx):
|
||||
ctx.data['address'] = ctx.text
|
||||
ctx.state = ''
|
||||
ctx.reply('Сохранил')
|
||||
|
||||
def fallback(self, ctx):
|
||||
ctx.reply('Не понял')
|
||||
```
|
||||
|
||||
Порядок разбора сообщения: команда, текст кнопки, состояние диалога, `fallback`. Модуль с интерфейсом нужно импортировать при старте, например в `AppConfig.ready()`.
|
||||
|
||||
## Запуск
|
||||
|
||||
1. В админке добавить бота с токеном и выбрать интерфейс.
|
||||
2. Действие "Опубликовать" проверяет токен, подтягивает имя бота и включает бота. "Отозвать" выключает.
|
||||
3. `make tg_run` запускает раннер. Смена интерфейса, публикация и отзыв применяются без перезапуска.
|
||||
|
||||
`TELEGRAM_API_URL` направляет бота на свой Bot API сервер (сервис `telegram` в `docker-compose.yaml`).
|
||||
+86
-2
@@ -1,3 +1,87 @@
|
||||
from django.contrib import admin
|
||||
from datetime import timedelta
|
||||
|
||||
# Register your models here.
|
||||
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):
|
||||
interface = forms.ChoiceField(required=False)
|
||||
|
||||
class Meta:
|
||||
model = Bot
|
||||
fields = ['token', 'interface']
|
||||
widgets = {'token': forms.PasswordInput(render_value=True)}
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.fields['interface'].choices = [('', '---------')] + interfaces.choices()
|
||||
|
||||
|
||||
@admin.register(Bot)
|
||||
class BotAdmin(admin.ModelAdmin):
|
||||
form = BotForm
|
||||
list_display = ('__str__', 'name', 'interface_title', 'is_active', 'status')
|
||||
readonly_fields = ('username', 'name', 'is_active', 'status', 'heartbeat_at', 'last_error', 'created_at')
|
||||
actions = ('publish', 'revoke')
|
||||
|
||||
@admin.display(description='интерфейс')
|
||||
def interface_title(self, bot):
|
||||
cls = interfaces.get(bot.interface)
|
||||
return (cls.title or bot.interface) if cls else '-'
|
||||
|
||||
@admin.display(description='состояние')
|
||||
def status(self, bot):
|
||||
if not bot.is_active:
|
||||
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 'ожидает раннер'
|
||||
|
||||
@admin.action(description='Опубликовать')
|
||||
def publish(self, request, queryset):
|
||||
for bot in queryset:
|
||||
try:
|
||||
bot.publish()
|
||||
except Exception as error:
|
||||
self.message_user(request, f'{bot}: {error}', messages.ERROR)
|
||||
else:
|
||||
self.message_user(request, f'{bot} опубликован')
|
||||
|
||||
@admin.action(description='Отозвать')
|
||||
def revoke(self, request, queryset):
|
||||
for bot in queryset:
|
||||
bot.revoke()
|
||||
self.message_user(request, f'Отозвано ботов: {len(queryset)}')
|
||||
|
||||
|
||||
class ReadOnlyAdmin(admin.ModelAdmin):
|
||||
def has_add_permission(self, request):
|
||||
return False
|
||||
|
||||
def has_change_permission(self, request, obj=None):
|
||||
return False
|
||||
|
||||
|
||||
@admin.register(User)
|
||||
class UserAdmin(ReadOnlyAdmin):
|
||||
list_display = ('tg_id', 'username', 'first_name', 'last_name', 'created_at')
|
||||
search_fields = ('tg_id', 'username', 'first_name', 'last_name')
|
||||
|
||||
|
||||
@admin.register(Chat)
|
||||
class ChatAdmin(ReadOnlyAdmin):
|
||||
list_display = ('__str__', 'bot', 'type', 'state', 'created_at', 'deleted_at')
|
||||
list_filter = ('bot', 'type')
|
||||
|
||||
|
||||
@admin.register(Message)
|
||||
class MessageAdmin(ReadOnlyAdmin):
|
||||
list_display = ('chat', 'direction', 'text', 'created_at')
|
||||
list_filter = ('direction', 'chat__bot')
|
||||
search_fields = ('text',)
|
||||
@@ -0,0 +1,15 @@
|
||||
import telebot
|
||||
from django.conf import settings
|
||||
from telebot import apihelper
|
||||
|
||||
|
||||
def configure_api_server():
|
||||
"""Направляет pyTelegramBotAPI на свой инстанс Bot API сервера, если он задан в TELEGRAM_API_URL."""
|
||||
url = getattr(settings, 'TELEGRAM_API_URL', '').rstrip('/')
|
||||
if url:
|
||||
apihelper.API_URL = url + '/bot{0}/{1}'
|
||||
apihelper.FILE_URL = url + '/file/bot{0}/{1}'
|
||||
|
||||
|
||||
def make_api(token: str) -> telebot.TeleBot:
|
||||
return telebot.TeleBot(token, threaded=False)
|
||||
@@ -3,3 +3,9 @@ from django.apps import AppConfig
|
||||
|
||||
class TelegramConfig(AppConfig):
|
||||
name = 'Telegram'
|
||||
verbose_name = 'Telegram'
|
||||
default_auto_field = 'django.db.models.BigAutoField'
|
||||
|
||||
def ready(self):
|
||||
from .api import configure_api_server
|
||||
configure_api_server()
|
||||
@@ -0,0 +1,61 @@
|
||||
from telebot import types
|
||||
|
||||
from . import interfaces, messages
|
||||
from .models import Bot, Chat, User
|
||||
|
||||
|
||||
class Context:
|
||||
"""Всё, что нужно обработчику интерфейса: входящее сообщение, собеседник, состояние диалога и ответ."""
|
||||
|
||||
def __init__(self, api, bot: Bot, chat: Chat, user: User, message: types.Message):
|
||||
self.api = api
|
||||
self.bot = bot
|
||||
self.chat = chat
|
||||
self.user = user
|
||||
self.message = message
|
||||
|
||||
@property
|
||||
def text(self) -> str:
|
||||
return self.message.text or ''
|
||||
|
||||
@property
|
||||
def command(self) -> str | None:
|
||||
"""'start' для '/start arg' и '/start@my_bot'."""
|
||||
if not self.text.startswith('/'):
|
||||
return None
|
||||
return self.text.split()[0][1:].split('@')[0].lower()
|
||||
|
||||
@property
|
||||
def args(self) -> str:
|
||||
"""Текст после команды."""
|
||||
parts = self.text.split(maxsplit=1)
|
||||
return parts[1] if self.command and len(parts) > 1 else ''
|
||||
|
||||
@property
|
||||
def state(self) -> str:
|
||||
return self.chat.state
|
||||
|
||||
@state.setter
|
||||
def state(self, value: str):
|
||||
self.chat.state = value or ''
|
||||
|
||||
@property
|
||||
def data(self) -> dict:
|
||||
return self.chat.data
|
||||
|
||||
def reset(self):
|
||||
"""Сбрасывает состояние и данные диалога."""
|
||||
self.chat.state = ''
|
||||
self.chat.data = {}
|
||||
|
||||
def routes(self) -> list[tuple[str, str]]:
|
||||
routes = []
|
||||
if self.command:
|
||||
routes.append((interfaces.COMMAND, self.command))
|
||||
routes.append((interfaces.TEXT, self.text))
|
||||
if self.state:
|
||||
routes.append((interfaces.STATE, self.state))
|
||||
return routes
|
||||
|
||||
def reply(self, text: str, keyboard: messages.Keyboard | None = None, parse_mode: str | None = None):
|
||||
return messages.send(self.chat, text, keyboard=keyboard, parse_mode=parse_mode, api=self.api)
|
||||
@@ -0,0 +1,66 @@
|
||||
import logging
|
||||
|
||||
from django.utils import timezone
|
||||
from telebot import types
|
||||
|
||||
from .context import Context
|
||||
from .models import Bot, Chat, Message, User
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
ALLOWED_UPDATES = ['message', 'my_chat_member']
|
||||
PRIVATE = 'private'
|
||||
LEFT_STATUSES = ('left', 'kicked')
|
||||
|
||||
|
||||
def handle_update(api, bot: Bot, update: types.Update):
|
||||
if update.my_chat_member:
|
||||
_handle_membership(bot, update.my_chat_member)
|
||||
elif update.message:
|
||||
_handle_message(api, bot, update.message)
|
||||
|
||||
|
||||
def _handle_membership(bot: Bot, member: types.ChatMemberUpdated):
|
||||
chat = _save_chat(bot, member.chat)
|
||||
if member.new_chat_member.status in LEFT_STATUSES:
|
||||
chat.deleted_at = timezone.now()
|
||||
chat.save(update_fields=['deleted_at'])
|
||||
|
||||
|
||||
def _handle_message(api, bot: Bot, message: types.Message):
|
||||
# первая итерация: только личные сообщения
|
||||
if message.chat.type != PRIVATE or not message.from_user:
|
||||
return
|
||||
user = _save_user(message.from_user)
|
||||
chat = _save_chat(bot, message.chat, user)
|
||||
Message.objects.create(
|
||||
chat=chat, user=user, direction=Message.Direction.IN, message_id=message.message_id, text=message.text or '',
|
||||
)
|
||||
|
||||
interface = bot.get_interface()
|
||||
if interface is None:
|
||||
logger.warning('%s has no interface, message ignored', bot)
|
||||
return
|
||||
interface.dispatch(Context(api, bot, chat, user, message))
|
||||
chat.save(update_fields=['state', 'data'])
|
||||
|
||||
|
||||
def _save_user(tg_user: types.User) -> User:
|
||||
user, _ = User.objects.update_or_create(
|
||||
tg_id=tg_user.id,
|
||||
defaults={
|
||||
'username': tg_user.username or '',
|
||||
'first_name': tg_user.first_name or '',
|
||||
'last_name': tg_user.last_name or '',
|
||||
'language_code': tg_user.language_code or '',
|
||||
},
|
||||
)
|
||||
return user
|
||||
|
||||
|
||||
def _save_chat(bot: Bot, tg_chat: types.Chat, user: User | None = None) -> Chat:
|
||||
defaults = {'type': tg_chat.type, 'title': tg_chat.title or '', 'deleted_at': None}
|
||||
if user:
|
||||
defaults['user'] = user
|
||||
chat, _ = Chat.objects.update_or_create(bot=bot, chat_id=tg_chat.id, defaults=defaults)
|
||||
return chat
|
||||
@@ -0,0 +1,94 @@
|
||||
"""Интерфейсы ботов.
|
||||
|
||||
Интерфейс - класс, определяющий поведение бота. Приложения наследуют Interface, размечают методы декораторами
|
||||
command/text/state и регистрируют класс через register. В админке боту выбирается один из зарегистрированных интерфейсов.
|
||||
|
||||
@register
|
||||
class Shop(Interface):
|
||||
title = 'Магазин'
|
||||
|
||||
@command('start')
|
||||
@text('главное меню')
|
||||
def menu(self, ctx):
|
||||
ctx.reply('Привет', keyboard=[['каталог']])
|
||||
|
||||
@state('await_address')
|
||||
def address(self, ctx):
|
||||
ctx.data['address'] = ctx.text
|
||||
ctx.state = ''
|
||||
|
||||
Порядок разбора сообщения: команда, текст кнопки, состояние диалога, fallback.
|
||||
"""
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from .context import Context
|
||||
|
||||
_registry: dict[str, type['Interface']] = {}
|
||||
|
||||
COMMAND = 'command'
|
||||
TEXT = 'text'
|
||||
STATE = 'state'
|
||||
|
||||
|
||||
def _route(kind: str, values: tuple[str, ...]):
|
||||
def decorator(func):
|
||||
func.__dict__.setdefault('_telegram_routes', []).extend((kind, value) for value in values)
|
||||
return func
|
||||
return decorator
|
||||
|
||||
|
||||
def command(*names: str):
|
||||
"""Обрабатывает команды /name. Имена без слеша, регистр не важен."""
|
||||
return _route(COMMAND, tuple(name.lstrip('/').lower() for name in names))
|
||||
|
||||
|
||||
def text(*texts: str):
|
||||
"""Обрабатывает сообщение с точно таким текстом, например нажатие кнопки клавиатуры."""
|
||||
return _route(TEXT, texts)
|
||||
|
||||
|
||||
def state(*states: str):
|
||||
"""Обрабатывает любое сообщение, пока чат находится в этом состоянии."""
|
||||
return _route(STATE, states)
|
||||
|
||||
|
||||
class Interface:
|
||||
title = ''
|
||||
_routes: dict[tuple[str, str], str] = {}
|
||||
|
||||
def __init_subclass__(cls, **kwargs):
|
||||
super().__init_subclass__(**kwargs)
|
||||
routes = {}
|
||||
for klass in reversed(cls.__mro__):
|
||||
for name, attr in vars(klass).items():
|
||||
for route in getattr(attr, '_telegram_routes', ()):
|
||||
routes[route] = name
|
||||
cls._routes = routes
|
||||
|
||||
def dispatch(self, ctx: 'Context'):
|
||||
for route in ctx.routes():
|
||||
name = self._routes.get(route)
|
||||
if name:
|
||||
return getattr(self, name)(ctx)
|
||||
return self.fallback(ctx)
|
||||
|
||||
def fallback(self, ctx: 'Context'):
|
||||
"""Вызывается, если сообщение не подошло ни под один обработчик."""
|
||||
|
||||
|
||||
def key(cls: type[Interface]) -> str:
|
||||
return f'{cls.__module__}.{cls.__qualname__}'
|
||||
|
||||
|
||||
def register(cls: type[Interface]) -> type[Interface]:
|
||||
_registry[key(cls)] = cls
|
||||
return cls
|
||||
|
||||
|
||||
def get(interface_key: str) -> type[Interface] | None:
|
||||
return _registry.get(interface_key)
|
||||
|
||||
|
||||
def choices() -> list[tuple[str, str]]:
|
||||
return [(k, cls.title or k) for k, cls in sorted(_registry.items())]
|
||||
File renamed without changes.
Whitespace-only changes.
@@ -0,0 +1,15 @@
|
||||
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,24 @@
|
||||
from telebot import types
|
||||
|
||||
from .models import Chat, Message
|
||||
|
||||
Keyboard = list[list[str]]
|
||||
|
||||
|
||||
def build_keyboard(keyboard: Keyboard | None):
|
||||
"""None - не менять клавиатуру у пользователя, [] - убрать клавиатуру."""
|
||||
if keyboard is None:
|
||||
return None
|
||||
if not keyboard:
|
||||
return types.ReplyKeyboardRemove()
|
||||
markup = types.ReplyKeyboardMarkup(resize_keyboard=True)
|
||||
for row in keyboard:
|
||||
markup.row(*row)
|
||||
return markup
|
||||
|
||||
|
||||
def send(chat: Chat, text: str, keyboard: Keyboard | None = None, parse_mode: str | None = None, api=None) -> Message:
|
||||
"""Отправляет сообщение в чат от имени его бота и сохраняет его в историю."""
|
||||
api = api or chat.bot.api()
|
||||
sent = api.send_message(chat.chat_id, text, reply_markup=build_keyboard(keyboard), parse_mode=parse_mode)
|
||||
return Message.objects.create(chat=chat, direction=Message.Direction.OUT, message_id=sent.message_id, text=text)
|
||||
@@ -0,0 +1,91 @@
|
||||
# Generated by Django 6.1.1 on 2026-09-24 12:00
|
||||
|
||||
import django.db.models.deletion
|
||||
from django.db import migrations, models
|
||||
|
||||
|
||||
class Migration(migrations.Migration):
|
||||
|
||||
initial = True
|
||||
|
||||
dependencies = [
|
||||
]
|
||||
|
||||
operations = [
|
||||
migrations.CreateModel(
|
||||
name='Bot',
|
||||
fields=[
|
||||
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
|
||||
('token', models.CharField(max_length=128, unique=True)),
|
||||
('username', models.CharField(blank=True, editable=False, max_length=64)),
|
||||
('name', models.CharField(blank=True, editable=False, max_length=128)),
|
||||
('interface', models.CharField(blank=True, help_text='Класс, определяющий поведение бота. Меняется на лету.', max_length=255)),
|
||||
('is_active', models.BooleanField(default=False, editable=False, verbose_name='опубликован')),
|
||||
('update_offset', models.BigIntegerField(default=0, editable=False)),
|
||||
('heartbeat_at', models.DateTimeField(blank=True, editable=False, null=True)),
|
||||
('last_error', models.TextField(blank=True, editable=False)),
|
||||
('created_at', models.DateTimeField(auto_now_add=True)),
|
||||
],
|
||||
options={
|
||||
'verbose_name': 'бот',
|
||||
'verbose_name_plural': 'боты',
|
||||
},
|
||||
),
|
||||
migrations.CreateModel(
|
||||
name='User',
|
||||
fields=[
|
||||
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
|
||||
('tg_id', models.BigIntegerField(unique=True)),
|
||||
('username', models.CharField(blank=True, max_length=64)),
|
||||
('first_name', models.CharField(blank=True, max_length=128)),
|
||||
('last_name', models.CharField(blank=True, max_length=128)),
|
||||
('language_code', models.CharField(blank=True, max_length=16)),
|
||||
('created_at', models.DateTimeField(auto_now_add=True)),
|
||||
('updated_at', models.DateTimeField(auto_now=True)),
|
||||
],
|
||||
options={
|
||||
'verbose_name': 'пользователь telegram',
|
||||
'verbose_name_plural': 'пользователи telegram',
|
||||
},
|
||||
),
|
||||
migrations.CreateModel(
|
||||
name='Chat',
|
||||
fields=[
|
||||
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
|
||||
('chat_id', models.BigIntegerField()),
|
||||
('type', models.CharField(max_length=16)),
|
||||
('title', models.CharField(blank=True, max_length=255)),
|
||||
('state', models.CharField(blank=True, help_text='Состояние диалога, задается интерфейсом бота', max_length=128)),
|
||||
('data', models.JSONField(blank=True, default=dict, help_text='Данные диалога, задаются интерфейсом бота')),
|
||||
('created_at', models.DateTimeField(auto_now_add=True)),
|
||||
('deleted_at', models.DateTimeField(blank=True, null=True)),
|
||||
('bot', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='chats', to='Telegram.bot')),
|
||||
('user', models.ForeignKey(blank=True, help_text='Собеседник в личном чате', null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='chats', to='Telegram.user')),
|
||||
],
|
||||
options={
|
||||
'verbose_name': 'чат',
|
||||
'verbose_name_plural': 'чаты',
|
||||
},
|
||||
),
|
||||
migrations.CreateModel(
|
||||
name='Message',
|
||||
fields=[
|
||||
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
|
||||
('direction', models.CharField(choices=[('in', 'входящее'), ('out', 'исходящее')], max_length=3)),
|
||||
('message_id', models.BigIntegerField(blank=True, null=True)),
|
||||
('text', models.TextField(blank=True)),
|
||||
('created_at', models.DateTimeField(auto_now_add=True)),
|
||||
('chat', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='messages', to='Telegram.chat')),
|
||||
('user', models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='messages', to='Telegram.user')),
|
||||
],
|
||||
options={
|
||||
'verbose_name': 'сообщение',
|
||||
'verbose_name_plural': 'сообщения',
|
||||
'ordering': ['-created_at'],
|
||||
},
|
||||
),
|
||||
migrations.AddConstraint(
|
||||
model_name='chat',
|
||||
constraint=models.UniqueConstraint(fields=('bot', 'chat_id'), name='telegram_chat_unique'),
|
||||
),
|
||||
]
|
||||
+102
-4
@@ -1,31 +1,129 @@
|
||||
from django.db import models
|
||||
|
||||
from . import interfaces
|
||||
from .api import make_api
|
||||
|
||||
|
||||
class Bot(models.Model):
|
||||
"""
|
||||
Пользователь заводит ботов указывая токен. Имя бота можно получить из API.
|
||||
Из админки нужен интерфейс публикации и отключения бота, а также отображение актуального состояния.
|
||||
"""
|
||||
...
|
||||
|
||||
token = models.CharField(max_length=128, unique=True)
|
||||
username = models.CharField(max_length=64, blank=True, editable=False)
|
||||
name = models.CharField(max_length=128, blank=True, editable=False)
|
||||
interface = models.CharField(
|
||||
max_length=255, blank=True,
|
||||
help_text='Класс, определяющий поведение бота. Меняется на лету.',
|
||||
)
|
||||
is_active = models.BooleanField('опубликован', default=False, editable=False)
|
||||
update_offset = models.BigIntegerField(default=0, editable=False)
|
||||
heartbeat_at = models.DateTimeField(null=True, blank=True, editable=False)
|
||||
last_error = models.TextField(blank=True, editable=False)
|
||||
created_at = models.DateTimeField(auto_now_add=True)
|
||||
|
||||
class Meta:
|
||||
verbose_name = 'бот'
|
||||
verbose_name_plural = 'боты'
|
||||
|
||||
def __str__(self):
|
||||
return f'@{self.username}' if self.username else f'bot #{self.pk}'
|
||||
|
||||
def api(self):
|
||||
return make_api(self.token)
|
||||
|
||||
def get_interface(self) -> interfaces.Interface | None:
|
||||
cls = interfaces.get(self.interface)
|
||||
return cls() if cls else None
|
||||
|
||||
def publish(self):
|
||||
"""Проверяет токен, подтягивает имя бота и включает его обслуживание раннером."""
|
||||
api = self.api()
|
||||
me = api.get_me()
|
||||
api.delete_webhook()
|
||||
self.username = me.username or ''
|
||||
self.name = me.first_name or ''
|
||||
self.is_active = True
|
||||
self.last_error = ''
|
||||
self.save()
|
||||
|
||||
def revoke(self):
|
||||
self.is_active = False
|
||||
self.save(update_fields=['is_active'])
|
||||
|
||||
|
||||
class User(models.Model):
|
||||
"""
|
||||
Когда клиент из тг пишет боту личное сообщение, то пользователя нужно завести в базу.
|
||||
"""
|
||||
...
|
||||
|
||||
tg_id = models.BigIntegerField(unique=True)
|
||||
username = models.CharField(max_length=64, blank=True)
|
||||
first_name = models.CharField(max_length=128, blank=True)
|
||||
last_name = models.CharField(max_length=128, blank=True)
|
||||
language_code = models.CharField(max_length=16, blank=True)
|
||||
created_at = models.DateTimeField(auto_now_add=True)
|
||||
updated_at = models.DateTimeField(auto_now=True)
|
||||
|
||||
class Meta:
|
||||
verbose_name = 'пользователь telegram'
|
||||
verbose_name_plural = 'пользователи telegram'
|
||||
|
||||
def __str__(self):
|
||||
return f'@{self.username}' if self.username else str(self.tg_id)
|
||||
|
||||
|
||||
class Chat(models.Model):
|
||||
"""
|
||||
При добавлении бота в группу, нужно добавить запись об этом в базу.
|
||||
При удалении, нужно пометить запись как удаленную, но не удалять.
|
||||
|
||||
Чат принадлежит боту: у одного пользователя с разными ботами разные чаты и разное состояние диалога.
|
||||
"""
|
||||
...
|
||||
|
||||
bot = models.ForeignKey(Bot, on_delete=models.CASCADE, related_name='chats')
|
||||
chat_id = models.BigIntegerField()
|
||||
type = models.CharField(max_length=16)
|
||||
title = models.CharField(max_length=255, blank=True)
|
||||
user = models.ForeignKey(
|
||||
User, on_delete=models.SET_NULL, null=True, blank=True, related_name='chats',
|
||||
help_text='Собеседник в личном чате',
|
||||
)
|
||||
state = models.CharField(max_length=128, blank=True, help_text='Состояние диалога, задается интерфейсом бота')
|
||||
data = models.JSONField(default=dict, blank=True, help_text='Данные диалога, задаются интерфейсом бота')
|
||||
created_at = models.DateTimeField(auto_now_add=True)
|
||||
deleted_at = models.DateTimeField(null=True, blank=True)
|
||||
|
||||
class Meta:
|
||||
verbose_name = 'чат'
|
||||
verbose_name_plural = 'чаты'
|
||||
constraints = [models.UniqueConstraint(fields=['bot', 'chat_id'], name='telegram_chat_unique')]
|
||||
|
||||
def __str__(self):
|
||||
return self.title or str(self.user or self.chat_id)
|
||||
|
||||
|
||||
class Message(models.Model):
|
||||
"""
|
||||
Исходящие и входящие сообщения между администратором и пользователем.
|
||||
"""
|
||||
...
|
||||
|
||||
class Direction(models.TextChoices):
|
||||
IN = 'in', 'входящее'
|
||||
OUT = 'out', 'исходящее'
|
||||
|
||||
chat = models.ForeignKey(Chat, on_delete=models.CASCADE, related_name='messages')
|
||||
user = models.ForeignKey(User, on_delete=models.SET_NULL, null=True, blank=True, related_name='messages')
|
||||
direction = models.CharField(max_length=3, choices=Direction.choices)
|
||||
message_id = models.BigIntegerField(null=True, blank=True)
|
||||
text = models.TextField(blank=True)
|
||||
created_at = models.DateTimeField(auto_now_add=True)
|
||||
|
||||
class Meta:
|
||||
verbose_name = 'сообщение'
|
||||
verbose_name_plural = 'сообщения'
|
||||
ordering = ['-created_at']
|
||||
|
||||
def __str__(self):
|
||||
return f'{self.get_direction_display()}: {self.text[:50]}'
|
||||
@@ -0,0 +1,90 @@
|
||||
"""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)
|
||||
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()
|
||||
+147
-2
@@ -1,3 +1,148 @@
|
||||
from django.test import TestCase
|
||||
from types import SimpleNamespace
|
||||
|
||||
# Create your tests here.
|
||||
from django.test import TestCase
|
||||
from telebot import types
|
||||
|
||||
from . import interfaces
|
||||
from .dispatcher import handle_update
|
||||
from .interfaces import Interface, command, register, state, text
|
||||
from .models import Bot, Chat, Message, User
|
||||
|
||||
|
||||
class FakeApi:
|
||||
def __init__(self):
|
||||
self.sent = []
|
||||
|
||||
def send_message(self, chat_id, text, reply_markup=None, parse_mode=None):
|
||||
self.sent.append((chat_id, text, reply_markup))
|
||||
return SimpleNamespace(message_id=len(self.sent))
|
||||
|
||||
|
||||
@register
|
||||
class EchoInterface(Interface):
|
||||
title = 'Echo'
|
||||
|
||||
@command('start')
|
||||
@text('меню')
|
||||
def menu(self, ctx):
|
||||
ctx.reset()
|
||||
ctx.reply('menu', keyboard=[['ask']])
|
||||
|
||||
@text('ask')
|
||||
def ask(self, ctx):
|
||||
ctx.state = 'await_name'
|
||||
ctx.reply('name?')
|
||||
|
||||
@state('await_name')
|
||||
def name(self, ctx):
|
||||
ctx.data['name'] = ctx.text
|
||||
ctx.state = ''
|
||||
ctx.reply(f'hi {ctx.text}')
|
||||
|
||||
def fallback(self, ctx):
|
||||
ctx.reply(f'echo {ctx.text}')
|
||||
|
||||
|
||||
@register
|
||||
class ChildInterface(EchoInterface):
|
||||
@command('start')
|
||||
def child_start(self, ctx):
|
||||
ctx.reply('child')
|
||||
|
||||
|
||||
def message_update(update_id, text, chat_type='private'):
|
||||
return types.Update.de_json({
|
||||
'update_id': update_id,
|
||||
'message': {
|
||||
'message_id': update_id,
|
||||
'date': 0,
|
||||
'text': text,
|
||||
'chat': {'id': 42, 'type': chat_type, 'title': 'group' if chat_type != 'private' else None},
|
||||
'from': {'id': 7, 'is_bot': False, 'first_name': 'Ann', 'username': 'ann'},
|
||||
},
|
||||
})
|
||||
|
||||
|
||||
def member_update(update_id, status):
|
||||
return types.Update.de_json({
|
||||
'update_id': update_id,
|
||||
'my_chat_member': {
|
||||
'chat': {'id': -100, 'type': 'group', 'title': 'group'},
|
||||
'from': {'id': 7, 'is_bot': False, 'first_name': 'Ann'},
|
||||
'date': 0,
|
||||
'old_chat_member': {'user': {'id': 1, 'is_bot': True, 'first_name': 'bot'}, 'status': 'left'},
|
||||
'new_chat_member': {'user': {'id': 1, 'is_bot': True, 'first_name': 'bot'}, 'status': status},
|
||||
},
|
||||
})
|
||||
|
||||
|
||||
class DispatcherTests(TestCase):
|
||||
def setUp(self):
|
||||
self.api = FakeApi()
|
||||
self.bot = Bot.objects.create(token='1:test', interface=interfaces.key(EchoInterface))
|
||||
|
||||
def send(self, text, update_id=1, **kwargs):
|
||||
handle_update(self.api, self.bot, message_update(update_id, text, **kwargs))
|
||||
return self.api.sent[-1][1] if self.api.sent else None
|
||||
|
||||
def test_command_text_state_fallback(self):
|
||||
self.assertEqual(self.send('/start'), 'menu')
|
||||
self.assertEqual(self.send('ask'), 'name?')
|
||||
self.assertEqual(self.send('Bob'), 'hi Bob')
|
||||
self.assertEqual(self.send('whatever'), 'echo whatever')
|
||||
self.assertEqual(Chat.objects.get().data, {'name': 'Bob'})
|
||||
|
||||
def test_command_with_bot_name_and_args(self):
|
||||
self.assertEqual(self.send('/START@my_bot payload'), 'menu')
|
||||
|
||||
def test_saves_user_chat_and_messages(self):
|
||||
self.send('/start')
|
||||
user = User.objects.get()
|
||||
self.assertEqual((user.tg_id, user.username), (7, 'ann'))
|
||||
chat = Chat.objects.get()
|
||||
self.assertEqual((chat.chat_id, chat.user, chat.type), (42, user, 'private'))
|
||||
self.assertEqual(
|
||||
list(Message.objects.order_by('pk').values_list('direction', 'text')),
|
||||
[('in', '/start'), ('out', 'menu')],
|
||||
)
|
||||
|
||||
def test_ignores_group_messages(self):
|
||||
self.assertIsNone(self.send('/start', chat_type='group'))
|
||||
self.assertFalse(Message.objects.exists())
|
||||
|
||||
def test_interface_switch_applies_immediately(self):
|
||||
self.bot.interface = interfaces.key(ChildInterface)
|
||||
self.assertEqual(self.send('/start'), 'child')
|
||||
self.assertEqual(self.send('меню'), 'menu')
|
||||
|
||||
def test_no_interface(self):
|
||||
self.bot.interface = ''
|
||||
self.assertIsNone(self.send('/start'))
|
||||
self.assertEqual(Message.objects.count(), 1)
|
||||
|
||||
def test_membership_marks_chat_deleted(self):
|
||||
handle_update(self.api, self.bot, member_update(1, 'member'))
|
||||
self.assertIsNone(Chat.objects.get().deleted_at)
|
||||
handle_update(self.api, self.bot, member_update(2, 'kicked'))
|
||||
self.assertIsNotNone(Chat.objects.get().deleted_at)
|
||||
self.assertEqual(Chat.objects.count(), 1)
|
||||
|
||||
|
||||
class InterfaceRegistryTests(TestCase):
|
||||
def test_registered_interfaces_are_choices(self):
|
||||
self.assertIn((interfaces.key(EchoInterface), 'Echo'), interfaces.choices())
|
||||
|
||||
class Unregistered(Interface):
|
||||
pass
|
||||
|
||||
self.assertIsNone(interfaces.get(interfaces.key(Unregistered)))
|
||||
|
||||
|
||||
class AdminTests(TestCase):
|
||||
def test_bot_admin_pages(self):
|
||||
from django.contrib.auth import get_user_model
|
||||
admin = get_user_model().objects.create_superuser('root', 'secret-pass-123')
|
||||
self.client.force_login(admin)
|
||||
bot = Bot.objects.create(token='1:test')
|
||||
for url in ('/admin/Telegram/bot/', '/admin/Telegram/bot/add/', f'/admin/Telegram/bot/{bot.pk}/change/'):
|
||||
self.assertEqual(self.client.get(url).status_code, 200, url)
|
||||
@@ -1,3 +0,0 @@
|
||||
from django.shortcuts import render
|
||||
|
||||
# Create your views here.
|
||||
@@ -28,3 +28,37 @@
|
||||
- 20Gb диска
|
||||
|
||||
Делается под Debian 13.
|
||||
|
||||
## Использование
|
||||
|
||||
Зависимости: ansible-core, ansible-runner (и cryptography, которая приходит с ansible-core). Django не нужен.
|
||||
На машине, откуда идет управление, должен быть `ssh`. `sshpass` не нужен: вход по паролю идет через `ssh_askpass`.
|
||||
|
||||
```python
|
||||
from serverus import Host, Server, generate_keypair, generate_password, generate_ssh_port
|
||||
|
||||
keys = generate_keypair()
|
||||
root_password = generate_password()
|
||||
ssh_port = generate_ssh_port()
|
||||
|
||||
# свежий сервер: вход по паролю на 22 порт
|
||||
server = Server(Host('1.2.3.4', private_key=keys.private_key))
|
||||
result = server.bootstrap('пароль от провайдера', keys.public_key, root_password, ssh_port)
|
||||
|
||||
# дальше только по ключу на новом порту
|
||||
server = Server(Host('1.2.3.4', port=ssh_port, private_key=keys.private_key))
|
||||
server.facts().data # {'uname': ..., 'distribution': 'Debian', 'version': '13'}
|
||||
server.release(generate_password())
|
||||
```
|
||||
|
||||
Каждая операция возвращает `Result(ok, status, data, error)`. Ключи, пароли и артефакты ansible живут во временной папке только на время запуска. Хранить ключи и пароли - задача вызывающей стороны.
|
||||
|
||||
- `bootstrap` - проверяет Debian 13, прописывает ключ, меняет пароль root, обновляет систему, включает nftables и переносит ssh на новый порт с входом только по ключу. Порт меняется в два шага с проверкой входа по ключу, чтобы не потерять доступ.
|
||||
- `release` - ставит новый пароль root, сбрасывает nftables, удаляет drop-in sshd (возвращаются настройки провайдера) и authorized_keys.
|
||||
- `facts` - `uname -a` и версия ОС.
|
||||
|
||||
Настройки ssh пишутся в `/etc/ssh/sshd_config.d/00-serverus.conf`, основной конфиг не трогаем.
|
||||
|
||||
Плейбуки лежат в `playbooks/`, общие шаги в `playbooks/tasks/`. Новая операция - это плейбук и метод в `Server` или его наследнике.
|
||||
|
||||
Тесты: `python -m unittest serverus.tests` из `src/`.
|
||||
@@ -0,0 +1,15 @@
|
||||
from .host import Host
|
||||
from .runner import Result, run_playbook
|
||||
from .credentials import KeyPair, generate_keypair, generate_password, generate_ssh_port
|
||||
from .server import Server
|
||||
|
||||
__all__ = [
|
||||
'Host',
|
||||
'KeyPair',
|
||||
'Result',
|
||||
'Server',
|
||||
'generate_keypair',
|
||||
'generate_password',
|
||||
'generate_ssh_port',
|
||||
'run_playbook',
|
||||
]
|
||||
@@ -0,0 +1,41 @@
|
||||
import secrets
|
||||
import string
|
||||
from dataclasses import dataclass
|
||||
|
||||
from cryptography.hazmat.primitives import serialization
|
||||
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
|
||||
|
||||
PASSWORD_ALPHABET = string.ascii_letters + string.digits
|
||||
SSH_PORT_RANGE = (20000, 60000)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class KeyPair:
|
||||
private_key: str
|
||||
public_key: str
|
||||
|
||||
def __repr__(self):
|
||||
return f'KeyPair({self.public_key})'
|
||||
|
||||
|
||||
def generate_keypair(comment: str = 'serverus') -> KeyPair:
|
||||
key = Ed25519PrivateKey.generate()
|
||||
private = key.private_bytes(
|
||||
serialization.Encoding.PEM,
|
||||
serialization.PrivateFormat.OpenSSH,
|
||||
serialization.NoEncryption(),
|
||||
).decode()
|
||||
public = key.public_key().public_bytes(
|
||||
serialization.Encoding.OpenSSH,
|
||||
serialization.PublicFormat.OpenSSH,
|
||||
).decode()
|
||||
return KeyPair(private_key=private, public_key=f'{public} {comment}')
|
||||
|
||||
|
||||
def generate_password(length: int = 32) -> str:
|
||||
return ''.join(secrets.choice(PASSWORD_ALPHABET) for _ in range(length))
|
||||
|
||||
|
||||
def generate_ssh_port() -> int:
|
||||
low, high = SSH_PORT_RANGE
|
||||
return low + secrets.randbelow(high - low)
|
||||
@@ -0,0 +1,16 @@
|
||||
from dataclasses import dataclass
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Host:
|
||||
"""Адрес сервера и способ входа на него."""
|
||||
|
||||
address: str
|
||||
port: int = 22
|
||||
user: str = 'root'
|
||||
private_key: str | None = None
|
||||
"""Приватный ключ в формате OpenSSH. Без него вход возможен только по паролю."""
|
||||
|
||||
def __repr__(self):
|
||||
# ключ не должен попадать в логи
|
||||
return f'Host({self.user}@{self.address}:{self.port})'
|
||||
@@ -0,0 +1,72 @@
|
||||
# extravars: login_password, public_key, new_password, ssh_port
|
||||
#
|
||||
# Порт ssh меняется в два шага, чтобы не потерять сервер:
|
||||
# сначала sshd слушает старый и новый порт, после проверки входа по ключу на новом порту старый закрывается.
|
||||
|
||||
- name: Prepare fresh server (password login)
|
||||
hosts: target
|
||||
gather_facts: true
|
||||
vars:
|
||||
ansible_password: "{{ login_password }}"
|
||||
tasks:
|
||||
- name: Require Debian 13
|
||||
ansible.builtin.assert:
|
||||
that:
|
||||
- ansible_facts.distribution == 'Debian'
|
||||
- ansible_facts.distribution_major_version == '13'
|
||||
fail_msg: "Supported only Debian 13, got {{ ansible_facts.distribution }} {{ ansible_facts.distribution_version }}"
|
||||
|
||||
- name: Create .ssh directory
|
||||
ansible.builtin.file:
|
||||
path: /root/.ssh
|
||||
state: directory
|
||||
mode: "0700"
|
||||
|
||||
- name: Authorize key
|
||||
ansible.builtin.lineinfile:
|
||||
path: /root/.ssh/authorized_keys
|
||||
line: "{{ public_key }}"
|
||||
create: true
|
||||
mode: "0600"
|
||||
|
||||
- ansible.builtin.import_tasks: tasks/root_password.yml
|
||||
|
||||
- name: Upgrade system
|
||||
ansible.builtin.apt:
|
||||
update_cache: true
|
||||
upgrade: dist
|
||||
|
||||
- ansible.builtin.import_tasks: tasks/firewall.yml
|
||||
vars:
|
||||
firewall_tcp_ports: [ "{{ ansible_port }}", "{{ ssh_port }}" ]
|
||||
|
||||
- ansible.builtin.import_tasks: tasks/sshd.yml
|
||||
vars:
|
||||
sshd_ports: [ "{{ ansible_port }}", "{{ ssh_port }}" ]
|
||||
sshd_keys_only: false
|
||||
|
||||
- name: Harden ssh (key login on new port)
|
||||
hosts: target
|
||||
gather_facts: false
|
||||
vars:
|
||||
ansible_port: "{{ ssh_port }}"
|
||||
tasks:
|
||||
- name: Check key login on new port
|
||||
ansible.builtin.wait_for_connection:
|
||||
timeout: 60
|
||||
|
||||
- ansible.builtin.import_tasks: tasks/firewall.yml
|
||||
vars:
|
||||
firewall_tcp_ports: [ "{{ ssh_port }}" ]
|
||||
|
||||
- ansible.builtin.import_tasks: tasks/sshd.yml
|
||||
vars:
|
||||
sshd_ports: [ "{{ ssh_port }}" ]
|
||||
sshd_keys_only: true
|
||||
|
||||
- name: Reset ssh connection
|
||||
ansible.builtin.meta: reset_connection
|
||||
|
||||
- name: Check hardened ssh
|
||||
ansible.builtin.wait_for_connection:
|
||||
timeout: 60
|
||||
@@ -0,0 +1,16 @@
|
||||
- name: Collect server facts
|
||||
hosts: target
|
||||
gather_facts: true
|
||||
tasks:
|
||||
- name: Read uname
|
||||
ansible.builtin.command: uname -a
|
||||
register: uname
|
||||
changed_when: false
|
||||
|
||||
- name: Return facts
|
||||
ansible.builtin.set_stats:
|
||||
aggregate: false
|
||||
data:
|
||||
uname: "{{ uname.stdout }}"
|
||||
distribution: "{{ ansible_facts.distribution }}"
|
||||
version: "{{ ansible_facts.distribution_major_version }}"
|
||||
@@ -0,0 +1,39 @@
|
||||
# extravars: new_password
|
||||
- name: Release server to owner
|
||||
hosts: target
|
||||
gather_facts: false
|
||||
tasks:
|
||||
- ansible.builtin.import_tasks: tasks/root_password.yml
|
||||
|
||||
- name: Restore default nftables config
|
||||
ansible.builtin.copy:
|
||||
dest: /etc/nftables.conf
|
||||
mode: "0755"
|
||||
content: |
|
||||
#!/usr/sbin/nft -f
|
||||
flush ruleset
|
||||
|
||||
- name: Flush and disable nftables
|
||||
ansible.builtin.systemd_service:
|
||||
name: nftables
|
||||
enabled: false
|
||||
state: stopped
|
||||
|
||||
- name: Remove sshd drop-in
|
||||
ansible.builtin.file:
|
||||
path: /etc/ssh/sshd_config.d/00-serverus.conf
|
||||
state: absent
|
||||
|
||||
- name: Validate sshd config
|
||||
ansible.builtin.command: /usr/sbin/sshd -t
|
||||
changed_when: false
|
||||
|
||||
- name: Remove authorized keys
|
||||
ansible.builtin.file:
|
||||
path: /root/.ssh/authorized_keys
|
||||
state: absent
|
||||
|
||||
- name: Restart ssh
|
||||
ansible.builtin.systemd_service:
|
||||
name: ssh
|
||||
state: restarted
|
||||
@@ -0,0 +1,18 @@
|
||||
# vars: firewall_tcp_ports, firewall_udp_ports (optional)
|
||||
- name: Install nftables
|
||||
ansible.builtin.apt:
|
||||
name: nftables
|
||||
state: present
|
||||
|
||||
- name: Write nftables config
|
||||
ansible.builtin.template:
|
||||
src: templates/nftables.conf.j2
|
||||
dest: /etc/nftables.conf
|
||||
mode: "0755"
|
||||
validate: /usr/sbin/nft -c -f %s
|
||||
|
||||
- name: Apply nftables
|
||||
ansible.builtin.systemd_service:
|
||||
name: nftables
|
||||
enabled: true
|
||||
state: restarted
|
||||
@@ -0,0 +1,6 @@
|
||||
# vars: new_password
|
||||
- name: Set root password
|
||||
ansible.builtin.shell: chpasswd
|
||||
args:
|
||||
stdin: "root:{{ new_password }}"
|
||||
no_log: true
|
||||
@@ -0,0 +1,16 @@
|
||||
# vars: sshd_ports, sshd_keys_only
|
||||
- name: Write sshd drop-in
|
||||
ansible.builtin.template:
|
||||
src: templates/sshd.conf.j2
|
||||
dest: /etc/ssh/sshd_config.d/00-serverus.conf
|
||||
mode: "0644"
|
||||
validate: /usr/sbin/sshd -t -f %s
|
||||
|
||||
- name: Validate full sshd config
|
||||
ansible.builtin.command: /usr/sbin/sshd -t
|
||||
changed_when: false
|
||||
|
||||
- name: Restart ssh
|
||||
ansible.builtin.systemd_service:
|
||||
name: ssh
|
||||
state: restarted
|
||||
@@ -0,0 +1,28 @@
|
||||
#!/usr/sbin/nft -f
|
||||
# {{ ansible_managed }}
|
||||
|
||||
flush ruleset
|
||||
|
||||
table inet filter {
|
||||
chain input {
|
||||
type filter hook input priority filter; policy drop;
|
||||
|
||||
iif "lo" accept
|
||||
ct state established,related accept
|
||||
ct state invalid drop
|
||||
meta l4proto { icmp, ipv6-icmp } accept
|
||||
|
||||
tcp dport { {{ firewall_tcp_ports | join(', ') }} } accept
|
||||
{% if firewall_udp_ports | default([]) %}
|
||||
udp dport { {{ firewall_udp_ports | join(', ') }} } accept
|
||||
{% endif %}
|
||||
}
|
||||
|
||||
chain forward {
|
||||
type filter hook forward priority filter; policy drop;
|
||||
}
|
||||
|
||||
chain output {
|
||||
type filter hook output priority filter; policy accept;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
# {{ ansible_managed }}
|
||||
# Удаление этого файла возвращает настройки ssh, которые были до serverus.
|
||||
{% for port in sshd_ports %}
|
||||
Port {{ port }}
|
||||
{% endfor %}
|
||||
{% if sshd_keys_only %}
|
||||
PermitRootLogin prohibit-password
|
||||
PasswordAuthentication no
|
||||
KbdInteractiveAuthentication no
|
||||
AuthenticationMethods publickey
|
||||
{% endif %}
|
||||
@@ -0,0 +1,95 @@
|
||||
import os
|
||||
import sys
|
||||
import tempfile
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Protocol
|
||||
|
||||
import ansible_runner
|
||||
|
||||
from .host import Host
|
||||
|
||||
PLAYBOOKS_DIR = Path(__file__).resolve().parent / 'playbooks'
|
||||
TARGET = 'target'
|
||||
# ansible-playbook ставится рядом с интерпретатором, но venv может быть не активирован
|
||||
PATH = os.pathsep.join([str(Path(sys.executable).parent), os.environ.get('PATH', '')])
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Result:
|
||||
ok: bool
|
||||
status: str
|
||||
data: dict = field(default_factory=dict)
|
||||
"""Данные, которые плейбук вернул через set_stats."""
|
||||
error: str = ''
|
||||
|
||||
|
||||
class Runner(Protocol):
|
||||
def __call__(self, playbook: str, host: Host, extravars: dict) -> Result: ...
|
||||
|
||||
|
||||
def run_playbook(playbook: str, host: Host, extravars: dict) -> Result:
|
||||
"""Выполняет плейбук из serverus/playbooks на одном хосте.
|
||||
|
||||
Все файлы запуска (ключ, extravars, артефакты) живут во временной папке и удаляются после выполнения.
|
||||
"""
|
||||
with tempfile.TemporaryDirectory(prefix='serverus-') as private_dir:
|
||||
result = ansible_runner.run(
|
||||
private_data_dir=private_dir,
|
||||
project_dir=str(PLAYBOOKS_DIR),
|
||||
playbook=f'{playbook}.yml',
|
||||
inventory={'all': {'hosts': {TARGET: _host_vars(host, private_dir)}}},
|
||||
extravars=extravars,
|
||||
envvars={
|
||||
'ANSIBLE_NOCOLOR': '1',
|
||||
'PATH': PATH,
|
||||
# свои ssh control сокеты на каждый запуск: иначе чужое открытое соединение пропустит без авторизации
|
||||
'ANSIBLE_SSH_CONTROL_PATH_DIR': os.path.join(private_dir, 'cp'),
|
||||
},
|
||||
quiet=True,
|
||||
)
|
||||
events = list(result.events)
|
||||
with result.stdout as stdout:
|
||||
output = stdout.read()
|
||||
ok = result.status == 'successful'
|
||||
return Result(
|
||||
ok=ok,
|
||||
status=result.status,
|
||||
data=_stats_data(events),
|
||||
error='' if ok else _error(events) or output.strip()[-1000:],
|
||||
)
|
||||
|
||||
|
||||
def _host_vars(host: Host, private_dir: str) -> dict:
|
||||
host_vars = {
|
||||
'ansible_host': host.address,
|
||||
'ansible_port': host.port,
|
||||
'ansible_user': host.user,
|
||||
'ansible_python_interpreter': 'auto_silent',
|
||||
# вход по паролю без sshpass
|
||||
'ansible_ssh_password_mechanism': 'ssh_askpass',
|
||||
# TODO: хранить known_hosts сервера и включить проверку ключа хоста
|
||||
'ansible_ssh_common_args': '-o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o LogLevel=ERROR',
|
||||
}
|
||||
if host.private_key:
|
||||
key_file = Path(private_dir) / 'id_host'
|
||||
key_file.write_text(host.private_key.rstrip('\n') + '\n')
|
||||
os.chmod(key_file, 0o600)
|
||||
host_vars['ansible_ssh_private_key_file'] = str(key_file)
|
||||
return host_vars
|
||||
|
||||
|
||||
def _stats_data(events: list[dict]) -> dict:
|
||||
for event in reversed(events):
|
||||
if event.get('event') == 'playbook_on_stats':
|
||||
return event.get('event_data', {}).get('artifact_data', {}) or {}
|
||||
return {}
|
||||
|
||||
|
||||
def _error(events: list[dict]) -> str:
|
||||
for event in reversed(events):
|
||||
if event.get('event') in ('runner_on_failed', 'runner_on_unreachable'):
|
||||
event_data = event.get('event_data', {})
|
||||
res = event_data.get('res', {})
|
||||
return f"{event_data.get('task', '')}: {res.get('msg') or res.get('stderr') or 'failed'}"
|
||||
return ''
|
||||
@@ -0,0 +1,39 @@
|
||||
from .host import Host
|
||||
from .runner import Result, Runner, run_playbook
|
||||
|
||||
|
||||
class Server:
|
||||
"""Фасад над плейбуками базовой настройки сервера.
|
||||
|
||||
Специализации (панель, прокси, докер-платформа) наследуют этот класс и добавляют свои операции.
|
||||
"""
|
||||
|
||||
def __init__(self, host: Host, runner: Runner = run_playbook):
|
||||
self.host = host
|
||||
self._runner = runner
|
||||
|
||||
def run(self, playbook: str, **extravars) -> Result:
|
||||
return self._runner(playbook, self.host, extravars)
|
||||
|
||||
def facts(self) -> Result:
|
||||
"""Возвращает data: {'uname': ..., 'distribution': ..., 'version': ...}."""
|
||||
return self.run('facts')
|
||||
|
||||
def bootstrap(self, password: str, public_key: str, new_password: str, ssh_port: int) -> Result:
|
||||
"""Первичная настройка свежего сервера (Debian 13).
|
||||
|
||||
Входит по паролю на host.port, прописывает public_key, меняет пароль root на new_password, обновляет систему,
|
||||
переносит ssh на ssh_port с входом только по ключу и включает nftables.
|
||||
Для проверки входа по ключу host.private_key должен быть парой к public_key.
|
||||
"""
|
||||
return self.run(
|
||||
'bootstrap',
|
||||
login_password=password,
|
||||
public_key=public_key,
|
||||
new_password=new_password,
|
||||
ssh_port=ssh_port,
|
||||
)
|
||||
|
||||
def release(self, new_password: str) -> Result:
|
||||
"""Возвращает сервер владельцу: откатывает ssh и nftables, удаляет ключи, ставит пароль root new_password."""
|
||||
return self.run('release', new_password=new_password)
|
||||
@@ -0,0 +1,65 @@
|
||||
import unittest
|
||||
|
||||
from . import Host, Result, Server, generate_keypair, generate_password, generate_ssh_port, run_playbook
|
||||
from .credentials import SSH_PORT_RANGE
|
||||
|
||||
|
||||
class FakeRunner:
|
||||
def __init__(self):
|
||||
self.calls = []
|
||||
|
||||
def __call__(self, playbook, host, extravars):
|
||||
self.calls.append((playbook, host, extravars))
|
||||
return Result(ok=True, status='successful')
|
||||
|
||||
|
||||
class CredentialsTests(unittest.TestCase):
|
||||
def test_keypair(self):
|
||||
pair = generate_keypair('test')
|
||||
self.assertTrue(pair.private_key.startswith('-----BEGIN OPENSSH PRIVATE KEY-----'))
|
||||
self.assertTrue(pair.public_key.startswith('ssh-ed25519 '))
|
||||
self.assertTrue(pair.public_key.endswith(' test'))
|
||||
self.assertNotIn('PRIVATE', repr(pair))
|
||||
|
||||
def test_password(self):
|
||||
self.assertEqual(len(generate_password()), 32)
|
||||
self.assertNotEqual(generate_password(), generate_password())
|
||||
|
||||
def test_ssh_port(self):
|
||||
low, high = SSH_PORT_RANGE
|
||||
self.assertTrue(low <= generate_ssh_port() < high)
|
||||
|
||||
def test_host_repr_hides_key(self):
|
||||
self.assertNotIn('secret', repr(Host('1.2.3.4', private_key='secret')))
|
||||
|
||||
|
||||
class ServerTests(unittest.TestCase):
|
||||
def test_operations_call_playbooks(self):
|
||||
runner = FakeRunner()
|
||||
server = Server(Host('1.2.3.4'), runner=runner)
|
||||
server.facts()
|
||||
server.bootstrap('old', 'ssh-ed25519 AAA', 'new', 40000)
|
||||
server.release('other')
|
||||
self.assertEqual([call[0] for call in runner.calls], ['facts', 'bootstrap', 'release'])
|
||||
self.assertEqual(runner.calls[1][2], {
|
||||
'login_password': 'old',
|
||||
'public_key': 'ssh-ed25519 AAA',
|
||||
'new_password': 'new',
|
||||
'ssh_port': 40000,
|
||||
})
|
||||
|
||||
|
||||
class RunPlaybookTests(unittest.TestCase):
|
||||
def test_facts_on_localhost(self):
|
||||
result = run_playbook('facts', Host('127.0.0.1'), {'ansible_connection': 'local'})
|
||||
self.assertTrue(result.ok, result.error)
|
||||
self.assertIn('Linux', result.data['uname'])
|
||||
|
||||
def test_unreachable_host_reports_error(self):
|
||||
result = run_playbook('facts', Host('127.0.0.1', port=1, private_key='bad'), {})
|
||||
self.assertFalse(result.ok)
|
||||
self.assertTrue(result.error)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
unittest.main()
|
||||
Reference in new issue
Block a user