- Добавлен сервис CrmMySQLService для подключения к MySQL БД vtiger CRM - Обновлён get_draft() для прямого SQL запроса к MySQL вместо webservice API - Получение cf_2624 и всех данных контакта из MySQL - Обновлена документация и SQL файлы для n8n - Добавлено логирование для отладки Преимущества: - Проще: один SQL запрос вместо цепочки HTTP запросов - Быстрее: прямой запрос к БД - Надёжнее: не зависит от webservice API - Актуальнее: всегда свежие данные из БД
118 lines
4.0 KiB
Python
118 lines
4.0 KiB
Python
"""
|
||
CRM MySQL Service - Подключение к MySQL БД vtiger CRM
|
||
"""
|
||
import aiomysql
|
||
from typing import Optional, Dict, Any, List
|
||
from ..config import settings
|
||
import logging
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
class CrmMySQLService:
|
||
"""Сервис для работы с MySQL БД vtiger CRM"""
|
||
|
||
def __init__(self):
|
||
self.pool: Optional[aiomysql.Pool] = None
|
||
|
||
async def connect(self):
|
||
"""Подключение к MySQL БД vtiger CRM"""
|
||
try:
|
||
self.pool = await aiomysql.create_pool(
|
||
host=settings.mysql_crm_host,
|
||
port=settings.mysql_crm_port,
|
||
user=settings.mysql_crm_user,
|
||
password=settings.mysql_crm_password,
|
||
db=settings.mysql_crm_db,
|
||
autocommit=True,
|
||
minsize=1,
|
||
maxsize=5
|
||
)
|
||
logger.info(f"✅ MySQL CRM DB connected: {settings.mysql_crm_host}:{settings.mysql_crm_port}/{settings.mysql_crm_db}")
|
||
except Exception as e:
|
||
logger.error(f"❌ MySQL CRM DB connection error: {e}")
|
||
raise
|
||
|
||
async def fetch_one(self, query: str, *args) -> Optional[Dict[str, Any]]:
|
||
"""
|
||
Выполнить SQL запрос и вернуть одну запись
|
||
|
||
Args:
|
||
query: SQL запрос с плейсхолдерами %s
|
||
*args: Параметры для запроса
|
||
|
||
Returns:
|
||
Dict с данными или None если не найдено
|
||
"""
|
||
if not self.pool:
|
||
await self.connect()
|
||
|
||
try:
|
||
async with self.pool.acquire() as conn:
|
||
async with conn.cursor(aiomysql.DictCursor) as cursor:
|
||
await cursor.execute(query, args)
|
||
result = await cursor.fetchone()
|
||
return dict(result) if result else None
|
||
except Exception as e:
|
||
logger.error(f"❌ Error executing query: {e}")
|
||
raise
|
||
|
||
async def fetch_all(self, query: str, *args) -> List[Dict[str, Any]]:
|
||
"""
|
||
Выполнить SQL запрос и вернуть все записи
|
||
|
||
Args:
|
||
query: SQL запрос с плейсхолдерами %s
|
||
*args: Параметры для запроса
|
||
|
||
Returns:
|
||
List[Dict] с данными
|
||
"""
|
||
if not self.pool:
|
||
await self.connect()
|
||
|
||
try:
|
||
async with self.pool.acquire() as conn:
|
||
async with conn.cursor(aiomysql.DictCursor) as cursor:
|
||
await cursor.execute(query, args)
|
||
results = await cursor.fetchall()
|
||
return [dict(row) for row in results] if results else []
|
||
except Exception as e:
|
||
logger.error(f"❌ Error executing query: {e}")
|
||
raise
|
||
|
||
async def execute(self, query: str, *args) -> int:
|
||
"""
|
||
Выполнить SQL запрос (INSERT, UPDATE, DELETE)
|
||
|
||
Args:
|
||
query: SQL запрос с плейсхолдерами %s
|
||
*args: Параметры для запроса
|
||
|
||
Returns:
|
||
Количество затронутых строк
|
||
"""
|
||
if not self.pool:
|
||
await self.connect()
|
||
|
||
try:
|
||
async with self.pool.acquire() as conn:
|
||
async with conn.cursor() as cursor:
|
||
await cursor.execute(query, args)
|
||
return cursor.rowcount
|
||
except Exception as e:
|
||
logger.error(f"❌ Error executing query: {e}")
|
||
raise
|
||
|
||
async def close(self):
|
||
"""Закрыть пул подключений"""
|
||
if self.pool:
|
||
self.pool.close()
|
||
await self.pool.wait_closed()
|
||
logger.info("MySQL CRM DB pool closed")
|
||
|
||
|
||
# Глобальный экземпляр
|
||
crm_mysql_service = CrmMySQLService()
|
||
|