Добавление всех наработок за период практики.

This commit is contained in:
SkyForces
2026-07-24 00:46:19 +03:00
parent 24943b8a73
commit ae4141ba28
56 changed files with 3987 additions and 136 deletions
View File
+3
View File
@@ -0,0 +1,3 @@
from .main import main
__all__ = ["main"]
+3
View File
@@ -0,0 +1,3 @@
from .main import main
main()
+88
View File
@@ -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)
+161
View File
@@ -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()
+42
View File
@@ -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
)
)
+20
View File
@@ -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
+38
View File
@@ -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
+641
View File
@@ -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
+73
View File
@@ -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
+531
View File
@@ -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 сделки: <code>/deal 123</code>",
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"<b>{html.escape(advance.target_stage_title)}</b>"
f"{final_note}.\n"
"Введите комментарий одним сообщением. "
"Чтобы продолжить без комментария, отправьте "
"<code>-</code>. Для отмены — <code>/cancel</code>."
),
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(
"Отправьте комментарий текстом или прочерк "
"<code>-</code>, чтобы продолжить без него.",
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")
+57
View File
@@ -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()
+34
View File
@@ -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
+84
View File
@@ -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"])
)
+272
View File
@@ -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"<b>Сделки: {title}</b>\n"
f"Всего сделок: <b>{total_deals}</b>\n"
f"Страница <b>{page + 1}</b> из <b>{total_pages}</b>"
)
@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"<b>Сделка #{deal_id}</b>",
f"<b>{title}</b>",
"",
f"Клиент: <code>{client}</code>",
*([f"Компания: <code>{company}</code>"] if company else []),
f"Телефон клиента: <code>{phone}</code>",
f"Стадия: <code>{stage}</code>",
f"Источник сделки: <code>{source}</code>",
f"Сумма: {cls.money(deal.get('OPPORTUNITY'), deal.get('CURRENCY_ID'))}",
f"Ответственный: <code>{assigned}</code>",
f"Дата создания: <code>{date}</code>"
]
if comments:
lines.extend(["", f"<b>Комментарий:</b>\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"<b>История сделки #{html.escape(deal_id)}</b>"]
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"• <b>{name}</b>",
f" Стадия: <code>{stage}</code>",
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"
)
],
]
)
+3
View File
@@ -0,0 +1,3 @@
from .app import create_app
__all__ = ["create_app"]
+69
View File
@@ -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
+98
View File
@@ -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,
)
+118
View File
@@ -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()
+57
View File
@@ -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),
)
+14
View File
@@ -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"))
+39
View File
@@ -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
+7
View File
@@ -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
+67
View File
@@ -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
+38
View File
@@ -0,0 +1,38 @@
<!doctype html>
<html lang="ru">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>Привязка Telegram</title>
<style>
body {
font: 16px sans-serif;
max-width: 560px;
margin: 48px auto;
padding: 0 20px;
}
a {
display: inline-block;
padding: 12px 18px;
color: white;
background: #168acd;
border-radius: 8px;
text-decoration: none;
}
small {
display: block;
margin-top: 16px;
color: #666;
}
</style>
</head>
<body>
<h1>Привязка Telegram</h1>
<p>{{ user.name }}, откройте бота и подтвердите привязку.</p>
<a href="{{ link.url }}" target="_blank" rel="noopener">Открыть Telegram</a>
<small>Ссылка одноразовая и действует до {{ link.expires_at.strftime('%H:%M
UTC') }}.</small>
</body>
</html>