diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..ebed0bc --- /dev/null +++ b/.dockerignore @@ -0,0 +1,20 @@ +.git/ +.agents/ +.codex/ + +.idea/ +.vscode/ +.venv/ +venv/ +env/ + +__pycache__/ +*.py[cod] +tests +docs +.ruff_cache/ + +.env +.env.* + +db/ diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..90bd831 --- /dev/null +++ b/.env.example @@ -0,0 +1,27 @@ +# PostgreSQL +POSTGRES_DB=bitrix_bot +POSTGRES_USER=postgres +POSTGRES_PASSWORD=change_admin_password +SITE_DB_PASSWORD=change_site_password +BOT_DB_PASSWORD=change_bot_password + +# Внешний nginx проксирует на 127.0.0.1:8000. +PUBLIC_BASE_URL=https://bot.example.ru +SITE_PUBLISHED_PORT=8000 +SITE_PORT=8000 +SITE_WORKERS=2 + +# Локальное приложение Битрикса. Нужно будет заменить ID и SECRET на реальные +# значения, полученные при регистрации приложения в Битриксе. +BITRIX_CLIENT_ID=local.example +BITRIX_CLIENT_SECRET=change_me +BITRIX_OAUTH_TOKEN_URL=https://oauth.bitrix.info/oauth/token/ +BITRIX_TAKE_TO_WORK_STAGE_ID=PREPARATION + +# Telegram +BOT_TOKEN=change_me +BOT_USERNAME=example_bot +BINDING_TOKEN_TTL_SECONDS=600 + +# Ключ шифрования токенов. Должен быть 32 байта в base64. +TOKEN_ENCRYPTION_KEY=change_me diff --git a/.gitattributes b/.gitattributes new file mode 100644 index 0000000..786c643 --- /dev/null +++ b/.gitattributes @@ -0,0 +1,7 @@ +* text=auto + +*.sh text eol=lf +*.sql text eol=lf +*.yaml text eol=lf +*.yml text eol=lf +Dockerfile text eol=lf \ No newline at end of file diff --git a/.gitignore b/.gitignore index dda4f8f..7746fd7 100644 --- a/.gitignore +++ b/.gitignore @@ -160,5 +160,9 @@ cython_debug/ # be found at https://github.com/github/gitignore/blob/main/Global/JetBrains.gitignore # and can be added to the global gitignore or merged into this file. For a more nuclear # option (not recommended) you can uncomment the following to ignore the entire idea folder. -#.idea/ +.idea/ +.vscode/ +.agents/ +.codex/ +.ruff_cache/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..5b89dc3 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,25 @@ +# Источник: https://jtprog.ru/posts/docker-base/ + +# В качестве родителя используем slim-образ с Python 3.13 +FROM python:3.13-slim + +# Просим Python не писать .pyc файлы и не не буферизовать stdin/stdout +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 + +# Задаем рабочую директорию +WORKDIR /srv/bitrix-deals-bot + +# Создаем системную группу и пользователя для запуска приложения +RUN groupadd --system runtime && useradd --system --gid runtime runtime + +# Копируем файл зависимостей и устанавливаем их +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +# Копируем исходный код приложения в контейнер +COPY --chown=runtime:runtime apps ./apps + +# Переходим на созданного пользователя для запуска приложения +USER runtime +CMD ["python", "-m", "apps.bot"] diff --git a/README.md b/README.md index 0fc1ccd..f1c9595 100644 --- a/README.md +++ b/README.md @@ -1,26 +1,156 @@ # BitrixDealsBot -BitrixDealsBot — это телеграм-бот, который позволяет пользователям +**BitrixDealsBot** — это телеграм-бот, который позволяет пользователям взаимодействовать с CRM-системой Bitrix24 для управления сделками и контактами. Цель бота — упростить процесс работы с CRM, предоставляя удобный интерфейс для обновления и отслеживания сделок прямо из Telegram. Задание выполняется в рамках учебной производственной практики для предприятия -Интерволга. +ООО "Интернет-агенство ИНТЕРВОЛГА". ## Формулировка задания + **Telegram-бот “Помощник менеджера CRM”** -Telegram-бот для менеджера по продажам. Бот помогает быстро смотреть новые лиды, +Telegram-бот для менеджера по продажам. Бот помогает быстро смотреть новые лиды, брать их в работу, менять статус и добавлять комментарии. **Стек:** Любой язык, любая БД, REST API Telegram, REST API Битрикс24 **Функции:** + - команда /leads показывает новые лиды (без ответственных); -- команда /lead ### показывает карточку лида по указанному ID: имя, телефон, источник, статус; +- команда /lead ### показывает карточку лида по указанному ID: имя, телефон, + источник, статус; - кнопка “Взять” устанавливает ответственного; -- кнопка “Позвонить позже” устанавливает ответственного и планирует звонок через 1 час; +- кнопка “Позвонить позже” устанавливает ответственного и планирует звонок через + 1 час; - кнопка “Закрыть” возвращает в список лидов /leads; - команда /history показывает историю действий; -- интеграция с Битрикс24 через webhook. \ No newline at end of file +- интеграция с Битрикс24 через webhook (Был осуществлен переход на OAuth). + +## Архитектура решения + +Основные архитектурные решения изложены в разделе ADR (Architecture Decision +Records) в папке `docs/adr`. [Главный документ](docs/adr/main.md) может +использоваться для навигации между заметками о принятых решениях. + +## Настройка + +> [!IMPORTANT] +> Для развертывания проекта требуется Docker и Docker Compose. В Windows +> рекомендуется использовать WSL2. Также требуется внешний nginx для +> проксирования запросов к сайту. + +Клонируйте репозиторий и перейдите в папку проекта: + +```bash +cd BitrixDealsBot +``` + +Скопируйте `.env.example` в `.env` и заполните значения. Ключ шифрования можно +создать командой: + +```bash +python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())" +``` + +`TOKEN_ENCRYPTION_KEY` должен быть одинаковым у сайта и бота. В БД OAuth-токены +попадают уже зашифрованными. + +Сгенерируйте или иным способом придумайте пароли для ролей администратора, +пользователя сайта и бота. В PostgreSQL роли создаются автоматически при первом +создании контейнера `db`. + +Также укажите URL, который будет использоваться в качестве базового адреса +сайта. Соответственно, для этого бы желательно иметь свой домен, соответствующие +DNS-записи и TLS-сертификат (например, от Let's Encrypt). + +Создайте локальное приложение в Битрикс24 и укажите в нем URL для привязки: + +```text +https://bot.example.ru/bitrix/bind +``` + +Битрикс сгенерирует `CLIENT_ID` и `CLIENT_SECRET`, которые нужно указать в +`.env`. В настройках приложения разрешите доступ к CRM и к минимальной +информации о пользователе. +Битрикс передает сайту OAuth-данные приложения. Сайт обменивает `REFRESH_ID` на +новую пару токенов и берет доверенные идентификаторы портала и пользователя из +ответа OAuth-сервера. Затем он проверяет пользователя через `user.current`, +сохраняет зашифрованную пару токенов и показывает ссылку на Telegram. + +Также создайте бота в Telegram через BotFather и укажите его токен в `.env`. +Дополнительно укажите в `.env` имя вашего бота, которое будет использоваться в +ссылках на него. + +Вы также можете изменить в '.env' стандартную стадию сделки, которая будет +использоваться при взятии сделки в работу. По умолчанию это стадия "C1:PREPARATION". +Если вы хотите использовать другую стадию, укажите ее код в +переменной `BITRIX_TAKE_TO_WORK_STAGE_ID`. + +## Запуск в Docker + +```bash +docker compose up --build -d +``` + +Создаются четыре контейнера: + +- `db` — PostgreSQL без опубликованного порта; +- `site` — Flask/Gunicorn на `127.0.0.1:8000` хостовой машины; +- `bot` — aiogram polling без входящего порта; +- `migrate` — контейнер для миграции БД, запускается при каждом + `docker compose up` и завершается после выполнения миграций. + +Инициализация БД выполняется автоматически только для нового volume. + +### Внешний nginx + +Nginx работает на хосте, вне Docker: + +```nginx +location / { + proxy_pass http://127.0.0.1:8000; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; +} +``` + +Если на вашем сервере есть панели типа FastPanel, Plesk, ISPmanager, то вы +можете создать сайт с обратным прокси прямо в интерфейсе панели. + +## Транзакции и доступ к БД + +Сайт подключается ролью `site_app`, бот — `bot_app`. Прямого доступа к таблицам +у них нет. + +- `binding.issue_v1` сохраняет OAuth-данные и выпускает ссылку одной + транзакцией. +- `binding.consume_v1` атомарно погашает ссылку и создает привязку. +- `oauth.*` выдает и обновляет токены только боту. +- короткая DB-аренда защищает refresh-токен от параллельного обновления. + +Сетевые запросы к Битриксу не выполняются внутри транзакций. +Назначения на сделки хранятся только в Битриксу. Локальная блокировка в процессе +бота не дает двум Telegram-пользователям одновременно взять одну сделку, а +повторная проверка `ASSIGNED_BY_ID` отсекает устаревшие кнопки. + +## Команды бота + +- `/start bind_` — привязать пользователя. +- `/deals` или `/leads` — показать сделки с фильтром и пагинацией. +- `/deal 123` — открыть карточку сделки. +- `/help` — показать справку. + +Фильтр «Мои сделки» показывает сделки всех стадий, где ответственным назначен +привязанный Битрикс-пользователь. + +Кнопка «Позвонить позже» создает дело с напоминанием в Битриксе через час. +История показывает пять последних переходов сделки по стадиям. Кнопка «Стать +ответственным и взять в работу» использует ID привязанного Битрикс-пользователя. +Ответственный также видит кнопку перехода на следующую стадию; финальный переход +выделен как завершение сделки. Перед переходом бот запрашивает необязательный +комментарий для таймлайна: пустая строка или прочерк означают переход без него. diff --git a/app/main.py b/app/main.py deleted file mode 100644 index e76c1e9..0000000 --- a/app/main.py +++ /dev/null @@ -1,120 +0,0 @@ -import asyncio -import html -import os -import decimal - -import httpx -from aiogram import Bot, Dispatcher -from aiogram.filters import Command -from aiogram.types import Message -from dotenv import load_dotenv - -load_dotenv() - -BOT_TOKEN = os.getenv("BOT_TOKEN") -BITRIX_WEBHOOK_URL = os.getenv("BITRIX_WEBHOOK_URL") - -dp = Dispatcher() - - -async def bitrix_call(method: str, params: dict | None = None) -> dict: - if not BITRIX_WEBHOOK_URL: - raise RuntimeError("BITRIX_WEBHOOK_URL is not set") - - base_url = BITRIX_WEBHOOK_URL.rstrip("/") + "/" - url = base_url + method - - async with httpx.AsyncClient(timeout=15) as client: - response = await client.post(url, json=params or {}) - response.raise_for_status() - data = response.json() - - if "error" in data: - description = data.get("error_description", data["error"]) - raise RuntimeError(f"Bitrix API error: {description}") - - return data - - -def format_deal(deal: dict) -> str: - deal_id = html.escape(str(deal.get("ID", "—"))) - title = html.escape(str(deal.get("TITLE", "Без названия"))) - stage = html.escape(str(deal.get("STAGE_ID", "—"))) - opportunity = html.escape(str(deal.get("OPPORTUNITY", "—"))) - currency = html.escape(str(deal.get("CURRENCY_ID", ""))) - date = html.escape(str(deal.get("DATE_CREATE", "—"))) - - return ( - f"#{deal_id} — {title}\n" - f"Стадия: {stage}\n" - f"Сумма: {decimal.Decimal(opportunity):,.2f} {currency}\n" - f"Дата создания: {date}" - ) - - -@dp.message(Command("start")) -async def start_handler(message: Message) -> None: - await message.answer( - "Привет. Команда /leads покажет последние сделки из Битрикс24." - ) - - -@dp.message(Command("leads")) -async def leads_handler(message: Message) -> None: - await message.answer("Запрашиваю сделки...") - - try: - data = await bitrix_call( - "crm.deal.list", - { - "order": {"DATE_CREATE": "DESC"}, - "filter": {}, - "select": [ - "ID", - "TITLE", - "STAGE_ID", - "OPPORTUNITY", - "CURRENCY_ID", - "DATE_CREATE", - ], - "start": 0, - }, - ) - - deals = data.get("result", []) - - if not deals: - await message.answer("Сделки не найдены.") - return - - text = "\n\n".join(format_deal(deal) for deal in deals[:10]) - - await message.answer( - f"Последние сделки:\n\n{text}", - parse_mode="HTML", - ) - - except httpx.HTTPStatusError as e: - await message.answer( - f"Ошибка HTTP при запросе к Битрикс24: {e.response.status_code}") - - except httpx.RequestError: - await message.answer("Не удалось подключиться к Битрикс24.") - - except RuntimeError as e: - await message.answer(f"Ошибка: {html.escape(str(e))}") - - except Exception: - await message.answer("Произошла неизвестная ошибка.") - - -async def main() -> None: - if not BOT_TOKEN: - raise RuntimeError("BOT_TOKEN is not set") - - bot = Bot(token=BOT_TOKEN) - await dp.start_polling(bot) - - -if __name__ == "__main__": - asyncio.run(main()) diff --git a/apps/__init__.py b/apps/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/bot/__init__.py b/apps/bot/__init__.py new file mode 100644 index 0000000..21db616 --- /dev/null +++ b/apps/bot/__init__.py @@ -0,0 +1,3 @@ +from .main import main + +__all__ = ["main"] diff --git a/apps/bot/__main__.py b/apps/bot/__main__.py new file mode 100644 index 0000000..5d6a810 --- /dev/null +++ b/apps/bot/__main__.py @@ -0,0 +1,3 @@ +from .main import main + +main() diff --git a/apps/bot/binding.py b/apps/bot/binding.py new file mode 100644 index 0000000..68f1c2b --- /dev/null +++ b/apps/bot/binding.py @@ -0,0 +1,88 @@ +import hashlib +import re + +from .database import BotDatabase +from .domain import Binding + +TOKEN_PATTERN = re.compile(r"^[A-Za-z0-9_-]{20,100}$") + + +def hash_token(token: str) -> bytes: + return hashlib.sha256(token.encode("utf-8")).digest() + + +class BotBindingRepository: + """Доступ бота только к функциям БД схемы binding. + Обертка над хранимыми функциями БД.""" + + def __init__(self, database: BotDatabase) -> None: + self.database = database + + async def consume( + self, + token_hash: bytes, + telegram_user_id: int, + telegram_chat_id: int + ) -> Binding | None: + async with self.database.transaction() as connection: + cursor = await connection.execute( + """ + SELECT * + FROM binding.consume_v1(%s, %s, %s) + """, + (token_hash, telegram_user_id, telegram_chat_id) + ) + row = await cursor.fetchone() + return self._binding(row) if row else None + + async def find( + self, + telegram_user_id: int, + member_id: str | None = None, + ) -> Binding | None: + async with self.database.transaction() as connection: + cursor = await connection.execute( + """ + SELECT * + FROM binding.find_by_telegram_v1(%s, %s) + """, + (telegram_user_id, member_id) + ) + row = await cursor.fetchone() + return self._binding(row) if row else None + + @staticmethod + def _binding(row: dict) -> Binding: + return Binding( + member_id=str(row["member_id"]), + domain=str(row["domain"]), + bitrix_user_id=int(row["bitrix_user_id"]), + telegram_user_id=int(row["telegram_user_id"]) + ) + + +class BindingService: + """Служба управления привязками пользователей.""" + + def __init__( + self, + repository: BotBindingRepository + ) -> None: + self.repository = repository + + async def consume( + self, + token: str, + telegram_user_id: int, + telegram_chat_id: int + ) -> Binding | None: + if not TOKEN_PATTERN.fullmatch(token): + return None + return await self.repository.consume( + hash_token(token), + telegram_user_id, + telegram_chat_id + ) + + async def find(self, telegram_user_id: int) -> Binding | None: + return await self.repository.find(telegram_user_id) diff --git a/apps/bot/bitrix.py b/apps/bot/bitrix.py new file mode 100644 index 0000000..fe57c10 --- /dev/null +++ b/apps/bot/bitrix.py @@ -0,0 +1,161 @@ +import asyncio +from datetime import UTC, datetime, timedelta +from typing import Any + +import httpx + +from .crypto import TokenCipher +from .domain import Binding, OAuthCredentials +from .oauth import BotOAuthRepository + +AUTH_ERRORS = {"expired_token", "invalid_token", "no_auth_found"} + + +class BitrixClient: + """REST-клиент Битрикса с OAuth-контекстом привязанного пользователя.""" + + def __init__( + self, + credentials: BotOAuthRepository, + cipher: TokenCipher, + client_id: str, + client_secret: str, + oauth_token_url: str, + client: httpx.AsyncClient | None = None + ) -> None: + self.credentials = credentials + self.cipher = cipher + self.client_id = client_id + self.client_secret = client_secret + self.oauth_token_url = oauth_token_url + # Передача клиента для упрощения тестирования. + self._client = client or httpx.AsyncClient(timeout=15) + self._owns_client = client is None + + async def call( + self, + binding: Binding, + method: str, + params: dict[str, Any] | None = None + ) -> dict[str, Any]: + credentials = await self.credentials.get(binding) + if not credentials: + raise RuntimeError("OAuth-данные пользователя не найдены") + + # Отправляем запрос и в случае истечения токена запрашиваем обновление. + data = await self._request(credentials, method, params) + if str(data.get("error") or "").lower() in AUTH_ERRORS: + credentials = await self._refresh(credentials, binding) + data = await self._request(credentials, method, params) + + if "error" in data: + description = data.get("error_description", data["error"]) + raise RuntimeError(f"Bitrix API error: {description}") + return data + + async def _request( + self, + credentials: OAuthCredentials, + method: str, + params: dict[str, Any] | None + ) -> dict[str, Any]: + payload = dict(params or {}) + payload["auth"] = self.cipher.decrypt(credentials.access_token) + response = await self._client.post( + f"https://{credentials.domain}/rest/{method}.json", + json=payload + ) + try: + data = response.json() + except ValueError: + response.raise_for_status() + raise RuntimeError("Bitrix вернул некорректный ответ") from None + + # Битрикс присылает полезное описание ошибки и при HTTP 4xx. + if response.is_error and "error" not in data: + response.raise_for_status() + return data + + async def _refresh( + self, + credentials: OAuthCredentials, + binding: Binding + ) -> OAuthCredentials: + # Если другой процесс уже обновляет токен, ждем его завершения. + if not await self.credentials.claim_refresh(credentials): + return await self._wait_for_refresh(credentials, binding) + + try: + # Битрикс возвращает новую пару, поэтому обновляем оба токена. + try: + response = await self._client.get( + self.oauth_token_url, + params={ + "grant_type": "refresh_token", + "client_id": self.client_id, + "client_secret": self.client_secret, + "refresh_token": self.cipher.decrypt( + credentials.refresh_token) + } + ) + response.raise_for_status() + except httpx.HTTPError: + # Не включаем URL с OAuth-секретами в traceback. + raise RuntimeError("Не удалось обновить OAuth-токен") from None + + data = response.json() + if "error" in data: + raise RuntimeError( + "Bitrix OAuth error: " + + str(data.get("error_description") or data["error"]) + ) + + # Проверяем, что обновленный токен принадлежит тому же порталу + # и пользователю. + if data.get("member_id") not in {None, credentials.member_id}: + raise RuntimeError("Bitrix вернул токен другого портала") + if int(data.get("user_id", credentials.bitrix_user_id)) != ( + credentials.bitrix_user_id + ): + raise RuntimeError("Bitrix вернул токен другого пользователя") + + # Обновляем токены в базе и возвращаем новые данные. + expires_at = datetime.now(UTC) + timedelta( + seconds=int(data.get("expires_in", 3600)) + ) + saved = await self.credentials.finish_refresh( + credentials, + self.cipher.encrypt(str(data["access_token"])), + self.cipher.encrypt(str(data["refresh_token"])), + expires_at + ) + + # Если другой процесс успел обновить токен, ждем его завершения. + if not saved: + return await self._wait_for_refresh(credentials, binding) + + updated = await self.credentials.get(binding) + if not updated: + raise RuntimeError("Обновленные OAuth-данные не найдены") + return updated + + except Exception: + await self.credentials.release_refresh(credentials) + raise + + async def _wait_for_refresh( + self, + previous: OAuthCredentials, + binding: Binding + ) -> OAuthCredentials: + for _ in range(80): + await asyncio.sleep(0.2) + current = await self.credentials.get(binding) + if current and current.version > previous.version: + return current + + raise RuntimeError("Не удалось дождаться обновления OAuth-токена") + + async def close(self) -> None: + if self._owns_client: + await self._client.aclose() diff --git a/apps/bot/config.py b/apps/bot/config.py new file mode 100644 index 0000000..7f310ad --- /dev/null +++ b/apps/bot/config.py @@ -0,0 +1,42 @@ +import os +from dataclasses import dataclass + +DEFAULT_TAKE_TO_WORK_STAGE_ID = "PREPARATION" + + +def _required(name: str) -> str: + value = os.getenv(name) + if not value: + raise RuntimeError(f"{name} is not set") + return value + + +@dataclass(frozen=True) +class BotConfig: + """Настройки процесса Telegram-бота.""" + + bot_token: str + database_url: str + bitrix_client_id: str + bitrix_client_secret: str + token_encryption_key: str + oauth_token_url: str + take_to_work_stage_id: str = DEFAULT_TAKE_TO_WORK_STAGE_ID + + @classmethod + def from_env(cls) -> "BotConfig": + return cls( + bot_token=_required("BOT_TOKEN"), + database_url=_required("DATABASE_URL"), + bitrix_client_id=_required("BITRIX_CLIENT_ID"), + bitrix_client_secret=_required("BITRIX_CLIENT_SECRET"), + token_encryption_key=_required("TOKEN_ENCRYPTION_KEY"), + oauth_token_url=os.getenv( + "BITRIX_OAUTH_TOKEN_URL", + "https://oauth.bitrix.info/oauth/token/", + ), + take_to_work_stage_id=os.getenv( + "BITRIX_TAKE_TO_WORK_STAGE_ID", + DEFAULT_TAKE_TO_WORK_STAGE_ID + ) + ) diff --git a/apps/bot/crypto.py b/apps/bot/crypto.py new file mode 100644 index 0000000..2eb831a --- /dev/null +++ b/apps/bot/crypto.py @@ -0,0 +1,20 @@ +from cryptography.fernet import Fernet, InvalidToken + + +class TokenCipher: + """Шифрует OAuth-токены перед хранением в БД.""" + + def __init__(self, key: str) -> None: + try: + self._fernet = Fernet(key.encode("ascii")) + except (ValueError, UnicodeEncodeError) as error: + raise RuntimeError("TOKEN_ENCRYPTION_KEY is invalid") from error + + def encrypt(self, value: str) -> bytes: + return self._fernet.encrypt(value.encode("utf-8")) + + def decrypt(self, value: bytes) -> str: + try: + return self._fernet.decrypt(value).decode("utf-8") + except InvalidToken as error: + raise RuntimeError("Не удалось расшифровать OAuth-токен") from error diff --git a/apps/bot/database.py b/apps/bot/database.py new file mode 100644 index 0000000..b3c2d63 --- /dev/null +++ b/apps/bot/database.py @@ -0,0 +1,38 @@ +from collections.abc import AsyncGenerator +from contextlib import asynccontextmanager + +from psycopg import AsyncConnection +from psycopg.rows import dict_row +from psycopg_pool import AsyncConnectionPool + + +class BotDatabase: + """Пул соединений БД для асинхронного процесса бота.""" + + def __init__(self, database_url: str) -> None: + # Пул может содержать в себе максимум 5 соединений. + self.pool = AsyncConnectionPool( + conninfo=database_url, + min_size=1, + max_size=5, + open=False, + # Фабрика для представления строк БД как словарей. + kwargs={"row_factory": dict_row} + ) + + async def open(self) -> None: + await self.pool.open(wait=True) + + async def close(self) -> None: + await self.pool.close() + + @asynccontextmanager + async def transaction(self) -> AsyncGenerator[AsyncConnection]: + async with self.pool.connection() as connection: + async with connection.transaction(): + yield connection + + async def ping(self) -> bool: + async with self.pool.connection() as connection: + cursor = await connection.execute("SELECT 1") + return await cursor.fetchone() is not None diff --git a/apps/bot/deals.py b/apps/bot/deals.py new file mode 100644 index 0000000..2ee6328 --- /dev/null +++ b/apps/bot/deals.py @@ -0,0 +1,641 @@ +import asyncio +import logging +import time +from collections.abc import AsyncGenerator +from contextlib import asynccontextmanager +from dataclasses import dataclass +from datetime import UTC, datetime, timedelta + +from .bitrix import BitrixClient +from .domain import ( + DEALS_PER_PAGE, + Binding, + ClientInfo, + DealPage, + DealStage, + DealStageAdvance, + DealStageFilter, +) + +logger = logging.getLogger(__name__) + + +class DealAssignmentConflict(RuntimeError): + """Ответственный изменился после показа карточки.""" + + +class DealAdvanceForbidden(RuntimeError): + """Стадию может менять только ответственный за сделку.""" + + +class DealStageConflict(RuntimeError): + """Стадия изменилась после запроса комментария.""" + + +class DealCommentSaveError(RuntimeError): + """Стадия изменена, но комментарий не добавлен.""" + + def __init__(self, advance: DealStageAdvance) -> None: + super().__init__( + "Стадия изменена, но комментарий не удалось сохранить." + ) + self.advance = advance + + +@dataclass +class _DealLockEntry: + lock: asyncio.Lock + users: int = 0 + + +class DealService: + """Загрузка и изменение сделок Bitrix.""" + + deal_select = [ + "ID", + "TITLE", + "STAGE_ID", + "IS_NEW", + "OPPORTUNITY", + "CURRENCY_ID", + "DATE_CREATE", + "ASSIGNED_BY_ID", + "CONTACT_ID", + "COMPANY_ID", + "SOURCE_ID", + "COMMENTS", + ] + + def __init__( + self, + bitrix: BitrixClient, + take_to_work_stage_id: str + ) -> None: + self.bitrix = bitrix + self.take_to_work_stage_id = take_to_work_stage_id + self._deal_locks: dict[tuple[str, str], _DealLockEntry] = {} + self._deal_locks_guard = asyncio.Lock() + self._stage_cache: dict[ + tuple[str, int, int], tuple[float, tuple[DealStage, ...]] + ] = {} + + async def list_by_stage( + self, + binding: Binding, + stage_key: str = "new", + page: int = 0, + limit: int = DEALS_PER_PAGE + ) -> DealPage: + stage_filters = await self.stage_filters(binding) + stage_filter = self._select_stage_filter(stage_filters, stage_key) + bitrix_filter = {} + if stage_filter.stage_id: + bitrix_filter["STAGE_ID"] = stage_filter.stage_id + if stage_filter.assigned_to_viewer: + bitrix_filter["ASSIGNED_BY_ID"] = binding.bitrix_user_id + + page = max(page, 0) + start_index = page * limit + end_index = start_index + limit + loaded_deals = [] + bitrix_start: int | None = 0 + total_deals = 0 + + # Битрикс и Telegram используют страницы разного размера. + while bitrix_start is not None and len(loaded_deals) < end_index: + # https://apidocs.bitrix24.ru/api-reference/crm/deals/crm-deal-list.html + data = await self.bitrix.call( + binding, + "crm.deal.list", + { + "order": {"DATE_CREATE": "DESC"}, + "filter": bitrix_filter, + "select": self.deal_select, + "start": bitrix_start + } + ) + loaded_deals.extend(data.get("result", [])) + total_deals = int(data.get("total", len(loaded_deals))) + bitrix_start = data.get("next") + + deals = loaded_deals[start_index:end_index] + total_pages = max(1, (total_deals + limit - 1) // limit) + return DealPage( + deals, + page, + total_deals, + total_pages, + stage_filter, + stage_filters + ) + + async def stage_filters( + self, + binding: Binding + ) -> tuple[DealStageFilter, ...]: + stages = await self._stages(binding, category_id=0) + filters = [ + DealStageFilter( + key="mine", + title="Мои сделки", + stage_id=None, + assigned_to_viewer=True, + ) + ] + filters.extend( + DealStageFilter(stage.stage_id, stage.title, stage.stage_id) + for stage in stages + ) + filters.append(DealStageFilter("all", "Все", None)) + return tuple(filters) + + def _select_stage_filter( + self, + filters: tuple[DealStageFilter, ...], + stage_key: str + ) -> DealStageFilter: + initial = next( + (item for item in filters if item.stage_id is not None), + filters[-1] + ) + # Callback `new` означает первую стадию, полученную из Битрикса. + if stage_key == "new": + return initial + + selected = next((item for item in filters if item.key == stage_key), + None) + return selected or initial + + async def _stages( + self, + binding: Binding, + category_id: int + ) -> tuple[DealStage, ...]: + key = (binding.member_id, binding.bitrix_user_id, category_id) + cached = self._stage_cache.get(key) + if cached and cached[0] > time.monotonic(): + return cached[1] + + entity_id = "DEAL_STAGE" if category_id == 0 else ( + f"DEAL_STAGE_{category_id}" + ) + items = [] + bitrix_start: int | None = 0 + # https://apidocs.bitrix24.ru/api-reference/crm/status/crm-status-list.html + while bitrix_start is not None: + data = await self.bitrix.call( + binding, + "crm.status.list", + { + "order": {"SORT": "ASC"}, + "filter": {"ENTITY_ID": entity_id}, + "start": bitrix_start + } + ) + items.extend(data.get("result", [])) + bitrix_start = data.get("next") + + stages = [] + seen_stage_ids = set() + for item in items: + raw_stage_id = str(item.get("STATUS_ID") or "") + if not raw_stage_id: + continue + + stage_id = raw_stage_id + prefix = f"C{category_id}:" + if category_id and not stage_id.startswith(prefix): + stage_id = prefix + stage_id + if stage_id in seen_stage_ids: + continue + + semantics = self._stage_semantics(item) + stages.append( + DealStage( + stage_id=stage_id, + title=str(item.get("NAME") or raw_stage_id), + semantics=semantics, + ) + ) + seen_stage_ids.add(stage_id) + + result = tuple(stages) + self._stage_cache[key] = (time.monotonic() + 300, result) + return result + + async def _stage_map( + self, + binding: Binding, + category_id: int + ) -> dict[str, str]: + stages = await self._stages(binding, category_id) + return {stage.stage_id: stage.title for stage in stages} + + @staticmethod + def _stage_semantics(item: dict) -> str: + semantics = str(item.get("SEMANTICS") or "").upper() + if semantics in {"P", "S", "F"}: + return semantics + + extra_semantics = str( + (item.get("EXTRA") or {}).get("SEMANTICS") or "" + ).lower() + return { + "process": "P", + "success": "S", + "failure": "F", + }.get(extra_semantics, "P") + + async def stage_name( + self, + binding: Binding, + category_id: int, + stage_id: str + ) -> str | None: + stages = await self._stage_map(binding, category_id) + return stages.get(stage_id) + + async def get(self, binding: Binding, deal_id: str) -> dict | None: + # https://apidocs.bitrix24.ru/api-reference/crm/deals/crm-deal-get.html + data = await self.bitrix.call(binding, "crm.deal.get", {"id": deal_id}) + deal = data.get("result") + if not deal: + return None + + client = await self.get_client_info(binding, deal) + deal["CLIENT_NAME"] = client.name + deal["CLIENT_COMPANY"] = client.company + deal["CLIENT_PHONE"] = client.phone + deal["SOURCE_NAME"] = await self.get_source_name( + binding, str(deal.get("SOURCE_ID") or "") + ) + deal["STAGE_NAME"] = await self.stage_name( + binding, + int(deal.get("CATEGORY_ID") or 0), + str(deal.get("STAGE_ID") or "") + ) + next_stage = await self._next_stage(binding, deal) + if next_stage: + deal["NEXT_STAGE_ID"] = next_stage.stage_id + deal["NEXT_STAGE_NAME"] = next_stage.title + deal["NEXT_STAGE_IS_FINAL"] = next_stage.is_final + return deal + + async def prepare_stage_advance( + self, + binding: Binding, + deal_id: str + ) -> DealStageAdvance: + deal = await self._get_raw(binding, deal_id) + if not deal: + raise RuntimeError("Сделка не найдена") + self._ensure_responsible(binding, deal) + + next_stage = await self._next_stage(binding, deal) + if not next_stage: + raise RuntimeError("Сделка уже находится на финальной стадии") + + return DealStageAdvance( + deal_id=deal_id, + current_stage_id=str(deal.get("STAGE_ID") or ""), + target_stage_id=next_stage.stage_id, + target_stage_title=next_stage.title, + is_final=next_stage.is_final, + ) + + async def advance_stage( + self, + binding: Binding, + deal_id: str, + expected_stage_id: str, + expected_target_stage_id: str, + comment: str | None, + ) -> DealStageAdvance: + async with self._deal_lock(binding, deal_id): + current = await self._get_raw(binding, deal_id) + if not current: + raise RuntimeError("Сделка не найдена") + self._ensure_responsible(binding, current) + + current_stage_id = str(current.get("STAGE_ID") or "") + if current_stage_id != expected_stage_id: + raise DealStageConflict( + "Стадия уже изменилась. Обновите карточку сделки." + ) + + next_stage = await self._next_stage(binding, current) + if ( + not next_stage + or next_stage.stage_id != expected_target_stage_id + ): + raise DealStageConflict( + "Набор стадий изменился. Откройте сделку заново." + ) + + advance = DealStageAdvance( + deal_id=deal_id, + current_stage_id=current_stage_id, + target_stage_id=next_stage.stage_id, + target_stage_title=next_stage.title, + is_final=next_stage.is_final, + ) + # https://apidocs.bitrix24.ru/api-reference/crm/deals/crm-deal-update.html + await self.bitrix.call( + binding, + "crm.deal.update", + { + "id": deal_id, + "fields": {"STAGE_ID": next_stage.stage_id}, + "params": {"REGISTER_HISTORY_EVENT": "Y"}, + }, + ) + + updated = await self._get_raw(binding, deal_id) + if str((updated or {}).get( + "STAGE_ID") or "") != next_stage.stage_id: + raise DealStageConflict( + "Стадия изменилась одновременно с обновлением." + ) + + if comment: + try: + # Комментарий добавляется в таймлайн, не затирая COMMENTS. + await self.bitrix.call( + binding, + "crm.timeline.comment.add", + { + "fields": { + "ENTITY_ID": int(deal_id), + "ENTITY_TYPE": "deal", + "COMMENT": comment, + } + }, + ) + except Exception as error: + raise DealCommentSaveError(advance) from error + + return advance + + async def _next_stage( + self, + binding: Binding, + deal: dict + ) -> DealStage | None: + current_stage_id = str(deal.get("STAGE_ID") or "") + current_semantics = str( + deal.get("STAGE_SEMANTIC_ID") or "" + ).upper() + if current_semantics in {"S", "F"}: + return None + + stages = await self._stages( + binding, + int(deal.get("CATEGORY_ID") or 0), + ) + current_index = next( + ( + index + for index, stage in enumerate(stages) + if stage.stage_id == current_stage_id + ), + None, + ) + if current_index is None or stages[current_index].is_final: + return None + if current_index + 1 >= len(stages): + return None + return stages[current_index + 1] + + @staticmethod + def _ensure_responsible(binding: Binding, deal: dict) -> None: + responsible_id = str(deal.get("ASSIGNED_BY_ID") or "") + if responsible_id != str(binding.bitrix_user_id): + raise DealAdvanceForbidden( + "Переводить сделку может только ответственный за нее." + ) + + async def history( + self, + binding: Binding, + deal_id: str, + limit: int = 5 + ) -> list[dict]: + # https://apidocs.bitrix24.ru/api-reference/crm/crm-stage-history-list.html + data = await self.bitrix.call( + binding, + "crm.stagehistory.list", + { + "entityTypeId": 2, + "order": {"ID": "DESC"}, + "filter": {"OWNER_ID": int(deal_id)}, + "select": [ + "ID", + "TYPE_ID", + "CATEGORY_ID", + "STAGE_ID", + "CREATED_TIME" + ], + "start": 0 + } + ) + result = data.get("result") or {} + events = list(result.get("items") or [])[:limit] + for event in events: + event["STAGE_NAME"] = await self.stage_name( + binding, + int(event.get("CATEGORY_ID") or 0), + str(event.get("STAGE_ID") or "") + ) + return events + + async def remind_to_call( + self, + binding: Binding, + deal_id: str + ) -> datetime: + deadline = datetime.now(UTC) + timedelta(hours=1) + # https://apidocs.bitrix24.ru/api-reference/crm/timeline/activities/todo/crm-activity-todo-add.html + await self.bitrix.call( + binding, + "crm.activity.todo.add", + { + "ownerTypeId": 2, + "ownerId": int(deal_id), + "deadline": deadline.isoformat(), + "title": "Позвонить клиенту", + "description": f"Отложенный звонок по сделке #{deal_id}", + "responsibleId": binding.bitrix_user_id, + "pingOffsets": [0] + } + ) + return deadline + + async def take_to_work( + self, + binding: Binding, + deal_id: str, + expected_responsible_id: str + ) -> bool: + async with self._deal_lock(binding, deal_id): + current = await self._get_raw(binding, deal_id) + if not current: + raise RuntimeError("Сделка не найдена") + + responsible_id = str(current.get("ASSIGNED_BY_ID") or "") + target_id = str(binding.bitrix_user_id) + if responsible_id == target_id: + return False + if responsible_id != expected_responsible_id: + raise DealAssignmentConflict( + "Ответственный уже изменился. Обновите карточку сделки." + ) + + target_stage_id = await self._take_to_work_stage_id( + binding, + int(current.get("CATEGORY_ID") or 0) + ) + + # https://apidocs.bitrix24.com/api-reference/crm/deals/crm-deal-update.html + await self.bitrix.call( + binding, + "crm.deal.update", + { + "id": deal_id, + "fields": { + "ASSIGNED_BY_ID": binding.bitrix_user_id, + "STAGE_ID": target_stage_id + }, + "params": {"REGISTER_HISTORY_EVENT": "Y"} + }, + ) + + # REST Bitrix не поддерживает условный UPDATE, поэтому проверяем результат. + updated = await self._get_raw(binding, deal_id) + if str((updated or {}).get("ASSIGNED_BY_ID") or "") != target_id: + raise DealAssignmentConflict( + "Ответственный изменился одновременно с назначением." + ) + return True + + async def _take_to_work_stage_id( + self, + binding: Binding, + category_id: int + ) -> str: + stages = await self._stage_map(binding, category_id) + candidates = [self.take_to_work_stage_id] + if category_id and ":" not in self.take_to_work_stage_id: + candidates.append( + f"C{category_id}:{self.take_to_work_stage_id}" + ) + + target = next((item for item in candidates if item in stages), None) + if target: + return target + + raise RuntimeError( + "Стадия для взятия в работу " + f"{self.take_to_work_stage_id} не найдена в Битриксе" + ) + + async def _get_raw(self, binding: Binding, deal_id: str) -> dict | None: + data = await self.bitrix.call(binding, "crm.deal.get", {"id": deal_id}) + return data.get("result") or None + + @asynccontextmanager + async def _deal_lock( + self, + binding: Binding, + deal_id: str + ) -> AsyncGenerator[None]: + key = (binding.member_id, deal_id) + async with self._deal_locks_guard: + entry = self._deal_locks.get(key) + if entry is None: + entry = _DealLockEntry(asyncio.Lock()) + self._deal_locks[key] = entry + entry.users += 1 + + await entry.lock.acquire() + try: + yield + finally: + entry.lock.release() + async with self._deal_locks_guard: + entry.users -= 1 + if entry.users == 0: + self._deal_locks.pop(key, None) + + async def get_client_info(self, binding: Binding, deal: dict) -> ClientInfo: + contact = None + company = None + + contact_id = str(deal.get("CONTACT_ID") or "") + if contact_id: + # https://apidocs.bitrix24.com/api-reference/crm/contacts/crm-contact-get.html + contact = await self._entity(binding, "crm.contact.get", contact_id) + + company_id = str(deal.get("COMPANY_ID") or "") + if company_id: + # https://apidocs.bitrix24.com/api-reference/crm/companies/crm-company-get.html + company = await self._entity(binding, "crm.company.get", company_id) + + return ClientInfo( + name=self._contact_name(contact), + company=self._company_name(company), + phone=( + self._phone_from_entity(contact) or self._phone_from_entity( + company) + ) + ) + + async def get_source_name(self, binding: Binding, + source_id: str) -> str | None: + if not source_id: + return None + + try: + # https://apidocs.bitrix24.ru/api-reference/crm/status/crm-status-list.html + data = await self.bitrix.call( + binding, + "crm.status.list", + { + "filter": { + "ENTITY_ID": "SOURCE", + "STATUS_ID": source_id + } + } + ) + except Exception: + logger.exception("Failed to load Bitrix source name") + return None + + sources = data.get("result", []) + return str(sources[0].get("NAME") or "") or None if sources else None + + async def _entity(self, binding: Binding, method: str, + entity_id: str) -> dict: + data = await self.bitrix.call(binding, method, {"id": entity_id}) + return data.get("result", {}) or {} + + @staticmethod + def _phone_from_entity(entity: dict | None) -> str | None: + phones = (entity or {}).get("PHONE") or [] + return str(phones[0].get("VALUE") or "") or None if phones else None + + @staticmethod + def _contact_name(contact: dict | None) -> str | None: + if not contact: + return None + parts = [ + str(contact.get("LAST_NAME") or "").strip(), + str(contact.get("NAME") or "").strip(), + str(contact.get("SECOND_NAME") or "").strip() + ] + return " ".join(part for part in parts if part) or None + + @staticmethod + def _company_name(company: dict | None) -> str | None: + if not company: + return None + return str(company.get("TITLE") or "").strip() or None diff --git a/apps/bot/domain.py b/apps/bot/domain.py new file mode 100644 index 0000000..c2802bd --- /dev/null +++ b/apps/bot/domain.py @@ -0,0 +1,73 @@ +from dataclasses import dataclass +from datetime import datetime + +DEALS_PER_PAGE = 5 +MAX_DEAL_MESSAGE_LENGTH = 3900 + + +@dataclass(frozen=True) +class Binding: + member_id: str + domain: str + bitrix_user_id: int + telegram_user_id: int + + +@dataclass(frozen=True) +class OAuthCredentials: + member_id: str + domain: str + bitrix_user_id: int + access_token: bytes + refresh_token: bytes + expires_at: datetime + version: int + + +@dataclass(frozen=True) +class ClientInfo: + name: str | None = None + company: str | None = None + phone: str | None = None + + +@dataclass(frozen=True) +class DealStageFilter: + key: str + title: str + stage_id: str | None + assigned_to_viewer: bool = False + + +@dataclass(frozen=True) +class DealStage: + stage_id: str + title: str + semantics: str + + @property + def is_final(self) -> bool: + return self.semantics in {"S", "F"} + + +@dataclass(frozen=True) +class DealStageAdvance: + deal_id: str + current_stage_id: str + target_stage_id: str + target_stage_title: str + is_final: bool + + +@dataclass(frozen=True) +class DealPage: + deals: list[dict] + page: int + total_deals: int + total_pages: int + stage_filter: DealStageFilter + stage_filters: tuple[DealStageFilter, ...] + + @property + def has_next(self) -> bool: + return self.page + 1 < self.total_pages diff --git a/apps/bot/handlers.py b/apps/bot/handlers.py new file mode 100644 index 0000000..a3a7c6d --- /dev/null +++ b/apps/bot/handlers.py @@ -0,0 +1,531 @@ +import html +import logging + +import httpx +from aiogram import F, Router +from aiogram.exceptions import TelegramAPIError, TelegramBadRequest +from aiogram.filters import Command, CommandObject +from aiogram.fsm.context import FSMContext +from aiogram.fsm.state import State, StatesGroup +from aiogram.types import ( + CallbackQuery, + ForceReply, + InlineKeyboardMarkup, + Message, +) + +from .binding import BindingService +from .deals import ( + DealAssignmentConflict, + DealCommentSaveError, + DealService, +) +from .domain import Binding, DealPage +from .middleware import BindingRequiredMiddleware +from .presentation import DealFormatter, DealKeyboards + +logger = logging.getLogger(__name__) + + +class DealAdvanceStates(StatesGroup): + waiting_comment = State() + + +def normalize_comment(value: str | None) -> str | None: + comment = (value or "").strip() + if not comment: + return None + if all( + character.isspace() or character in "-‐‑‒–—―" + for character in comment + ): + return None + return comment + + +class StartBotHandlers: + """Справка и погашение одноразовой ссылки.""" + + def __init__(self, bindings: BindingService) -> None: + self.bindings = bindings + self.router = Router(name="start") + self.router.message.register(self.start, Command("start", "help")) + + async def start(self, message: Message, command: CommandObject) -> None: + payload = (command.args or "").strip() + # Если команда /start пришла с параметром bind_*, + # то это одноразовая ссылка для привязки. + if command.command == "start" and payload.startswith("bind_"): + await self.bind(message, payload.removeprefix("bind_")) + return + + # pyrefly: ignore [missing-attribute] + binding = await self.bindings.find(message.from_user.id) + if binding: + await message.answer( + "Аккаунт привязан. /deals покажет сделки, а /deal ID откроет карточку." + ) + else: + await message.answer( + "Открой приложение в Bitrix24 и нажми кнопку привязки Telegram." + ) + + async def bind(self, message: Message, token: str) -> None: + if message.chat.type != "private": + await message.answer( + "Привязку нужно открыть в личном чате с ботом.") + return + + binding = await self.bindings.consume( + token, + message.from_user.id, # pyrefly: ignore [missing-attribute] + message.chat.id + ) + if not binding: + await message.answer("Ссылка недействительна или уже использована.") + return + + await message.answer( + "Telegram успешно привязан к Bitrix24. Теперь доступна команда /deals." + ) + + +class DealBotHandlers: + """Команды и inline-кнопки для сделок.""" + + def __init__( + self, + service: DealService, + bindings: BindingService | None = None + ) -> None: + self.service = service + self.router = Router(name="deals") + self.router.message.register(self.deals, Command("leads", "deals")) + self.router.message.register(self.deal_by_command, + Command("lead", "deal")) + self.router.message.register( + self.cancel_stage_advance, + DealAdvanceStates.waiting_comment, + Command("cancel"), + ) + self.router.message.register( + self.advance_stage, + DealAdvanceStates.waiting_comment, + F.text, + ) + self.router.message.register( + self.require_stage_comment, + DealAdvanceStates.waiting_comment, + ) + self.router.callback_query.register( + self.deals_page, F.data.startswith("deals:page:") + ) + self.router.callback_query.register( + self.deal_by_button, F.data.startswith("deal:view:") + ) + self.router.callback_query.register( + self.assign_responsible, F.data.startswith("deal:assign:") + ) + self.router.callback_query.register( + self.request_stage_advance, + F.data.startswith("deal:advance:"), + ) + self.router.callback_query.register( + self.show_history, F.data.startswith("deal:history:") + ) + self.router.callback_query.register( + self.remind_to_call, F.data.startswith("deal:remind:") + ) + + # Подвязываем middleware, который проверяет наличие привязки к Битриксу. + if bindings: + middleware = BindingRequiredMiddleware(bindings) + self.router.message.middleware(middleware) + self.router.callback_query.middleware(middleware) + + async def deals(self, message: Message, binding: Binding) -> None: + try: + await self.send_deals_page(message, binding, stage_key="new", + page=0) + except Exception as error: + await self.answer_error(message, error) + + async def deal_by_command( + self, + message: Message, + command: CommandObject, + binding: Binding + ) -> None: + deal_id = (command.args or "").strip() + if not deal_id.isdigit(): + await message.answer( + "Укажи ID сделки: /deal 123", + parse_mode="HTML", + ) + return + + try: + await self.send_deal(message, binding, deal_id) + except Exception as error: + await self.answer_error(message, error) + + async def deals_page(self, callback: CallbackQuery, + binding: Binding) -> None: + parts = (callback.data or "deals:page:new:0").split(":") + stage_key = parts[2] if len(parts) > 2 else "new" + page = int(parts[3]) if len(parts) > 3 else 0 + + try: + deal_page = await self.service.list_by_stage(binding, stage_key, + page) + changed = await self.edit_callback_message( + callback, + self.page_text(deal_page), + DealKeyboards.deals_page( + deal_page.deals, + deal_page.stage_filters, + deal_page.stage_filter.key, + deal_page.page, + deal_page.has_next + ), + unchanged_text="Список уже актуален." + ) + if changed: + await callback.answer() + + except Exception as error: + await self.answer_callback_error(callback, error) + + async def deal_by_button(self, callback: CallbackQuery, + binding: Binding) -> None: + deal_id = (callback.data or "").split(":")[-1] + try: + deal = await self.service.get(binding, deal_id) + if not deal: + await callback.answer("Сделка не найдена.", show_alert=True) + return + + changed = await self.edit_callback_message( + callback, + DealFormatter.deal_details(deal), + DealKeyboards.deal_card( + deal, + binding.bitrix_user_id, + ), + unchanged_text="Карточка уже открыта." + ) + if changed: + await callback.answer() + + except Exception as error: + await self.answer_callback_error(callback, error) + + async def assign_responsible( + self, callback: CallbackQuery, binding: Binding + ) -> None: + parts = (callback.data or "").split(":") + if len(parts) != 4: + await callback.answer( + "Карточка устарела. Откройте сделку заново.", + show_alert=True + ) + return + + deal_id, expected_responsible_id = parts[2], parts[3] + try: + assigned = await self.service.take_to_work( + binding, + deal_id, + expected_responsible_id + ) + deal = await self.service.get(binding, deal_id) + if not deal: + await callback.answer( + "Сделка обновлена, но повторно не найдена.", + show_alert=True + ) + return + + changed = await self.edit_callback_message( + callback, + DealFormatter.deal_details(deal), + DealKeyboards.deal_card( + deal, + binding.bitrix_user_id, + ), + unchanged_text="Сделка уже отображается актуально." + ) + if changed: + text = ( + "Сделка переведена в работу." + if assigned + else "Вы уже ответственный за эту сделку." + ) + await callback.answer(text) + + except DealAssignmentConflict as error: + await callback.answer(str(error), show_alert=True) + except Exception as error: + await self.answer_callback_error(callback, error) + + async def request_stage_advance( + self, + callback: CallbackQuery, + binding: Binding, + state: FSMContext, + ) -> None: + deal_id = (callback.data or "").split(":")[-1] + if not deal_id.isdigit() or not isinstance(callback.message, Message): + await callback.answer( + "Карточка устарела. Откройте сделку заново.", + show_alert=True, + ) + return + + try: + advance = await self.service.prepare_stage_advance( + binding, + deal_id, + ) + await state.set_state(DealAdvanceStates.waiting_comment) + await state.set_data( + { + "deal_id": advance.deal_id, + "current_stage_id": advance.current_stage_id, + "target_stage_id": advance.target_stage_id, + "target_stage_title": advance.target_stage_title, + "is_final": advance.is_final, + } + ) + + final_note = " (финальная)" if advance.is_final else "" + await callback.message.answer( + ( + f"Следующая стадия: " + f"{html.escape(advance.target_stage_title)}" + f"{final_note}.\n" + "Введите комментарий одним сообщением. " + "Чтобы продолжить без комментария, отправьте " + "-. Для отмены — /cancel." + ), + parse_mode="HTML", + reply_markup=ForceReply( + selective=True, + input_field_placeholder="Комментарий или -", + ), + ) + await callback.answer("Жду комментарий.") + except Exception as error: + await self.answer_callback_error(callback, error) + + async def advance_stage( + self, + message: Message, + binding: Binding, + state: FSMContext, + ) -> None: + data = await state.get_data() + await state.clear() + deal_id = str(data.get("deal_id") or "") + expected_stage_id = str(data.get("current_stage_id") or "") + target_stage_id = str(data.get("target_stage_id") or "") + if not deal_id or not expected_stage_id or not target_stage_id: + await message.answer( + "Запрос устарел. Откройте карточку сделки заново." + ) + return + + comment = normalize_comment(message.text) + try: + advance = await self.service.advance_stage( + binding, + deal_id, + expected_stage_id, + target_stage_id, + comment, + ) + if advance.is_final: + result_text = ( + "Сделка переведена на финальную стадию " + f"«{html.escape(advance.target_stage_title)}»." + ) + else: + result_text = ( + "Сделка переведена на стадию " + f"«{html.escape(advance.target_stage_title)}»." + ) + if comment: + result_text += " Комментарий добавлен в таймлайн." + else: + result_text += " Переход выполнен без комментария." + await message.answer(result_text, parse_mode="HTML") + await self.send_deal(message, binding, deal_id) + except DealCommentSaveError: + await message.answer( + ( + "Стадия изменена, но комментарий не удалось " + "сохранить в Битриксе." + ) + ) + try: + await self.send_deal(message, binding, deal_id) + except Exception: + logger.exception( + "Failed to refresh deal after comment save error" + ) + except Exception as error: + await self.answer_error(message, error) + + @staticmethod + async def cancel_stage_advance( + message: Message, + state: FSMContext, + ) -> None: + await state.clear() + await message.answer("Переход на следующую стадию отменен.") + + @staticmethod + async def require_stage_comment(message: Message) -> None: + await message.answer( + "Отправьте комментарий текстом или прочерк " + "-, чтобы продолжить без него.", + parse_mode="HTML", + ) + + async def show_history(self, callback: CallbackQuery, + binding: Binding) -> None: + deal_id = (callback.data or "").split(":")[-1] + try: + events = await self.service.history(binding, deal_id) + changed = await self.edit_callback_message( + callback, + DealFormatter.deal_history(deal_id, events), + DealKeyboards.deal_history(deal_id), + unchanged_text="История уже открыта.", + ) + if changed: + await callback.answer() + except Exception as error: + await self.answer_callback_error(callback, error) + + async def remind_to_call( + self, + callback: CallbackQuery, + binding: Binding + ) -> None: + deal_id = (callback.data or "").split(":")[-1] + try: + await self.service.remind_to_call(binding, deal_id) + await callback.answer( + "Напоминание создано в Битриксе на час позже.", + show_alert=True + ) + except Exception as error: + await self.answer_callback_error(callback, error) + + async def send_deals_page( + self, + message: Message, + binding: Binding, + stage_key: str, + page: int + ) -> None: + deal_page = await self.service.list_by_stage(binding, stage_key, page) + await message.answer( + self.page_text(deal_page), + reply_markup=DealKeyboards.deals_page( + deal_page.deals, + deal_page.stage_filters, + deal_page.stage_filter.key, + deal_page.page, + deal_page.has_next + ), + parse_mode="HTML" + ) + + async def send_deal(self, message: Message, binding: Binding, + deal_id: str) -> None: + deal = await self.service.get(binding, deal_id) + if not deal: + await message.answer("Сделка не найдена.") + return + + await message.answer( + DealFormatter.deal_details(deal), + reply_markup=DealKeyboards.deal_card( + deal, + binding.bitrix_user_id + ), + parse_mode="HTML" + ) + + @staticmethod + def page_text(deal_page: DealPage) -> str: + title = DealFormatter.list_title( + deal_page.stage_filter, + deal_page.page, + deal_page.total_deals, + deal_page.total_pages + ) + if not deal_page.deals: + return f"{title}\n\nСделки в этом фильтре не найдены." + return title + + @staticmethod + async def edit_callback_message( + callback: CallbackQuery, + text: str, + reply_markup: InlineKeyboardMarkup, + unchanged_text: str + ) -> bool: + if not isinstance(callback.message, Message): + await callback.answer("Не удалось обновить сообщение.", + show_alert=True) + return False + + try: + await callback.message.edit_text( + text, + reply_markup=reply_markup, + parse_mode="HTML" + ) + return True + + except TelegramBadRequest as error: + if "message is not modified" in str(error).lower(): + await callback.answer(unchanged_text) + return False + raise + + @staticmethod + async def answer_error(message: Message, error: Exception) -> None: + if isinstance(error, httpx.HTTPStatusError): + logger.exception("Bitrix HTTP error") + await message.answer( + f"Ошибка HTTP Битрикс24: {error.response.status_code}") + elif isinstance(error, httpx.RequestError): + logger.exception("Bitrix connection error") + await message.answer("Не удалось подключиться к Битрикс24.") + elif isinstance(error, TelegramAPIError): + logger.exception("Telegram API error") + await message.answer("Telegram не смог выполнить действие.") + elif isinstance(error, RuntimeError): + logger.exception("Runtime error") + await message.answer(f"Ошибка: {html.escape(str(error))}") + else: + logger.exception("Unexpected bot error") + await message.answer("Произошла неизвестная ошибка.") + + @classmethod + async def answer_callback_error( + cls, + callback: CallbackQuery, + error: Exception + ) -> None: + if isinstance(callback.message, Message): + await cls.answer_error(callback.message, error) + try: + await callback.answer("Не удалось выполнить действие.", + show_alert=True) + except TelegramAPIError: + logger.exception("Failed to answer callback after error") diff --git a/apps/bot/main.py b/apps/bot/main.py new file mode 100644 index 0000000..5dc9162 --- /dev/null +++ b/apps/bot/main.py @@ -0,0 +1,57 @@ +import asyncio +import logging + +from aiogram import Bot, Dispatcher +from dotenv import load_dotenv + +from .binding import BindingService, BotBindingRepository +from .bitrix import BitrixClient +from .config import BotConfig +from .crypto import TokenCipher +from .database import BotDatabase +from .deals import DealService +from .handlers import DealBotHandlers, StartBotHandlers +from .oauth import BotOAuthRepository + + +async def run() -> None: + load_dotenv() + logging.basicConfig(level=logging.INFO) + logging.getLogger("httpx").setLevel(logging.WARNING) + logging.getLogger("httpcore").setLevel(logging.WARNING) + config = BotConfig.from_env() + + database = BotDatabase(config.database_url) + await database.open() + + bitrix = BitrixClient( + BotOAuthRepository(database), + TokenCipher(config.token_encryption_key), + config.bitrix_client_id, + config.bitrix_client_secret, + config.oauth_token_url, + ) + + bindings = BindingService(BotBindingRepository(database)) + deals = DealService( + bitrix, + config.take_to_work_stage_id, + ) + + dispatcher = Dispatcher() + dispatcher.include_router(StartBotHandlers(bindings).router) + dispatcher.include_router(DealBotHandlers(deals, bindings).router) + + try: + await dispatcher.start_polling(Bot(token=config.bot_token)) + finally: + await bitrix.close() + await database.close() + + +def main() -> None: + asyncio.run(run()) + + +if __name__ == "__main__": + main() diff --git a/apps/bot/middleware.py b/apps/bot/middleware.py new file mode 100644 index 0000000..bdb7850 --- /dev/null +++ b/apps/bot/middleware.py @@ -0,0 +1,34 @@ +from collections.abc import Awaitable, Callable +from typing import Any + +from aiogram import BaseMiddleware +from aiogram.types import CallbackQuery, Message, TelegramObject + +from .binding import BindingService + + +class BindingRequiredMiddleware(BaseMiddleware): + """Не пропускает CRM-команды до привязки аккаунта.""" + + def __init__(self, bindings: BindingService) -> None: + self.bindings = bindings + + async def __call__( + self, + handler: Callable[[TelegramObject, dict[str, Any]], Awaitable[Any]], + event: TelegramObject, + data: dict[str, Any] + ) -> Any: + user = data.get("event_from_user") + if user: + binding = await self.bindings.find(user.id) + if binding: + data["binding"] = binding + return await handler(event, data) + + text = "Сначала привяжи Telegram через приложение в Bitrix24." + if isinstance(event, CallbackQuery): + await event.answer(text, show_alert=True) + elif isinstance(event, Message): + await event.answer(text) + return None diff --git a/apps/bot/oauth.py b/apps/bot/oauth.py new file mode 100644 index 0000000..b98c93c --- /dev/null +++ b/apps/bot/oauth.py @@ -0,0 +1,84 @@ +from datetime import datetime + +from .database import BotDatabase +from .domain import Binding, OAuthCredentials + + +class BotOAuthRepository: + """Транзакционные операции с OAuth-данными Битрикса. + Обертка над хранимыми функциями БД.""" + + def __init__(self, database: BotDatabase) -> None: + self.database = database + + async def get(self, binding: Binding) -> OAuthCredentials | None: + async with self.database.transaction() as connection: + cursor = await connection.execute( + "SELECT * FROM oauth.get_credentials_v1(%s, %s)", + (binding.member_id, binding.bitrix_user_id) + ) + row = await cursor.fetchone() + # pyrefly: ignore [bad-argument-type] + return self._credentials(row) if row else None + + async def claim_refresh(self, credentials: OAuthCredentials) -> bool: + async with self.database.transaction() as connection: + cursor = await connection.execute( + "SELECT oauth.claim_refresh_v1(%s, %s, %s)", + ( + credentials.member_id, + credentials.bitrix_user_id, + credentials.version + ) + ) + row = await cursor.fetchone() + + # pyrefly: ignore [missing-attribute] + return bool(row and next(iter(row.values()))) + + async def finish_refresh( + self, + credentials: OAuthCredentials, + access_token: bytes, + refresh_token: bytes, + expires_at: datetime + ) -> bool: + async with self.database.transaction() as connection: + cursor = await connection.execute( + "SELECT oauth.finish_refresh_v1(%s, %s, %s, %s, %s, %s)", + ( + credentials.member_id, + credentials.bitrix_user_id, + credentials.version, + access_token, + refresh_token, + expires_at + ) + ) + row = await cursor.fetchone() + + # pyrefly: ignore [missing-attribute] + return bool(row and next(iter(row.values()))) + + async def release_refresh(self, credentials: OAuthCredentials) -> None: + async with self.database.transaction() as connection: + await connection.execute( + "SELECT oauth.release_refresh_v1(%s, %s, %s)", + ( + credentials.member_id, + credentials.bitrix_user_id, + credentials.version + ) + ) + + @staticmethod + def _credentials(row: dict) -> OAuthCredentials: + return OAuthCredentials( + member_id=str(row["member_id"]), + domain=str(row["domain"]), + bitrix_user_id=int(row["bitrix_user_id"]), + access_token=bytes(row["access_token"]), + refresh_token=bytes(row["refresh_token"]), + expires_at=row["expires_at"], + version=int(row["version"]) + ) diff --git a/apps/bot/presentation.py b/apps/bot/presentation.py new file mode 100644 index 0000000..2577a51 --- /dev/null +++ b/apps/bot/presentation.py @@ -0,0 +1,272 @@ +import decimal +import html +from typing import Any + +from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup + +from .domain import ( + MAX_DEAL_MESSAGE_LENGTH, + DealStageFilter, +) + + +class DealFormatter: + """Тексты карточек и списков Telegram.""" + + @staticmethod + def truncate(text: str, limit: int) -> str: + if len(text) <= limit: + return text + return text[: limit - 3].rstrip() + "..." + + @staticmethod + def list_title( + stage_filter: DealStageFilter, + page: int, + total_deals: int, + total_pages: int + ) -> str: + title = html.escape(stage_filter.title) + return ( + f"Сделки: {title}\n" + f"Всего сделок: {total_deals}\n" + f"Страница {page + 1} из {total_pages}" + ) + + @staticmethod + def money(value: Any, currency: Any) -> str: + try: + amount = decimal.Decimal(str(value or "0")) + formatted = f"{amount:,.2f}".replace(",", " ") + except decimal.InvalidOperation: + formatted = html.escape(str(value or "0")) + return f"{formatted} {html.escape(str(currency or ''))}".strip() + + @classmethod + def deal_summary(cls, deal: dict) -> str: + summary = " · ".join( + ( + f"#{deal.get('ID', '-')}", + str(deal.get("TITLE") or "Без названия"), + cls.money(deal.get("OPPORTUNITY"), deal.get("CURRENCY_ID")) + ) + ) + return cls.truncate(summary, 60) + + @classmethod + def deal_details(cls, deal: dict) -> str: + deal_id = html.escape(str(deal.get("ID", "-"))) + title = html.escape(str(deal.get("TITLE") or "Без названия")) + stage = html.escape( + str(deal.get("STAGE_NAME") or deal.get("STAGE_ID") or "-") + ) + source = html.escape( + str(deal.get("SOURCE_NAME") or deal.get("SOURCE_ID") or "-") + ) + assigned = html.escape(str(deal.get("ASSIGNED_BY_ID") or "не назначен")) + date = html.escape(str(deal.get("DATE_CREATE") or "-")) + client = html.escape(str(deal.get("CLIENT_NAME") or "не указан")) + company = html.escape(str(deal.get("CLIENT_COMPANY") or "")) + phone = html.escape(str(deal.get("CLIENT_PHONE") or "не найден")) + comments = html.escape(str(deal.get("COMMENTS") or "")).strip() + + lines = [ + f"Сделка #{deal_id}", + f"{title}", + "", + f"Клиент: {client}", + *([f"Компания: {company}"] if company else []), + f"Телефон клиента: {phone}", + f"Стадия: {stage}", + f"Источник сделки: {source}", + f"Сумма: {cls.money(deal.get('OPPORTUNITY'), deal.get('CURRENCY_ID'))}", + f"Ответственный: {assigned}", + f"Дата создания: {date}" + ] + if comments: + lines.extend(["", f"Комментарий:\n{comments}"]) + + return cls.truncate("\n".join(lines), MAX_DEAL_MESSAGE_LENGTH) + + @classmethod + def deal_history(cls, deal_id: str, events: list[dict]) -> str: + lines = [f"История сделки #{html.escape(deal_id)}"] + if not events: + lines.extend(["", "Изменения стадий пока не найдены."]) + return "\n".join(lines) + + event_names = { + "1": "Сделка создана", + "2": "Переход на стадию", + "3": "Переход на финальную стадию", + "5": "Изменение воронки" + } + for event in events: + date = html.escape( + str(event.get("CREATED_TIME") or "дата не указана")) + event_type = str(event.get("TYPE_ID") or "") + name = event_names.get(event_type, "Изменение стадии") + stage = html.escape(str( + event.get("STAGE_NAME") + or event.get("STAGE_ID") + or "не указана" + )) + lines.extend( + [ + "", + f"• {name}", + f" Стадия: {stage}", + f" {date}" + ] + ) + + return cls.truncate("\n".join(lines), MAX_DEAL_MESSAGE_LENGTH) + + +class DealKeyboards: + @staticmethod + def deals_page( + deals: list[dict], + stage_filters: tuple[DealStageFilter, ...], + stage_key: str, + page: int, + has_next: bool + ) -> InlineKeyboardMarkup: + filter_buttons = [ + InlineKeyboardButton( + text=("✓ " if stage.key == stage_key else "") + stage.title, + callback_data=f"deals:page:{stage.key}:0" + ) + for stage in stage_filters + ] + rows = [ + [ + InlineKeyboardButton( + text=DealFormatter.deal_summary(deal), + callback_data=f"deal:view:{deal['ID']}" + ) + ] + for deal in deals + ] + + navigation = [] + if page > 0: + navigation.append( + InlineKeyboardButton( + text="Назад", + callback_data=f"deals:page:{stage_key}:{page - 1}" + ) + ) + if has_next: + navigation.append( + InlineKeyboardButton( + text="Вперед", + callback_data=f"deals:page:{stage_key}:{page + 1}" + ) + ) + if navigation: + rows.append(navigation) + + rows.append( + [ + InlineKeyboardButton( + text="Обновить", + callback_data=f"deals:page:{stage_key}:{page}" + ) + ] + ) + rows.extend( + filter_buttons[index: index + 2] + for index in range(0, len(filter_buttons), 2) + ) + + return InlineKeyboardMarkup(inline_keyboard=rows) + + @staticmethod + def deal_card( + deal: dict, + viewer_bitrix_user_id: int + ) -> InlineKeyboardMarkup: + deal_id = str(deal["ID"]) + responsible_id = str(deal.get("ASSIGNED_BY_ID") or "") + rows = [] + + if ( + str(deal.get("IS_NEW") or "").upper() == "Y" + and responsible_id != str(viewer_bitrix_user_id) + ): + rows.append( + [ + InlineKeyboardButton( + text="Стать ответственным и взять в работу", + callback_data=( + f"deal:assign:{deal_id}:{responsible_id}") + ) + ] + ) + + if ( + responsible_id == str(viewer_bitrix_user_id) + and deal.get("NEXT_STAGE_ID") + ): + next_stage_name = DealFormatter.truncate( + str(deal.get("NEXT_STAGE_NAME") or "следующая стадия"), + 42, + ) + if deal.get("NEXT_STAGE_IS_FINAL"): + button_text = f"Завершить: {next_stage_name}" + else: + button_text = f"Следующая стадия: {next_stage_name}" + rows.append( + [ + InlineKeyboardButton( + text=button_text, + callback_data=f"deal:advance:{deal_id}", + ) + ] + ) + + rows.append( + [ + InlineKeyboardButton( + text="Позвонить позже", + callback_data=f"deal:remind:{deal_id}" + ) + ] + ) + rows.append( + [ + InlineKeyboardButton( + text="История изменений", + callback_data=f"deal:history:{deal_id}" + ) + ] + ) + rows.append( + [ + InlineKeyboardButton( + text="К списку сделок", + callback_data="deals:page:new:0" + ) + ] + ) + return InlineKeyboardMarkup(inline_keyboard=rows) + + @staticmethod + def deal_history(deal_id: str) -> InlineKeyboardMarkup: + return InlineKeyboardMarkup( + inline_keyboard=[ + [ + InlineKeyboardButton( + text="К сделке", + callback_data=f"deal:view:{deal_id}" + ) + ], + [ + InlineKeyboardButton( + text="К списку сделок", + callback_data="deals:page:new:0" + ) + ], + ] + ) diff --git a/apps/site/__init__.py b/apps/site/__init__.py new file mode 100644 index 0000000..b94a1e8 --- /dev/null +++ b/apps/site/__init__.py @@ -0,0 +1,3 @@ +from .app import create_app + +__all__ = ["create_app"] diff --git a/apps/site/app.py b/apps/site/app.py new file mode 100644 index 0000000..f56e24a --- /dev/null +++ b/apps/site/app.py @@ -0,0 +1,69 @@ +import atexit +import logging + +import httpx +from flask import Flask, jsonify +from dotenv import load_dotenv +from werkzeug.middleware.proxy_fix import ProxyFix + +from .binding import BindingService, SiteBindingRepository +from .bitrix import BitrixAuthError, BitrixClient +from .config import SiteConfig +from .crypto import TokenCipher +from .database import SiteDatabase +from .routes import SiteInputError, create_blueprint + +logger = logging.getLogger(__name__) + + +def create_app() -> Flask: + load_dotenv() + logging.getLogger("httpx").setLevel(logging.WARNING) + logging.getLogger("httpcore").setLevel(logging.WARNING) + config = SiteConfig.from_env() + app = Flask(__name__) + app.config["PUBLIC_BASE_URL"] = config.public_base_url + + if config.trust_proxy: + # Используется конфигурация, при которой снаружи контейнера находится + # обратный прокси nginx. + app.wsgi_app = ProxyFix( + app.wsgi_app, + x_for=1, + x_proto=1, + x_host=1, + ) + + database = SiteDatabase(config.database_url) + database.open() + bitrix = BitrixClient( + config.bitrix_client_id, + config.bitrix_client_secret, + config.oauth_token_url, + ) + bindings = BindingService( + SiteBindingRepository(database), + TokenCipher(config.token_encryption_key), + config.bot_username, + config.binding_ttl_seconds, + ) + + app.extensions["database"] = database + app.register_blueprint(create_blueprint(bindings, bitrix)) + atexit.register(database.close) + atexit.register(bitrix.close) + + @app.errorhandler(SiteInputError) + def input_error(error: SiteInputError): + return jsonify(error=str(error)), 400 + + @app.errorhandler(BitrixAuthError) + def auth_error(error: BitrixAuthError): + return jsonify(error=str(error)), 403 + + @app.errorhandler(httpx.HTTPError) + def bitrix_error(error: httpx.HTTPError): + logger.exception("Bitrix request failed") + return jsonify(error="Не удалось проверить пользователя Bitrix24"), 502 + + return app diff --git a/apps/site/binding.py b/apps/site/binding.py new file mode 100644 index 0000000..41dd274 --- /dev/null +++ b/apps/site/binding.py @@ -0,0 +1,98 @@ +import hashlib +import secrets +from dataclasses import dataclass +from datetime import UTC, datetime, timedelta + +from .database import SiteDatabase +from .crypto import TokenCipher + + +def hash_token(token: str) -> bytes: + return hashlib.sha256(token.encode("utf-8")).digest() + + +@dataclass(frozen=True) +class BindingLink: + url: str + expires_at: datetime + + +class SiteBindingRepository: + """Доступ сайта только к функциям выпуска токенов.""" + + def __init__(self, database: SiteDatabase) -> None: + self.database = database + + def issue( + self, + member_id: str, + domain: str, + bitrix_user_id: int, + token_hash: bytes, + token_expires_at: datetime, + access_token: bytes, + refresh_token: bytes, + oauth_expires_at: datetime, + ) -> None: + """Сохраняет в БД информацию о токене, выданном пользователю.""" + query = """ + SELECT * + FROM binding.issue_v1(%s, %s, %s, %s, %s, %s, %s, %s) \ + """ + with self.database.transaction() as connection: + connection.execute( + query, + ( + member_id, + domain, + bitrix_user_id, + token_hash, + token_expires_at, + access_token, + refresh_token, + oauth_expires_at, + ), + ).fetchone() + + +class BindingService: + def __init__( + self, + repository: SiteBindingRepository, + cipher: TokenCipher, + bot_username: str, + ttl_seconds: int, + ) -> None: + self.repository = repository + self.cipher = cipher + self.bot_username = bot_username + self.ttl_seconds = ttl_seconds + + def issue( + self, + member_id: str, + domain: str, + bitrix_user_id: int, + access_token: str, + refresh_token: str, + auth_expires_seconds: int, + ) -> BindingLink: + """Выдает ссылку для привязки аккаунта.""" + token = secrets.token_urlsafe(32) + now = datetime.now(UTC) + token_expires_at = now + timedelta(seconds=self.ttl_seconds) + oauth_expires_at = now + timedelta(seconds=auth_expires_seconds) + self.repository.issue( + member_id, + domain, + bitrix_user_id, + hash_token(token), + token_expires_at, + self.cipher.encrypt(access_token), + self.cipher.encrypt(refresh_token), + oauth_expires_at, + ) + return BindingLink( + url=f"https://t.me/{self.bot_username}?start=bind_{token}", + expires_at=token_expires_at, + ) diff --git a/apps/site/bitrix.py b/apps/site/bitrix.py new file mode 100644 index 0000000..1c25509 --- /dev/null +++ b/apps/site/bitrix.py @@ -0,0 +1,118 @@ +from dataclasses import dataclass +from urllib.parse import urlparse + +import httpx + + +class BitrixAuthError(RuntimeError): + pass + + +@dataclass(frozen=True) +class BitrixUser: + id: int + name: str + + +@dataclass(frozen=True) +class BitrixAuth: + member_id: str + domain: str + access_token: str + refresh_token: str + expires_in: int + user: BitrixUser + + +class BitrixClient: + """Получает доверенный OAuth-контекст и проверяет пользователя.""" + + def __init__( + self, + client_id: str, + client_secret: str, + oauth_token_url: str, + client: httpx.Client | None = None, + ) -> None: + self.client_id = client_id + self.client_secret = client_secret + self.oauth_token_url = oauth_token_url + # Клиент передается как внешняя зависимость для модульного тестирования. + self._client = client or httpx.Client(timeout=15) + # Соответственно, если клиент внешний, то класс + # этим ресурсом не управляет. + self._owns_client = client is None + + def authorize(self, refresh_token: str) -> BitrixAuth: + try: + # Обмениваем рефреш-токен на новую пару токенов. + # https://apidocs.bitrix24.com/settings/oauth/auto-renewal.html + # https://apidocs.bitrix24.com/settings/oauth/simple-way.html + response = self._client.get( + self.oauth_token_url, + params={ + "grant_type": "refresh_token", + "client_id": self.client_id, + "client_secret": self.client_secret, + "refresh_token": refresh_token, + }, + ) + response.raise_for_status() + except httpx.HTTPError: + # URL запроса содержит секреты, поэтому не пробрасываем его выше. + raise BitrixAuthError("Не удалось обновить OAuth-токен") from None + + data = response.json() + if "error" in data: + raise BitrixAuthError( + str(data.get("error_description") or data["error"])) + + # Получаем эндпоинт, с которым связаны наши токены. + endpoint = str(data.get("client_endpoint") or "") + parsed_endpoint = urlparse(endpoint) + if parsed_endpoint.scheme != "https" or not parsed_endpoint.hostname: + raise BitrixAuthError("Bitrix вернул некорректный REST endpoint") + + # Сохраняем токен доступа и проверяем пользователя. + access_token = str(data["access_token"]) + user = self._current_user(endpoint, access_token) + expected_user_id = data.get("user_id") + if expected_user_id is not None and user.id != int(expected_user_id): + raise BitrixAuthError( + "OAuth-токен принадлежит другому пользователю") + + return BitrixAuth( + member_id=str(data["member_id"]), + domain=parsed_endpoint.hostname.lower(), + access_token=access_token, + refresh_token=str(data["refresh_token"]), + expires_in=int(data.get("expires_in", 3600)), + user=user, + ) + + def _current_user(self, endpoint: str, access_token: str) -> BitrixUser: + """Получение информации о пользователе для проверки работоспособности.""" + # https://apidocs.bitrix24.com/api-reference/user/user-current.html + response = self._client.post( + endpoint.rstrip("/") + "/user.current.json", + data={"auth": access_token}, + ) + response.raise_for_status() + data = response.json() + if "error" in data or not data.get("result"): + raise BitrixAuthError("Bitrix не подтвердил текущего пользователя") + + user = data["result"] + name = " ".join( + part + for part in ( + str(user.get("NAME") or "").strip(), + str(user.get("LAST_NAME") or "").strip(), + ) + if part + ) + return BitrixUser(id=int(user["ID"]), name=name or f"ID {user['ID']}") + + def close(self) -> None: + if self._owns_client: + self._client.close() diff --git a/apps/site/config.py b/apps/site/config.py new file mode 100644 index 0000000..bb9fcd0 --- /dev/null +++ b/apps/site/config.py @@ -0,0 +1,57 @@ +import os +from dataclasses import dataclass +from urllib.parse import urlparse + + +def _required(name: str) -> str: + value = os.getenv(name) + if not value: + raise RuntimeError(f"{name} is not set") + return value + + +def _as_bool(value: str | None, default: bool = False) -> bool: + if value is None: + return default + return value.lower() in {"1", "true", "yes", "on"} + + +@dataclass(frozen=True) +class SiteConfig: + """Настройки HTTP-приложения.""" + + database_url: str + public_base_url: str + bot_username: str + token_encryption_key: str + bitrix_client_id: str + bitrix_client_secret: str + oauth_token_url: str + binding_ttl_seconds: int = 600 + trust_proxy: bool = True + + @classmethod + def from_env(cls) -> "SiteConfig": + public_base_url = _required("PUBLIC_BASE_URL").rstrip("/") + parsed_url = urlparse(public_base_url) + if parsed_url.scheme not in {"http", "https"} or not parsed_url.netloc: + raise RuntimeError("PUBLIC_BASE_URL must be an absolute URL") + + ttl = int(os.getenv("BINDING_TOKEN_TTL_SECONDS", "600")) + if not 60 <= ttl <= 3600: + raise RuntimeError("BINDING_TOKEN_TTL_SECONDS must be 60..3600") + + return cls( + database_url=_required("DATABASE_URL"), + public_base_url=public_base_url, + bot_username=_required("BOT_USERNAME").lstrip("@"), + token_encryption_key=_required("TOKEN_ENCRYPTION_KEY"), + bitrix_client_id=_required("BITRIX_CLIENT_ID"), + bitrix_client_secret=_required("BITRIX_CLIENT_SECRET"), + oauth_token_url=os.getenv( + "BITRIX_OAUTH_TOKEN_URL", + "https://oauth.bitrix.info/oauth/token/", + ), + binding_ttl_seconds=ttl, + trust_proxy=_as_bool(os.getenv("TRUST_PROXY"), default=True), + ) diff --git a/apps/site/crypto.py b/apps/site/crypto.py new file mode 100644 index 0000000..e858e3a --- /dev/null +++ b/apps/site/crypto.py @@ -0,0 +1,14 @@ +from cryptography.fernet import Fernet + + +class TokenCipher: + """Шифрует OAuth-токены перед передачей в БД.""" + + def __init__(self, key: str) -> None: + try: + self._fernet = Fernet(key.encode("ascii")) + except (ValueError, UnicodeEncodeError) as error: + raise RuntimeError("TOKEN_ENCRYPTION_KEY is invalid") from error + + def encrypt(self, value: str) -> bytes: + return self._fernet.encrypt(value.encode("utf-8")) diff --git a/apps/site/database.py b/apps/site/database.py new file mode 100644 index 0000000..12b6839 --- /dev/null +++ b/apps/site/database.py @@ -0,0 +1,39 @@ +from collections.abc import Generator +from contextlib import contextmanager + +from psycopg import Connection +from psycopg.rows import dict_row +from psycopg_pool import ConnectionPool + + +class SiteDatabase: + """Данный класс представляет собой обертку над пулом соединений + с базой данных PostgreSQL.""" + + def __init__(self, database_url: str) -> None: + # Пул может содержать в себе максимум 5 соединений. + self.pool = ConnectionPool( + conninfo=database_url, + min_size=1, + max_size=5, + open=False, + # Фабрика для представления строк БД как словарей. + kwargs={"row_factory": dict_row}, + ) + + def open(self) -> None: + self.pool.open(wait=True) + + def close(self) -> None: + self.pool.close() + + @contextmanager + def transaction(self) -> Generator[Connection]: + with self.pool.connection() as connection: + with connection.transaction(): + yield connection + + def ping(self) -> bool: + """Простая проверка подключения к БД.""" + with self.pool.connection() as connection: + return connection.execute("SELECT 1").fetchone() is not None diff --git a/apps/site/gunicorn.conf.py b/apps/site/gunicorn.conf.py new file mode 100644 index 0000000..cc06043 --- /dev/null +++ b/apps/site/gunicorn.conf.py @@ -0,0 +1,7 @@ +import os + +bind = f"{os.getenv('SITE_HOST', '0.0.0.0')}:{os.getenv('SITE_PORT', '8000')}" +workers = int(os.getenv("SITE_WORKERS", "2")) +accesslog = "-" +errorlog = "-" +timeout = 30 diff --git a/apps/site/routes.py b/apps/site/routes.py new file mode 100644 index 0000000..44c66f4 --- /dev/null +++ b/apps/site/routes.py @@ -0,0 +1,67 @@ +from collections.abc import Mapping +from typing import Any + +from flask import Blueprint, current_app, jsonify, render_template, request + +from .binding import BindingService +from .bitrix import BitrixClient + + +class SiteInputError(ValueError): + pass + + +def _field(payload: Mapping[str, Any], name: str) -> str: + """Проверка существования обязательного поля с именем name.""" + for key in (name, name.lower(), name.upper()): + value = payload.get(key) + if value is not None and str(value).strip(): + return str(value).strip() + raise SiteInputError(f"Не передано поле {name}") + + +def create_blueprint( + bindings: BindingService, + bitrix: BitrixClient, +) -> Blueprint: + blueprint = Blueprint("site", __name__) + + @blueprint.get("/") + def index(): + return jsonify( + service="bitrix-telegram-binding", + public_url=current_app.config["PUBLIC_BASE_URL"], + ) + + @blueprint.get("/health") + def health(): + database = current_app.extensions["database"] + try: + available = database.ping() + except Exception: + available = False + status = "ok" if available else "error" + return jsonify(status=status), 200 if available else 503 + + @blueprint.post("/bitrix/bind") + def bind(): + payload = request.get_json(silent=True) or request.form + refresh_token = _field(payload, "REFRESH_ID") + + # OAuth-ответ дает доверенные ID портала и пользователя. + auth = bitrix.authorize(refresh_token) + link = bindings.issue( + auth.member_id, + auth.domain, + auth.user.id, + auth.access_token, + auth.refresh_token, + auth.expires_in, + ) + return render_template( + "binding.html", + user=auth.user, + link=link, + ) + + return blueprint diff --git a/apps/site/templates/binding.html b/apps/site/templates/binding.html new file mode 100644 index 0000000..2042a60 --- /dev/null +++ b/apps/site/templates/binding.html @@ -0,0 +1,38 @@ + + + + + + Привязка Telegram + + + +

Привязка Telegram

+

{{ user.name }}, откройте бота и подтвердите привязку.

+Открыть Telegram +Ссылка одноразовая и действует до {{ link.expires_at.strftime('%H:%M + UTC') }}. + + diff --git a/compose.yaml b/compose.yaml new file mode 100644 index 0000000..5f52d72 --- /dev/null +++ b/compose.yaml @@ -0,0 +1,100 @@ +# Источник: https://jtprog.ru/posts/docker-base/ + +# Четыре моих основных сервиса: база данных, миграция, сайт и бот. Сайт и бот +# используют одну базу данных, но от имени разных пользователей. +services: + db: + image: postgres:17-alpine + restart: unless-stopped + environment: + POSTGRES_DB: ${POSTGRES_DB:-bitrix_bot} + POSTGRES_USER: ${POSTGRES_USER:-postgres} + POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?POSTGRES_PASSWORD is required} + SITE_DB_PASSWORD: ${SITE_DB_PASSWORD:?SITE_DB_PASSWORD is required} + BOT_DB_PASSWORD: ${BOT_DB_PASSWORD:?BOT_DB_PASSWORD is required} + volumes: + # Данные базы данных будут храниться в volume, чтобы при + # пересоздании контейнера данные не терялись. + - postgres_data:/var/lib/postgresql/data + # Скрипты инициализации базы данных будут монтироваться в контейнер (только для чтения), + # чтобы при пересоздании контейнера они не терялись. + - ./db/init:/docker-entrypoint-initdb.d:ro + # Проверка готовности базы данных. + # Контейнеры сайта и бота будут ждать, пока база данных не станет доступной. + healthcheck: + test: [ "CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-postgres} -d ${POSTGRES_DB:-bitrix_bot}" ] + interval: 5s + timeout: 3s + retries: 10 + + # Сервис для выполнения миграций базы данных. + migrate: + image: postgres:17-alpine + environment: + PGHOST: db + PGPORT: 5432 + PGUSER: ${POSTGRES_USER:-postgres} + PGDATABASE: ${POSTGRES_DB:-bitrix_bot} + PGPASSWORD: ${POSTGRES_PASSWORD} + command: [ "sh", "/scripts/migrate.sh" ] + volumes: + - ./db/migrations:/migrations:ro + - ./db/migrate.sh:/scripts/migrate.sh:ro + depends_on: + db: + condition: service_healthy + + site: + # Сборка отдельного образа для сайта из Dockerfile в корне проекта. + build: . + restart: unless-stopped + # Gunicorn будет запускать приложение Flask. + command: + - gunicorn + - --config + - apps/site/gunicorn.conf.py + - apps.site:create_app() + environment: + DATABASE_URL: postgresql://site_app:${SITE_DB_PASSWORD}@db:5432/${POSTGRES_DB:-bitrix_bot} + PUBLIC_BASE_URL: ${PUBLIC_BASE_URL} + BOT_USERNAME: ${BOT_USERNAME} + BITRIX_CLIENT_ID: ${BITRIX_CLIENT_ID} + BITRIX_CLIENT_SECRET: ${BITRIX_CLIENT_SECRET} + BITRIX_OAUTH_TOKEN_URL: ${BITRIX_OAUTH_TOKEN_URL:-https://oauth.bitrix.info/oauth/token/} + TOKEN_ENCRYPTION_KEY: ${TOKEN_ENCRYPTION_KEY} + BINDING_TOKEN_TTL_SECONDS: ${BINDING_TOKEN_TTL_SECONDS:-600} + TRUST_PROXY: "true" + SITE_HOST: 0.0.0.0 + SITE_PORT: ${SITE_PORT:-8000} + # Количество воркеров Gunicorn. + SITE_WORKERS: ${SITE_WORKERS:-2} + ports: + # Публикуем порт сайта на хост-машине. + - "127.0.0.1:${SITE_PUBLISHED_PORT:-8000}:${SITE_PORT:-8000}" + depends_on: + db: + condition: service_healthy + migrate: + condition: service_completed_successfully + + bot: + build: . + restart: unless-stopped + command: [ "python", "-m", "apps.bot" ] + environment: + DATABASE_URL: postgresql://bot_app:${BOT_DB_PASSWORD}@db:5432/${POSTGRES_DB:-bitrix_bot} + BOT_TOKEN: ${BOT_TOKEN} + BITRIX_CLIENT_ID: ${BITRIX_CLIENT_ID} + BITRIX_CLIENT_SECRET: ${BITRIX_CLIENT_SECRET} + BITRIX_OAUTH_TOKEN_URL: ${BITRIX_OAUTH_TOKEN_URL:-https://oauth.bitrix.info/oauth/token/} + TOKEN_ENCRYPTION_KEY: ${TOKEN_ENCRYPTION_KEY} + BITRIX_TAKE_TO_WORK_STAGE_ID: ${BITRIX_TAKE_TO_WORK_STAGE_ID:-PREPARATION} + depends_on: + db: + condition: service_healthy + migrate: + condition: service_completed_successfully + +# Именованный volume для хранения базы данных. +volumes: + postgres_data: diff --git a/db/init/01-users.sh b/db/init/01-users.sh new file mode 100644 index 0000000..43a25b0 --- /dev/null +++ b/db/init/01-users.sh @@ -0,0 +1,25 @@ +#!/bin/sh +set -eu + +: "${SITE_DB_PASSWORD:?SITE_DB_PASSWORD is required}" +: "${BOT_DB_PASSWORD:?BOT_DB_PASSWORD is required}" + +psql -v ON_ERROR_STOP=1 \ + --username "$POSTGRES_USER" \ + --dbname "$POSTGRES_DB" \ + --set=site_password="$SITE_DB_PASSWORD" \ + --set=bot_password="$BOT_DB_PASSWORD" <<'SQL' +DO $roles$ +BEGIN + IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'site_role') THEN + CREATE ROLE site_role NOLOGIN; + END IF; + IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'bot_role') THEN + CREATE ROLE bot_role NOLOGIN; + END IF; +END +$roles$; + +CREATE ROLE site_app LOGIN PASSWORD :'site_password' IN ROLE site_role; +CREATE ROLE bot_app LOGIN PASSWORD :'bot_password' IN ROLE bot_role; +SQL diff --git a/db/migrate.sh b/db/migrate.sh new file mode 100644 index 0000000..9c6b0ac --- /dev/null +++ b/db/migrate.sh @@ -0,0 +1,65 @@ +#!/bin/sh +set -eu + +PSQL="psql --no-psqlrc --set=ON_ERROR_STOP=1" + +echo "Проверка таблицы миграций" + +$PSQL <<'SQL' +CREATE TABLE IF NOT EXISTS public.schema_migrations ( + version text PRIMARY KEY, + checksum text NOT NULL, + applied_at timestamptz NOT NULL DEFAULT now() +); + +REVOKE ALL ON public.schema_migrations FROM PUBLIC; +SQL + +find /migrations \ + -maxdepth 1 \ + -type f \ + -name '[0-9][0-9][0-9]_*.sql' | +sort | +while IFS= read -r file; do + version="$(basename "$file")" + checksum="$(sha256sum "$file" | cut -d ' ' -f 1)" + + saved_checksum="$( + $PSQL \ + --tuples-only \ + --no-align \ + --set=migration_version="$version" <<'SQL' +SELECT checksum +FROM public.schema_migrations +WHERE version = :'migration_version'; +SQL + )" + + if [ -n "$saved_checksum" ]; then + if [ "$saved_checksum" != "$checksum" ]; then + echo "Ошибка: применённая миграция $version была изменена" + exit 1 + fi + + echo "Пропуск $version" + continue + fi + + echo "Применение $version" + + { + echo "BEGIN;" + cat "$file" + echo "" + echo "INSERT INTO public.schema_migrations(version, checksum)" + echo "VALUES (:'migration_version', :'migration_checksum');" + echo "COMMIT;" + } | + $PSQL \ + --set=migration_version="$version" \ + --set=migration_checksum="$checksum" + + echo "Миграция $version применена" +done + +echo "Все миграции применены" diff --git a/db/migrations/001_initial.sql b/db/migrations/001_initial.sql new file mode 100644 index 0000000..d321b58 --- /dev/null +++ b/db/migrations/001_initial.sql @@ -0,0 +1,493 @@ +-- Подключение расширения pgcrypto для генерации UUID. +CREATE EXTENSION IF NOT EXISTS pgcrypto; +CREATE SCHEMA IF NOT EXISTS binding; +CREATE SCHEMA IF NOT EXISTS oauth; + +CREATE TABLE IF NOT EXISTS binding.portals ( + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + -- member_id - уникальный идентификатор портала из Битрикса. + member_id text NOT NULL UNIQUE, + -- domain - домен портала, например: example.bitrix24.ru. + domain text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now() +); + +CREATE TABLE IF NOT EXISTS binding.tokens ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + portal_id bigint NOT NULL REFERENCES binding.portals(id) ON DELETE CASCADE, + bitrix_user_id bigint NOT NULL, + -- Хеш токена, генерируется сервером. + token_hash bytea NOT NULL UNIQUE, + expires_at timestamptz NOT NULL, + consumed_at timestamptz, + revoked_at timestamptz, + created_at timestamptz NOT NULL DEFAULT now() +); + +CREATE INDEX IF NOT EXISTS binding_tokens_owner_idx + ON binding.tokens (portal_id, bitrix_user_id, created_at DESC); + +CREATE TABLE IF NOT EXISTS binding.user_bindings ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + portal_id bigint NOT NULL REFERENCES binding.portals(id) ON DELETE CASCADE, + bitrix_user_id bigint NOT NULL, + telegram_user_id bigint NOT NULL, + telegram_chat_id bigint NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + UNIQUE (portal_id, bitrix_user_id), + UNIQUE (portal_id, telegram_user_id) +); + +CREATE TABLE IF NOT EXISTS oauth.user_credentials ( + portal_id bigint NOT NULL REFERENCES binding.portals(id) ON DELETE CASCADE, + bitrix_user_id bigint NOT NULL, + access_token bytea NOT NULL, + refresh_token bytea NOT NULL, + expires_at timestamptz NOT NULL, + -- Двойной механизм защиты от гонок данных при обновлении токенов. + version bigint NOT NULL DEFAULT 1, + refresh_locked_until timestamptz, + updated_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (portal_id, bitrix_user_id) +); + +CREATE OR REPLACE FUNCTION binding.issue_v1( + p_member_id text, + p_domain text, + p_bitrix_user_id bigint, + p_token_hash bytea, + p_token_expires_at timestamptz, + p_access_token bytea, + p_refresh_token bytea, + p_oauth_expires_at timestamptz +) +RETURNS TABLE(token_id uuid, expires_at timestamptz) +LANGUAGE plpgsql +-- Задаем SECURITY DEFINER, чтобы функция выполнялась с правами владельца +-- схемы binding. +SECURITY DEFINER +-- Устанавливаем search_path в pg_catalog, чтобы нельзя было подменить +-- функции в схеме binding или oauth. +SET search_path = pg_catalog +AS $function$ +#variable_conflict error +DECLARE + v_portal_id bigint; + v_token_id uuid; +BEGIN + -- Добавляем данные о портале, если его ещё нет, + -- или обновляем домен, если портал уже существует. + INSERT INTO binding.portals(member_id, domain) + VALUES (p_member_id, lower(p_domain)) + ON CONFLICT (member_id) DO UPDATE + SET domain = EXCLUDED.domain, + updated_at = now() + RETURNING id INTO v_portal_id; + + -- Сохраняем OAuth-данные пользователя, если их ещё нет, + -- или обновляем их, если они уже существуют. + INSERT INTO oauth.user_credentials( + portal_id, + bitrix_user_id, + access_token, + refresh_token, + expires_at + ) + VALUES ( + v_portal_id, + p_bitrix_user_id, + p_access_token, + p_refresh_token, + p_oauth_expires_at + ) + ON CONFLICT (portal_id, bitrix_user_id) DO UPDATE + SET access_token = EXCLUDED.access_token, + refresh_token = EXCLUDED.refresh_token, + expires_at = EXCLUDED.expires_at, + version = oauth.user_credentials.version + 1, + refresh_locked_until = NULL, + updated_at = now(); + + -- Отзываем все предыдущие неиспользованные токены пользователя. + UPDATE binding.tokens + SET revoked_at = now() + WHERE portal_id = v_portal_id + AND bitrix_user_id = p_bitrix_user_id + AND consumed_at IS NULL + AND revoked_at IS NULL; + + -- Выпускаем новый токен привязки. + INSERT INTO binding.tokens( + portal_id, + bitrix_user_id, + token_hash, + expires_at + ) + VALUES ( + v_portal_id, + p_bitrix_user_id, + p_token_hash, + p_token_expires_at + ) + RETURNING id INTO v_token_id; + + -- Возвращаем идентификатор токена и срок его действия. + RETURN QUERY SELECT v_token_id, p_token_expires_at; +END +$function$; + +COMMENT ON FUNCTION binding.issue_v1( + text, + text, + bigint, + bytea, + timestamptz, + bytea, + bytea, + timestamptz +) +IS $doc$ +Сохраняет OAuth-данные пользователя и выпускает токен привязки. + +Гарантии: +- предыдущие неиспользованные токены отзываются; +- при обновлении OAuth-данных увеличивается их версия. + +Возвращает: +- token_id - идентификатор токена; +- expires_at - срок действия токена. +$doc$; + +CREATE OR REPLACE FUNCTION binding.consume_v1( + p_token_hash bytea, + p_telegram_user_id bigint, + p_telegram_chat_id bigint +) +RETURNS TABLE( + member_id text, + domain text, + bitrix_user_id bigint, + telegram_user_id bigint +) +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog +AS $function$ +#variable_conflict error +DECLARE + v_portal_id bigint; + v_bitrix_user_id bigint; +BEGIN + -- UPDATE не позволит двум запросам погасить один токен. + UPDATE binding.tokens AS token + SET consumed_at = now() + WHERE token.token_hash = p_token_hash + AND token.consumed_at IS NULL + AND token.revoked_at IS NULL + AND token.expires_at > now() + RETURNING token.portal_id, token.bitrix_user_id + INTO v_portal_id, v_bitrix_user_id; + + -- Если UPDATE не вернул ни одной строки, значит токен недействителен. + IF NOT FOUND THEN + RETURN; + END IF; + + -- Удаляем все привязки к Telegram для данного портала и пользователя Битрикса, + -- кроме той, которая соответствует текущему пользователю Битрикса. + DELETE FROM binding.user_bindings AS user_binding + WHERE user_binding.portal_id = v_portal_id + AND user_binding.telegram_user_id = p_telegram_user_id + AND user_binding.bitrix_user_id <> v_bitrix_user_id; + + -- Добавляем или обновляем привязку к Telegram для данного портала + -- и пользователя Битрикса. + INSERT INTO binding.user_bindings( + portal_id, + bitrix_user_id, + telegram_user_id, + telegram_chat_id + ) + VALUES ( + v_portal_id, + v_bitrix_user_id, + p_telegram_user_id, + p_telegram_chat_id + ) + ON CONFLICT ON CONSTRAINT user_bindings_pkey DO UPDATE + SET telegram_user_id = EXCLUDED.telegram_user_id, + telegram_chat_id = EXCLUDED.telegram_chat_id, + updated_at = now(); + + RETURN QUERY + SELECT + portal.member_id, + portal.domain, + v_bitrix_user_id, + p_telegram_user_id + FROM binding.portals AS portal + WHERE portal.id = v_portal_id; +END +$function$; + +COMMENT ON FUNCTION binding.consume_v1( + bytea, + bigint, + bigint +) +IS $doc$ +Погашает токен привязки и сохраняет привязку к Telegram. + +Гарантии: +- токен погашается только один раз; +- если токен недействителен, функция возвращает пустой результат; +- если токен действителен, функция возвращает данные портала и пользователя +Битрикса, а также сохраняет привязку к Telegram; +- если пользователь Битрикса уже был привязан к другому пользователю Telegram, +старая привязка удаляется. + +Возвращает: +- member_id - идентификатор портала; +- domain - домен портала; +- bitrix_user_id - идентификатор пользователя Битрикса; +- telegram_user_id - идентификатор пользователя Telegram. +$doc$; + +CREATE OR REPLACE FUNCTION binding.find_by_telegram_v1( + p_telegram_user_id bigint, + p_member_id text DEFAULT NULL +) +RETURNS TABLE( + member_id text, + domain text, + bitrix_user_id bigint, + telegram_user_id bigint +) +LANGUAGE sql +STABLE +SECURITY DEFINER +SET search_path = pg_catalog +AS $function$ + SELECT + portal.member_id, + portal.domain, + user_binding.bitrix_user_id, + user_binding.telegram_user_id + FROM binding.user_bindings AS user_binding + JOIN binding.portals AS portal ON portal.id = user_binding.portal_id + WHERE user_binding.telegram_user_id = p_telegram_user_id + AND (p_member_id IS NULL OR portal.member_id = p_member_id) + ORDER BY user_binding.updated_at DESC + LIMIT 1 +$function$; + +COMMENT ON FUNCTION binding.find_by_telegram_v1( + bigint, + text +) +IS $doc$ +Находит привязку к Telegram по идентификатору пользователя Telegram. + +Возвращает: +- member_id - идентификатор портала; +- domain - домен портала; +- bitrix_user_id - идентификатор пользователя Битрикса; +- telegram_user_id - идентификатор пользователя Telegram. +Если p_member_id не NULL, то поиск ограничивается указанным порталом. +$doc$; + +CREATE OR REPLACE FUNCTION oauth.get_credentials_v1( + p_member_id text, + p_bitrix_user_id bigint +) +RETURNS TABLE( + member_id text, + domain text, + bitrix_user_id bigint, + access_token bytea, + refresh_token bytea, + expires_at timestamptz, + version bigint +) +LANGUAGE sql +STABLE +SECURITY DEFINER +SET search_path = pg_catalog +AS $function$ + SELECT + portal.member_id, + portal.domain, + credentials.bitrix_user_id, + credentials.access_token, + credentials.refresh_token, + credentials.expires_at, + credentials.version + FROM oauth.user_credentials AS credentials + JOIN binding.portals AS portal ON portal.id = credentials.portal_id + WHERE portal.member_id = p_member_id + AND credentials.bitrix_user_id = p_bitrix_user_id +$function$; + +COMMENT ON FUNCTION oauth.get_credentials_v1( + text, + bigint +) +IS $doc$ +Находит OAuth-данные пользователя по идентификатору портала и идентификатору +пользователя Битрикса. + +Возвращает: +- member_id - идентификатор портала; +- domain - домен портала; +- bitrix_user_id - идентификатор пользователя Битрикса; +- access_token - токен доступа; +- refresh_token - токен обновления; +- expires_at - срок действия токена доступа; +- version - версия данных. +$doc$; + +CREATE OR REPLACE FUNCTION oauth.claim_refresh_v1( + p_member_id text, + p_bitrix_user_id bigint, + p_version bigint +) +RETURNS boolean +LANGUAGE sql +VOLATILE +SECURITY DEFINER +SET search_path = pg_catalog +AS $function$ + -- Создаем временную таблицу claimed, + -- которая будет содержать результат обновления (CTE). + WITH claimed AS ( + UPDATE oauth.user_credentials AS credentials + SET refresh_locked_until = now() + interval '30 seconds' + FROM binding.portals AS portal + WHERE credentials.portal_id = portal.id + AND portal.member_id = p_member_id + AND credentials.bitrix_user_id = p_bitrix_user_id + AND credentials.version = p_version + AND ( + credentials.refresh_locked_until IS NULL + OR credentials.refresh_locked_until < now() + ) + RETURNING 1 + ) + SELECT EXISTS(SELECT 1 FROM claimed) +$function$; + +COMMENT ON FUNCTION oauth.claim_refresh_v1( + text, + bigint, + bigint +) +IS $doc$ +Пытается захватить токен обновления для пользователя. + +Возвращает: +- true, если захват успешен; +- false, если захват не удался. +$doc$; + +CREATE OR REPLACE FUNCTION oauth.finish_refresh_v1( + p_member_id text, + p_bitrix_user_id bigint, + p_version bigint, + p_access_token bytea, + p_refresh_token bytea, + p_expires_at timestamptz +) +RETURNS boolean +LANGUAGE sql +VOLATILE +SECURITY DEFINER +SET search_path = pg_catalog +AS $function$ + WITH updated AS ( + UPDATE oauth.user_credentials AS credentials + SET access_token = p_access_token, + refresh_token = p_refresh_token, + expires_at = p_expires_at, + version = credentials.version + 1, + refresh_locked_until = NULL, + updated_at = now() + FROM binding.portals AS portal + WHERE credentials.portal_id = portal.id + AND portal.member_id = p_member_id + AND credentials.bitrix_user_id = p_bitrix_user_id + AND credentials.version = p_version + RETURNING 1 + ) + SELECT EXISTS(SELECT 1 FROM updated) +$function$; + +COMMENT ON FUNCTION oauth.finish_refresh_v1( + text, + bigint, + bigint, + bytea, + bytea, + timestamptz +) IS $doc$ +Завершает процесс обновления токена для пользователя. + +Возвращает: +- true, если обновление успешно завершено; +- false, если обновление не удалось. +$doc$; + +CREATE OR REPLACE FUNCTION oauth.release_refresh_v1( + p_member_id text, + p_bitrix_user_id bigint, + p_version bigint +) +RETURNS void +LANGUAGE sql +VOLATILE +SECURITY DEFINER +SET search_path = pg_catalog +AS $function$ + UPDATE oauth.user_credentials AS credentials + SET refresh_locked_until = NULL + FROM binding.portals AS portal + WHERE credentials.portal_id = portal.id + AND portal.member_id = p_member_id + AND credentials.bitrix_user_id = p_bitrix_user_id + AND credentials.version = p_version +$function$; + +COMMENT ON FUNCTION oauth.release_refresh_v1( + text, + bigint, + bigint +) IS $doc$ +Освобождает токен обновления для пользователя. + +Возвращает: +- void. +$doc$; + +-- Отзываем все права у PUBLIC. +REVOKE ALL ON ALL TABLES IN SCHEMA binding, oauth FROM PUBLIC; +REVOKE EXECUTE ON ALL FUNCTIONS IN SCHEMA binding, oauth FROM PUBLIC; + +-- Даем права на использование схемы и выполнение функций ролям site_role +-- и bot_role. +GRANT USAGE ON SCHEMA binding TO site_role, bot_role; +GRANT USAGE ON SCHEMA oauth TO bot_role; + +-- Даем права на выполнение функций ролям site_role и bot_role. +GRANT EXECUTE ON FUNCTION binding.issue_v1( + text, text, bigint, bytea, timestamptz, bytea, bytea, timestamptz +) TO site_role; +GRANT EXECUTE ON FUNCTION binding.consume_v1(bytea, bigint, bigint) TO bot_role; +GRANT EXECUTE ON FUNCTION binding.find_by_telegram_v1(bigint, text) TO bot_role; +GRANT EXECUTE ON FUNCTION oauth.get_credentials_v1(text, bigint) TO bot_role; +GRANT EXECUTE ON FUNCTION oauth.claim_refresh_v1(text, bigint, bigint) TO bot_role; +GRANT EXECUTE ON FUNCTION oauth.finish_refresh_v1( + text, bigint, bigint, bytea, bytea, timestamptz +) TO bot_role; +GRANT EXECUTE ON FUNCTION oauth.release_refresh_v1(text, bigint, bigint) + TO bot_role; diff --git a/docs/adr/001-apps-and-database.md b/docs/adr/001-apps-and-database.md new file mode 100644 index 0000000..bf496ef --- /dev/null +++ b/docs/adr/001-apps-and-database.md @@ -0,0 +1,64 @@ +# ADR-001: разделение приложения на сайт, Telegram-бот и базу данных + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Приложение включает страницу привязки пользователя Битрикс24, Telegram-интерфейс +менеджера и хранилище интеграционных данных. HTTP-сайт обрабатывает короткие +входящие запросы, тогда как Telegram-бот выполняет длительный polling и +параллельные REST-операции. + +## Решение + +Архитектура разделена на три основные структурные единицы: сайт привязки, +Telegram-бот и база данных, которая обслуживает остальные компоненты. Зона +ответственности каждого модуля определена отдельно. Сайт и бот не делят общий +код, а их права на уровне базы данных ограничены. + +HTTP-приложение `site` обслуживает только инициацию привязки и взаимодействует с +REST API и OAuth Битрикс24. Telegram-приложение `bot` обрабатывает команды +менеджера и взаимодействует как с REST API и OAuth Битрикс24, так и с Telegram +Bot API. PostgreSQL предоставляет обоим процессам устойчивый версионированный +контракт в виде `SECURITY DEFINER`-функций и поддерживает применение миграций. + +Программное решение регистрируется администратором портала как локальное +приложение с указанием ссылки на страницу привязки. Локальное приложение +отправляет на HTTPS-адрес `/bitrix/bind` идентификационные данные пользователя и +refresh-токен. Сайт обменивает его на новую пару OAuth-токенов, извлекает +доверенные `member_id`, `user_id` и `client_endpoint` из ответа Битрикс24, +проверяет пользователя методом `user.current` и только после этого выпускает +одноразовую ссылку Telegram. + +![UML-диаграмма компонентов программного решения](assets/report/architecture-components.png) + +*Рисунок ADR-001/1. UML-диаграмма компонентов программного решения* + +На схеме также обозначен сервис `migrate`. Он не является постоянно запущенным +модулем, но отвечает за миграции схемы БД. При перезапуске Docker Compose сервис +последовательно выполняет необходимые SQL-скрипты. Сайт имеет право выполнять +только функцию выпуска ссылки `binding.issue_v1`. Бот погашает ссылку, получает +привязку и обращается к OAuth-функциям. Прямые операции `SELECT`, `INSERT` и +`UPDATE` над таблицами для ролей приложений запрещены, поэтому граница базы +данных одновременно является границей доступа. + +| Компонент | Ответственность | Внешний интерфейс | +|------------|----------------------------------------------------------|------------------------------------------| +| nginx | Завершение TLS и проксирование только к сайту | HTTPS → 127.0.0.1:8000 | +| apps.site | OAuth-проверка пользователя и выпуск одноразовой ссылки | POST /bitrix/bind, GET /health | +| apps.bot | Команды Telegram, карточки и операции со сделками | Telegram Bot API, Битрикс REST | +| PostgreSQL | Привязки, токены, транзакционная синхронизация | binding.\* и oauth.\* | +| migrate | Однократное применение SQL-миграций до старта приложений | db/migrations/\*.sql | +| Битрикс24 | Источник CRM-данных и OAuth-контекста | oauth/token, user.current, crm.\* | +| Telegram | Пользовательский канал и доставка callback-событий | getUpdates, sendMessage, editMessageText | + +*Таблица ADR-001/1. Ответственность компонентов архитектуры* + +## Последствия + +Разделение уменьшает связанность и позволяет перезапускать или масштабировать +процессы независимо. Для коротких входящих запросов используется синхронный +Flask/Gunicorn, а для длительного polling и параллельных REST-операций — +asyncio/aiogram. Связь приложений формализована версионированными функциями +PostgreSQL вместо общего программного модуля. diff --git a/docs/adr/002-oauth-telegram-binding.md b/docs/adr/002-oauth-telegram-binding.md new file mode 100644 index 0000000..64a0a4b --- /dev/null +++ b/docs/adr/002-oauth-telegram-binding.md @@ -0,0 +1,46 @@ +# ADR-002: привязка пользователей Битрикс24 и Telegram + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Локальное приложение CRM отправляет на HTTPS-адрес `/bitrix/bind` +идентификационные данные пользователя и refresh-токен. Идентификаторам портала и +пользователя из входной формы доверять нельзя: контекст должен быть получен от +OAuth-сервера Битрикс24 и подтверждён методом `user.current`. + +## Решение + +Сайт использует refresh-токен для получения новой OAuth-пары, доверенных +`member_id`, `user_id` и `client_endpoint`. После этого `BitrixClient` сверяет +`user_id` с результатом `user.current`. + +Процесс привязки учётных записей представлен на диаграмме последовательности. + +![Диаграмма последовательности привязки Битрикс24 к Telegram](assets/report/binding-sequence.png) + +*Рисунок ADR-002/1. Диаграмма последовательности привязки Битрикс24 к Telegram* + +После проверки пользователя функцией `secrets.token_urlsafe(32)` формируется +одноразовый токен привязки. В БД записывается только SHA-256-хеш, поэтому +компрометация базы не позволяет восстановить действующую ссылку. Срок жизни +задаётся переменной окружения `BINDING_TOKEN_TTL_SECONDS`, ограничен диапазоном +от 60 до 3600 секунд и по умолчанию равен 600 секундам. При повторном выпуске +прежние непогашенные токены того же пользователя отзываются. + +Пользователь переходит по одноразовой ссылке в чат с Telegram-ботом. Бот +повторно вычисляет SHA-256-хеш и сверяет его с активными токенами. Если токен +существует, не истёк, не отозван и ещё не погашен, он помечается использованным, +а в таблице привязок создаётся или обновляется связь пользователя Битрикс24 с +аккаунтом Telegram. + +Погашение выполняется только в личном чате. Проверка токена и изменение привязки +выполняются функцией `binding.consume_v1` в одной транзакции. + +## Последствия + +Привязка не использует идентификаторы пользователя из недоверенной входной +формы. В базе хранится только хеш одноразового токена, а повторный выпуск ссылки +отзывает предыдущие непогашенные токены. Атомарное погашение не позволяет двум +запросам одновременно использовать одну ссылку. diff --git a/docs/adr/003-oauth-credential-lifecycle.md b/docs/adr/003-oauth-credential-lifecycle.md new file mode 100644 index 0000000..8b783cf --- /dev/null +++ b/docs/adr/003-oauth-credential-lifecycle.md @@ -0,0 +1,53 @@ +# ADR-003: защита и обновление OAuth-токенов + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Для выполнения REST-запросов приложение хранит `access_token` и `refresh_token`. +Битрикс24 возвращает новую пару токенов при каждом обновлении, поэтому +одновременное использование одного refresh-токена несколькими воркерами может +привести к потере актуальной пары. + +## Решение + +До передачи в PostgreSQL `access_token` и `refresh_token` шифруются алгоритмом +Fernet. Общий `TOKEN_ENCRYPTION_KEY` передаётся контейнерам `site` и `bot` через +переменные окружения, но не записывается в базу. Бот расшифровывает access-токен +непосредственно перед REST-запросом и не включает OAuth-параметры в тексты +ошибок. + +Принятые меры защиты сведены в таблицу. + +| Риск | Реализованная мера | Остаточный контроль | +|---------------------------------|--------------------------------------------|--------------------------------------| +| Утечка одноразовой ссылки из БД | Хранение SHA-256-хеша | Короткий TTL и однократное погашение | +| Чтение OAuth-токенов из БД | Fernet-шифрование до INSERT/UPDATE | Секретный ключ вне БД | +| Подмена пользователя | member_id/user_id из OAuth + user.current | Проверка HTTPS endpoint | +| Гонка refresh token | Версия и аренда refresh_locked_until | Повторное чтение до версии N+1 | +| Избыточные права приложений | Разные роли и EXECUTE только на функции | REVOKE для PUBLIC | +| Долгая транзакция | Сетевые запросы выполняются вне транзакции | Короткие контексты Psycopg | + +*Таблица ADR-003/1. Риски и меры защиты от них* + +При получении `expired_token`, `invalid_token` или `no_auth_found` клиент +пытается обновить пару токенов. Поле `version` реализует оптимистическую +проверку, а `refresh_locked_until` — короткую аренду продолжительностью 30 +секунд. Это предотвращает одновременное использование одного refresh-токена +несколькими воркерами. + +После успешного обновления токенов воркер освобождает аренду. Если она занята, +другой воркер ожидает обновления, после чего повторяет обращение к API с новой +парой токенов. + +![Диаграмма последовательности обновления OAuth-токена](assets/report/oauth-refresh-sequence.png) + +*Рисунок ADR-003/1. Диаграмма последовательности обновления OAuth-токена* + +## Последствия + +Сетевой запрос к OAuth выполняется вне транзакции PostgreSQL. Версия и аренда +координируют обновление между воркерами, а повторное чтение позволяет продолжить +работу с версией `N+1`. Дальнейшее развитие механизма защиты предусматривает +ротацию ключей шифрования. diff --git a/docs/adr/004-bot-layers.md b/docs/adr/004-bot-layers.md new file mode 100644 index 0000000..96d3728 --- /dev/null +++ b/docs/adr/004-bot-layers.md @@ -0,0 +1,61 @@ +# ADR-004: слоистая организация Telegram-бота + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Telegram-бот принимает команды и callback-запросы, проверяет привязку +пользователя, обращается к PostgreSQL и REST API Битрикс24, а затем формирует +HTML-сообщения и inline-клавиатуры. + +## Решение + +Бот организован по слоям: обработчики принимают события Telegram, middleware +добавляет проверенную привязку в контекст, сервисы реализуют прикладные +сценарии, `BitrixClient` отвечает за OAuth и HTTP, а классы представления +формируют HTML-тексты и inline-клавиатуры. + +![UML-диаграмма основных классов решения](assets/report/bot-class-diagram.png) + +*Рисунок ADR-004/1. UML-диаграмма основных классов решения* + +Классы `Binding`, `OAuthCredentials`, `ClientInfo`, `DealStageFilter` и +`DealPage` являются dataclass-моделями передачи данных. Они отделяют словари +REST-ответов и строки БД от интерфейсов сервисов. `DealPage` дополнительно +вычисляет признак `has_next`, который используется при построении кнопок +пагинации. + +Основные команды и callback-действия Telegram-бота приведены в таблице. + +| Ввод | Обработчик | Результат | +|-----------------------------------------|---------------------------------|--------------------------------------------| +| /start, /help | StartBotHandlers.start | Справка или состояние привязки | +| /start bind_<token> | StartBotHandlers.bind | Погашение одноразовой ссылки в личном чате | +| /deals, /leads | DealBotHandlers.deals | Первая страница сделок начальной стадии | +| /deal ID, /lead ID | DealBotHandlers.deal_by_command | Карточка сделки по идентификатору | +| deals:page:<stage>:<page> | deals_page | Фильтрация по стадии и пагинация | +| deal:view:<id> | deal_by_button | Карточка выбранной сделки | +| deal:assign:<id>:<expected> | assign_responsible | Назначение текущего Битрикс-пользователя | +| deal:remind:<id> | remind_to_call | Создание дела на звонок через час | +| deal:history:<id> | show_history | Пять последних переходов по стадиям | + +*Таблица ADR-004/1. Пользовательские команды и callback-действия* + +Навигация по страницам списка, переход к карточке и возврат выполняются кнопками +клавиатуры с редактированием исходного сообщения бота. Список сделок также +поддерживает фильтрацию, которая задаётся дополнительными кнопками на странице +просмотра списка. + +Перед выполнением CRM-команд `BindingRequiredMiddleware` ищет привязку по +Telegram user id. При успешной проверке объект `Binding` помещается в словарь +`data` и передаётся именованным параметром обработчика. Если привязки нет, +цепочка прерывается до REST-запроса, а пользователь получает инструкцию открыть +приложение в Битрикс24. + +## Последствия + +Конструкторы принимают зависимости явно, поэтому сервисы можно тестировать с +имитационными репозиториями и REST-клиентами. Проверка middleware ограждает +пользователя от ошибочного поведения и гарантирует наличие привязки перед +обращением к CRM. diff --git a/docs/adr/005-bitrix-deal-read-model.md b/docs/adr/005-bitrix-deal-read-model.md new file mode 100644 index 0000000..5c7c844 --- /dev/null +++ b/docs/adr/005-bitrix-deal-read-model.md @@ -0,0 +1,83 @@ +# ADR-005: получение списка и карточки сделки из Битрикс24 + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Приложение не копирует CRM-данные в локальную базу. Сведения о сделках, +контактах, компаниях, стадиях и истории запрашиваются через REST API +непосредственно в момент действия пользователя. PostgreSQL хранит только данные, +необходимые для идентификации пользователя и выполнения авторизованных запросов. + +Реализация работает с сущностью сделки и методами `crm.deal.*`. Команды `/leads` +и `/lead` используются как пользовательские псевдонимы `/deals` и `/deal`. + +## Решение + +Для получения и изменения данных применяются следующие методы REST API +Битрикс24. + +| Метод | Назначение | Ключевые параметры | +|-----------------------|-------------------------------|----------------------------------------| +| crm.deal.list | Список и пагинация | filter, select, order, start | +| crm.deal.get | Карточка и контроль состояния | id | +| crm.deal.update | Ответственный и стадия | id, fields, REGISTER_HISTORY_EVENT | +| crm.status.list | Стадии воронки и источники | ENTITY_ID, STATUS_ID | +| crm.contact.get | ФИО и телефон контакта | id | +| crm.company.get | Название и телефон компании | id | +| crm.stagehistory.list | История переходов | entityTypeId=2, OWNER_ID | +| crm.activity.todo.add | Отложенный звонок | ownerTypeId=2, deadline, responsibleId | + +*Таблица ADR-005/1. Используемые методы REST API Битрикс24* + +### Формирование списка + +Названия стадий не зашиты в интерфейсе. Метод `crm.status.list` получает +актуальную конфигурацию воронки, после чего первая стадия трактуется как +псевдофильтр `new`. Карта стадий кэшируется в памяти на 300 секунд отдельно для +портала, пользователя и категории. Дополнительно добавляется фильтр «Все», не +передающий `STAGE_ID` в Битрикс24. + +Размер страницы Telegram равен пяти сделкам, тогда как Битрикс24 может +возвращать другое количество элементов за запрос. `DealService` собирает +REST-страницы по полю `next` до тех пор, пока не сможет выделить диапазон +`[page * limit; page * limit + limit)`. Значение `total` используется для +расчёта общего числа страниц. Список сделок запрашивается методом +`crm.deal.list` с параметрами фильтрации. + +![Список сделок в интерфейсе Telegram](assets/report/deal-list-ui.png) + +*Рисунок ADR-005/1. Список сделок в интерфейсе Telegram* + +### Формирование карточки + +Получение карточки сделки продолжает сценарий работы со списком. + +![Диаграмма последовательности просмотра списка и карточки сделки](assets/report/deal-list-card-sequence.png) + +*Рисунок ADR-005/2. Диаграмма последовательности просмотра списка и карточки +сделки* + +Карточка загружается методом `crm.deal.get`, затем обогащается данными связанных +сущностей. Для контакта составляется ФИО и выбирается первый телефон; при +отсутствии телефона контакта проверяется компания. Идентификаторы источника и +стадии преобразуются в человекочитаемые названия. В итоговое сообщение +включаются сумма, валюта, ответственный, дата создания и комментарий. + +Все динамические строки перед включением в HTML-ответ Telegram проходят +`html.escape`. Длина карточки ограничена 3900 символами, что оставляет запас до +ограничения Telegram и предотвращает ошибку отправки из-за длинного комментария. +Кнопка назначения отображается только для новой сделки, если текущий +пользователь ещё не является ответственным. + +![Карточка сделки в интерфейсе Telegram](assets/report/deal-card-ui.png) + +*Рисунок ADR-005/3. Карточка сделки в интерфейсе Telegram* + +## Последствия + +Битрикс24 остаётся источником актуальных CRM-данных, а локальная база не требует +синхронизации сделок и связанных сущностей. В качестве дальнейшего развития +предусмотрен переход с устаревающих методов `crm.deal.*` на универсальные методы +`crm.item.*`. diff --git a/docs/adr/006-guarded-deal-mutations.md b/docs/adr/006-guarded-deal-mutations.md new file mode 100644 index 0000000..f9b59e7 --- /dev/null +++ b/docs/adr/006-guarded-deal-mutations.md @@ -0,0 +1,60 @@ +# ADR-006: изменение сделки с проверкой актуального состояния + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Состояние сделки может измениться в Битрикс24 после формирования карточки в +Telegram, но до нажатия callback-кнопки. Используемый метод `crm.deal.update` не +предоставляет условный `UPDATE`, поэтому перед изменением требуется проверить, +что показанное пользователю состояние остаётся актуальным. + +## Решение + +### Назначение ответственного и изменение стадии + +Callback-кнопка назначения содержит не только id сделки, но и `ASSIGNED_BY_ID`, +который был показан пользователю: +`deal:assign::`. Перед изменением сервис +повторно загружает сделку и сравнивает фактического ответственного с ожидаемым. +Если карточка устарела, REST-обновление не выполняется. + +Внутри одного процесса операции по паре `(member_id, deal_id)` последовательно +выполняются под `asyncio.Lock`. После `crm.deal.update` сервис повторно читает +сделку и убеждается, что ответственным стал `bitrix_user_id` привязанного +пользователя. + +![Диаграмма последовательности взятия сделки в работу](assets/report/deal-assignment-sequence.png) + +*Рисунок ADR-006/1. Диаграмма последовательности взятия сделки в работу* + +При обновлении одновременно передаются `ASSIGNED_BY_ID` связанного пользователя, +рабочая `STAGE_ID` и параметр `REGISTER_HISTORY_EVENT=Y`. После REST-запроса +выполняется контрольное чтение сделки. + +### Планирование звонка и просмотр истории + +Действие «Позвонить позже» создаёт в Битрикс24 дело типа `todo` с крайним сроком +через один час. Владельцем является сделка (`ownerTypeId=2`), а ответственным — +связанный пользователь Битрикс24. Массив `pingOffsets=[0]` включает напоминание +в момент наступления срока. + +История загружается методом `crm.stagehistory.list` с фильтром `OWNER_ID` и +сортировкой по убыванию идентификатора. В интерфейс выводятся первые пять +событий. Для каждого события идентификатор стадии преобразуется в название с +учётом `CATEGORY_ID`, после чего `DealFormatter` формирует защищённый HTML-текст +и клавиатуру возврата к карточке или списку. + +![Диаграмма последовательности планирования звонка и просмотра истории](assets/report/reminder-history-sequence.png) + +*Рисунок ADR-006/2. Диаграмма последовательности планирования звонка и просмотра +истории* + +## Последствия + +Локальный `asyncio.Lock` защищает только один процесс `bot`. При горизонтальном +масштабировании на несколько контейнеров потребуется распределённая блокировка +либо серверная условная операция. Повторная проверка REST-результата сохраняет +защиту от внешних изменений, но не делает два удалённых вызова одной транзакцией +Битрикс24. diff --git a/docs/adr/007-container-deployment.md b/docs/adr/007-container-deployment.md new file mode 100644 index 0000000..50bdd77 --- /dev/null +++ b/docs/adr/007-container-deployment.md @@ -0,0 +1,46 @@ +# ADR-007: развёртывание приложения с помощью Docker Compose + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Программное решение состоит из PostgreSQL, сервиса миграций, Telegram-бота и +Flask-сайта. База данных должна быть готова до запуска приложений, а +SQL-миграции должны выполняться последовательно с сохранением истории +применения. + +## Решение + +Для развёртывания используется Docker Compose из четырёх основных сервисов: +`db`, `site`, `bot` и `migrate`. Сервисы базы данных и миграций используют +готовые образы на основе Alpine Linux. Модули приложения собираются с помощью +Dockerfile на базе среды выполнения Python 3.13 и включают необходимые +библиотеки. + +Первым запускается контейнер `db` с PostgreSQL, который также импортирует +первичные настройки ролей. Контейнер PostgreSQL не публикует порт на хост и +остаётся доступным только внутри сети Compose. + +Затем запускается контейнер `migrate`. Он последовательно применяет SQL-скрипты +из `db/migrations` и сохраняет историю их применения в таблице +`public.schema_migrations`. Для каждой миграции хранится контрольная сумма. Если +уже применённый файл был изменён, выполнение завершается с ошибкой. + +После успешного завершения миграций запускаются Telegram-бот и Flask-сайт. +Условиями их запуска являются нормальное состояние контейнера базы данных и +успешное завершение контейнера `migrate`. + +Сайт доступен только через loopback-адрес `127.0.0.1`. Внешний nginx принимает +HTTPS-трафик и передаёт заголовки `X-Forwarded-*`. Бот не имеет входящего порта +и получает обновления методом long polling. + +Контейнеры приложения выполняются от системного пользователя `runtime`, а не от +`root`. + +## Последствия + +Миграции выполняются до запуска прикладных процессов. PostgreSQL не публикуется +наружу, сайт доступен извне только через HTTPS-прокси, а Telegram-бот не требует +входящего сетевого порта. Изменение уже применённой миграции обнаруживается по +несовпадению контрольной суммы. diff --git a/docs/adr/008-data-storage-and-access.md b/docs/adr/008-data-storage-and-access.md new file mode 100644 index 0000000..8e43971 --- /dev/null +++ b/docs/adr/008-data-storage-and-access.md @@ -0,0 +1,70 @@ +# ADR-008: хранение интеграционных данных и разграничение доступа + +**Статус:** Принято +**Дата:** 2026-07-23 + +## Контекст + +Основное назначение базы данных — хранение авторизационных данных, токенов +привязки, OAuth-токенов и связей между пользователями Битрикс24 и Telegram. +Сделки, контакты, компании, стадии и история остаются в Битрикс24 и в локальной +базе не дублируются. + +## Решение + +Модель хранения нормализована вокруг портала Битрикс24. + +![Схема базы данных](assets/report/database-schema.png) + +*Рисунок ADR-008/1. Схема базы данных* + +| Таблица | Ключевые данные | Назначение и ограничения | +|------------------------|-------------------------------------------------|-------------------------------------------------------| +| binding.portals | member_id, domain | Справочник порталов; member_id уникален | +| binding.tokens | token_hash, expires_at, consumed_at, revoked_at | Одноразовые ссылки; токен хранится только как хеш | +| binding.user_bindings | portal_id, bitrix_user_id, telegram_user_id | Однозначная привязка пользователей в пределах портала | +| oauth.user_credentials | access_token, refresh_token, version, lock | Зашифрованные OAuth-данные и координация обновления | + +*Таблица ADR-008/1. Назначение таблиц базы данных* + +Поле `member_id` является устойчивым внешним идентификатором портала, а числовой +`bitrix_user_id` имеет смысл только вместе с `portal_id`. Связи с `portals` +используют `ON DELETE CASCADE`: удаление портала автоматически удаляет его +ссылки, привязки и OAuth-данные. Составной первичный ключ +`oauth.user_credentials(portal_id, bitrix_user_id)` исключает две конкурирующие +записи учётных данных одного пользователя. Поля `consumed_at` и `revoked_at` +разделяют два независимых основания недействительности одноразовой ссылки. + +Для работы с данными в PostgreSQL созданы две роли: `site_role` и `bot_role`. +Они имеют разные права доступа к хранимым функциям и не могут редактировать +таблицы напрямую. + +| Функция | Вызывающая роль | Назначение | +|-----------------------------|-----------------|-------------------------------------------------------------------------------------| +| binding.issue_v1 | site_role | Сохранить OAuth-данные, отозвать старые ссылки и выпустить новую в одной транзакции | +| binding.consume_v1 | bot_role | Однократно погасить ссылку и создать привязку | +| binding.find_by_telegram_v1 | bot_role | Получить актуальную привязку Telegram | +| oauth.get_credentials_v1 | bot_role | Получить учётные данные связанного пользователя | +| oauth.claim_refresh_v1 | bot_role | Получить короткую аренду на обновление токена | +| oauth.finish_refresh_v1 | bot_role | Атомарно записать новую пару при совпадении версии | +| oauth.release_refresh_v1 | bot_role | Освободить аренду после ошибки | + +*Таблица ADR-008/2. Контракт хранимых функций PostgreSQL* + +Сайт использует роль `site_role` и может вызывать только функцию привязки. +Telegram-бот использует `bot_role` и может вызывать функции погашения токенов, +получения привязки и работы с OAuth. Роль `site_role` не имеет доступа к схеме +`oauth` и не может читать сохранённые токены; роль `bot_role` не может выпускать +новые ссылки от имени сайта. + +Все функции объявлены `SECURITY DEFINER` и фиксируют `search_path` в +`pg_catalog`, что уменьшает риск подмены объектов. Права `PUBLIC` на таблицы и +функции отозваны. + +## Последствия + +Конкурирующие вызовы `binding.consume_v1` для одного хеша не смогут одновременно +пройти условие `consumed_at IS NULL`: `UPDATE` блокирует строку, а после +завершения первой транзакции второй вызов видит уже установленное время +погашения. Проверка и изменение не разделены между приложением и базой, поэтому +отсутствует окно гонки между `SELECT` и `UPDATE`. diff --git a/docs/adr/assets/report/architecture-components.png b/docs/adr/assets/report/architecture-components.png new file mode 100644 index 0000000..42a4177 Binary files /dev/null and b/docs/adr/assets/report/architecture-components.png differ diff --git a/docs/adr/assets/report/binding-sequence.png b/docs/adr/assets/report/binding-sequence.png new file mode 100644 index 0000000..06601c3 Binary files /dev/null and b/docs/adr/assets/report/binding-sequence.png differ diff --git a/docs/adr/assets/report/bot-class-diagram.png b/docs/adr/assets/report/bot-class-diagram.png new file mode 100644 index 0000000..84e8533 Binary files /dev/null and b/docs/adr/assets/report/bot-class-diagram.png differ diff --git a/docs/adr/assets/report/database-schema.png b/docs/adr/assets/report/database-schema.png new file mode 100644 index 0000000..7a3062e Binary files /dev/null and b/docs/adr/assets/report/database-schema.png differ diff --git a/docs/adr/assets/report/deal-assignment-sequence.png b/docs/adr/assets/report/deal-assignment-sequence.png new file mode 100644 index 0000000..3296112 Binary files /dev/null and b/docs/adr/assets/report/deal-assignment-sequence.png differ diff --git a/docs/adr/assets/report/deal-card-ui.png b/docs/adr/assets/report/deal-card-ui.png new file mode 100644 index 0000000..0d8eb3d Binary files /dev/null and b/docs/adr/assets/report/deal-card-ui.png differ diff --git a/docs/adr/assets/report/deal-list-card-sequence.png b/docs/adr/assets/report/deal-list-card-sequence.png new file mode 100644 index 0000000..13d1dd9 Binary files /dev/null and b/docs/adr/assets/report/deal-list-card-sequence.png differ diff --git a/docs/adr/assets/report/deal-list-ui.png b/docs/adr/assets/report/deal-list-ui.png new file mode 100644 index 0000000..8dca83b Binary files /dev/null and b/docs/adr/assets/report/deal-list-ui.png differ diff --git a/docs/adr/assets/report/oauth-refresh-sequence.png b/docs/adr/assets/report/oauth-refresh-sequence.png new file mode 100644 index 0000000..ea4bf49 Binary files /dev/null and b/docs/adr/assets/report/oauth-refresh-sequence.png differ diff --git a/docs/adr/assets/report/reminder-history-sequence.png b/docs/adr/assets/report/reminder-history-sequence.png new file mode 100644 index 0000000..e4973c1 Binary files /dev/null and b/docs/adr/assets/report/reminder-history-sequence.png differ diff --git a/docs/adr/main.md b/docs/adr/main.md index 15cc34c..6f1031e 100644 --- a/docs/adr/main.md +++ b/docs/adr/main.md @@ -1,13 +1,43 @@ -# Архитектурные решения (ADR) +# Архитектурные решения BitrixDealsBot -## Введение +В ходе производственной практики разработано серверное приложение +BitrixDealsBot, которое предоставляет менеджеру по продажам интерфейс Telegram +для работы со сделками CRM Битрикс24. Разработанное решение не копирует +CRM-данные в локальную базу: все сведения о сделках, контактах, компаниях, +стадиях и истории запрашиваются через REST API непосредственно в момент действия +пользователя. PostgreSQL хранит только данные, необходимые для идентификации +пользователя Битрикс24 и его привязки к пользователю Telegram-бота, включая +OAuth-токены. -Данный раздел содержит архитектурные решения (ADR) для проекта. -Архитектурные решения описывают ключевые решения, принятые в процессе разработки -системы, включая выбор технологий, подходов и структурных решений. +В индивидуальном задании используется термин «лид», однако реализация работает с +сущностью сделки и методами `crm.deal.*`. Команды `/leads` и `/lead` сохранены +как пользовательские псевдонимы `/deals` и `/deal`, поэтому интерфейс остаётся +совместимым с формулировкой задания, а в документации используется технически +точное понятие «сделка». -## Оглавление +## Технологический стек -**В данном разделе представлены следующие архитектурные решения:** +| Уровень | Технология | Назначение | +|---------------|---------------------------|--------------------------------------------------------------------------| +| Язык | Python 3.13 | Серверная логика сайта и Telegram-бота | +| Telegram | aiogram 3 | Асинхронная маршрутизация команд, callback-запросов и опроса сервера | +| HTTP-сервер | Flask 3 + Gunicorn | Страница привязки пользователя Битрикс24 к конкретному аккаунту Telegram | +| Запросы к API | httpx | Синхронные и асинхронные запросы к OAuth и REST API Битрикс24 | +| Хранилище | PostgreSQL 17 + Psycopg 3 | Транзакции, хранимые функции и пулы соединений | +| Защита | Fernet + SHA-256 | Шифрование OAuth-токенов и хеширование одноразовых ссылок | +| Развёртывание | Docker Compose + nginx | Изоляция процессов, миграции, HTTPS и обратное проксирование | -_(В процессе разработки будут добавляться новые решения)_ \ No newline at end of file +*Таблица 1. Технологический стек решения* + +## Состав группы ADR + +| ADR | Архитектурное решение | Статус | +|----------------------------------------------|-----------------------------------------------------------|---------| +| [ADR-001](001-apps-and-database.md) | Разделение приложения на сайт, Telegram-бот и базу данных | Принято | +| [ADR-002](002-oauth-telegram-binding.md) | Привязка пользователей Битрикс24 и Telegram | Принято | +| [ADR-003](003-oauth-credential-lifecycle.md) | Защита и обновление OAuth-токенов | Принято | +| [ADR-004](004-bot-layers.md) | Слоистая организация Telegram-бота | Принято | +| [ADR-005](005-bitrix-deal-read-model.md) | Получение списка и карточки сделки из Битрикс24 | Принято | +| [ADR-006](006-guarded-deal-mutations.md) | Изменение сделки с проверкой актуального состояния | Принято | +| [ADR-007](007-container-deployment.md) | Развёртывание приложения с помощью Docker Compose | Принято | +| [ADR-008](008-data-storage-and-access.md) | Хранение интеграционных данных и разграничение доступа | Принято | diff --git a/requirements.txt b/requirements.txt index 71badf2..b965b7e 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,3 +1,8 @@ httpx~=0.28.1 aiogram~=3.29.1 -python-dotenv~=1.2.2 \ No newline at end of file +python-dotenv~=1.2.2 +Flask>=3.1,<4 +gunicorn>=23,<24 +psycopg[binary,pool]>=3.2,<4 +cryptography>=44,<48 +werkzeug>=3.1,<4 \ No newline at end of file