/
AAAAAAAAAA
/
Telegram1
Обзор
Документация
Войти
/
AAAAAAAAAA
/
Telegram1
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
amocrm_server.py
786 строк
31 KB
AAAAAAAAAA
upload files
05 фев 2026, 15:43
Верифицирован
05 фев 2026, 15:43
21e7500
Код
Авторство
О чём код?
import os import json import logging import hmac import hashlib import sqlite3 from datetime import datetime from typing import Optional from dotenv import load_dotenv from fastapi import FastAPI, Request, Header, HTTPException, BackgroundTasks from fastapi.responses import JSONResponse from fastapi.middleware.cors import CORSMiddleware import uvicorn import aiohttp import asyncio # ===== НАСТРОЙКА ===== load_dotenv() # Конфигурация из .env SECRET_KEY = os.getenv("AMOCRM_SECRET_KEY", "") PORT = int(os.getenv("AMOCRM_WEBHOOK_PORT", "8001")) BOT_TOKEN = os.getenv("BOT_TOKEN", "8498375927:AAHAN7-tTHJwaBumH71qDs3KabaB5PNoLb0") ADMIN_ID = int(os.getenv("ADMIN_ID", "7771463228")) DATABASE_PATH = os.getenv("DATABASE_PATH", "webinar_registrations.db") # Настройка логирования logging.basicConfig( format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', level=logging.INFO ) logger = logging.getLogger(__name__) app = FastAPI( title="amoCRM Webhook Server", description="Сервер для приема вебхуков от amoCRM с интеграцией в Telegram бота", version="2.0" ) # Настройка CORS app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) # ===== ВСПОМОГАТЕЛЬНЫЕ ФУНКЦИИ ===== def get_db_connection(): """Создает подключение к базе данных""" return sqlite3.connect(DATABASE_PATH, check_same_thread=False) def init_database(): """Инициализация таблиц для amoCRM если их нет""" conn = get_db_connection() cursor = conn.cursor() # Таблица для детальных логов вебхуков cursor.execute(''' CREATE TABLE IF NOT EXISTS amocrm_webhook_details ( id INTEGER PRIMARY KEY AUTOINCREMENT, webhook_id TEXT, account_id TEXT, account_subdomain TEXT, event_type TEXT NOT NULL, entity_type TEXT, entity_id TEXT, data TEXT NOT NULL, headers TEXT, signature_verified BOOLEAN DEFAULT FALSE, processing_status TEXT DEFAULT 'pending', processed_at TIMESTAMP, error_message TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ) ''') # Таблица для обработки лидов cursor.execute(''' CREATE TABLE IF NOT EXISTS amocrm_leads ( id INTEGER PRIMARY KEY AUTOINCREMENT, lead_id INTEGER NOT NULL UNIQUE, name TEXT, status_id INTEGER, status_name TEXT, pipeline_id INTEGER, pipeline_name TEXT, responsible_user_id INTEGER, responsible_user_name TEXT, price INTEGER, created_at TIMESTAMP, updated_at TIMESTAMP, last_webhook_at TIMESTAMP, webhook_count INTEGER DEFAULT 0, data TEXT, synced_with_bot BOOLEAN DEFAULT FALSE ) ''') # Таблица для обработки контактов cursor.execute(''' CREATE TABLE IF NOT EXISTS amocrm_contacts ( id INTEGER PRIMARY KEY AUTOINCREMENT, contact_id INTEGER NOT NULL UNIQUE, name TEXT, first_name TEXT, last_name TEXT, phone TEXT, email TEXT, responsible_user_id INTEGER, created_at TIMESTAMP, updated_at TIMESTAMP, last_webhook_at TIMESTAMP, webhook_count INTEGER DEFAULT 0, data TEXT, synced_with_bot BOOLEAN DEFAULT FALSE ) ''') # Таблица для событий, требующих уведомления бота cursor.execute(''' CREATE TABLE IF NOT EXISTS amocrm_bot_notifications ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_type TEXT NOT NULL, entity_type TEXT NOT NULL, entity_id TEXT NOT NULL, notification_data TEXT NOT NULL, status TEXT DEFAULT 'pending', sent_to_bot BOOLEAN DEFAULT FALSE, bot_response TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, processed_at TIMESTAMP ) ''') conn.commit() conn.close() logger.info("✅ Таблицы amoCRM инициализированы") def verify_signature(payload: bytes, signature: str) -> bool: """Проверка подписи вебхука от amoCRM""" if not SECRET_KEY: logger.warning("⚠️ SECRET_KEY не установлен, проверка подписи отключена") return True if not signature: logger.warning("⚠️ Подпись отсутствует в запросе") return False try: # amoCRM использует HMAC-SHA256 expected_signature = hmac.new( SECRET_KEY.encode('utf-8'), payload, hashlib.sha256 ).hexdigest() return hmac.compare_digest(expected_signature, signature) except Exception as e: logger.error(f"❌ Ошибка проверки подписи: {e}") return False def save_webhook_to_db(webhook_data: dict, headers: dict, signature_verified: bool): """Сохранение вебхука в базу данных""" try: conn = get_db_connection() cursor = conn.cursor() # Извлекаем основные данные account = webhook_data.get('account', {}) event_type = webhook_data.get('event_type', 'unknown') cursor.execute(''' INSERT INTO amocrm_webhook_details (webhook_id, account_id, account_subdomain, event_type, entity_type, entity_id, data, headers, signature_verified) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ''', ( webhook_data.get('webhook_id'), account.get('id'), account.get('subdomain'), event_type, webhook_data.get('entity_type'), webhook_data.get('entity_id'), json.dumps(webhook_data, ensure_ascii=False), json.dumps(dict(headers), ensure_ascii=False), signature_verified )) webhook_detail_id = cursor.lastrowid # Обработка в зависимости от типа события if event_type.startswith('lead'): process_lead_webhook(webhook_data, webhook_detail_id, cursor) elif event_type.startswith('contact'): process_contact_webhook(webhook_data, webhook_detail_id, cursor) elif event_type.startswith('task'): process_task_webhook(webhook_data, webhook_detail_id, cursor) conn.commit() conn.close() logger.info(f"✅ Вебхук сохранен в БД (ID: {webhook_detail_id})") return webhook_detail_id except Exception as e: logger.error(f"❌ Ошибка сохранения вебхука в БД: {e}") return None def process_lead_webhook(webhook_data: dict, webhook_id: int, cursor): """Обработка вебхуков связанных с лидами""" try: event_type = webhook_data.get('event_type', '') leads_data = webhook_data.get('leads', {}) if 'add' in leads_data and leads_data['add']: for lead in leads_data['add']: save_or_update_lead(lead, cursor) # Создаем уведомление для бота create_bot_notification( 'new_lead', 'lead', lead.get('id'), { 'lead_id': lead.get('id'), 'lead_name': lead.get('name', 'Без названия'), 'price': lead.get('price'), 'status_id': lead.get('status_id'), 'webhook_id': webhook_id }, cursor ) elif 'status' in leads_data and leads_data['status']: for lead in leads_data['status']: update_lead_status(lead, cursor) create_bot_notification( 'lead_status_changed', 'lead', lead.get('id'), { 'lead_id': lead.get('id'), 'old_status_id': lead.get('old_status_id'), 'new_status_id': lead.get('status_id'), 'webhook_id': webhook_id }, cursor ) elif 'update' in leads_data and leads_data['update']: for lead in leads_data['update']: update_lead_data(lead, cursor) except Exception as e: logger.error(f"❌ Ошибка обработки лида: {e}") def save_or_update_lead(lead_data: dict, cursor): """Сохранение или обновление лида в БД""" try: cursor.execute(''' INSERT OR REPLACE INTO amocrm_leads (lead_id, name, status_id, pipeline_id, price, responsible_user_id, created_at, updated_at, last_webhook_at, webhook_count, data) VALUES (?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP, COALESCE((SELECT webhook_count FROM amocrm_leads WHERE lead_id = ?), 0) + 1, ?) ''', ( lead_data.get('id'), lead_data.get('name'), lead_data.get('status_id'), lead_data.get('pipeline_id'), lead_data.get('price'), lead_data.get('responsible_user_id'), datetime.fromtimestamp(lead_data.get('created_at')).isoformat() if lead_data.get('created_at') else None, datetime.fromtimestamp(lead_data.get('updated_at')).isoformat() if lead_data.get('updated_at') else None, lead_data.get('id'), json.dumps(lead_data, ensure_ascii=False) )) except Exception as e: logger.error(f"❌ Ошибка сохранения лида: {e}") def process_contact_webhook(webhook_data: dict, webhook_id: int, cursor): """Обработка вебхуков связанных с контактами""" try: event_type = webhook_data.get('event_type', '') contacts_data = webhook_data.get('contacts', {}) if 'add' in contacts_data and contacts_data['add']: for contact in contacts_data['add']: save_or_update_contact(contact, cursor) create_bot_notification( 'new_contact', 'contact', contact.get('id'), { 'contact_id': contact.get('id'), 'contact_name': contact.get('name', 'Без имени'), 'webhook_id': webhook_id }, cursor ) except Exception as e: logger.error(f"❌ Ошибка обработки контакта: {e}") def save_or_update_contact(contact_data: dict, cursor): """Сохранение или обновление контакта в БД""" try: # Извлекаем телефоны и emails custom_fields = contact_data.get('custom_fields', []) phones = [] emails = [] for field in custom_fields: if field.get('code') == 'PHONE': phones = [p.get('value') for p in field.get('values', [])] elif field.get('code') == 'EMAIL': emails = [e.get('value') for e in field.get('values', [])] cursor.execute(''' INSERT OR REPLACE INTO amocrm_contacts (contact_id, name, first_name, last_name, phone, email, responsible_user_id, created_at, updated_at, last_webhook_at, webhook_count, data) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP, COALESCE((SELECT webhook_count FROM amocrm_contacts WHERE contact_id = ?), 0) + 1, ?) ''', ( contact_data.get('id'), contact_data.get('name'), contact_data.get('first_name'), contact_data.get('last_name'), json.dumps(phones), json.dumps(emails), contact_data.get('responsible_user_id'), datetime.fromtimestamp(contact_data.get('created_at')).isoformat() if contact_data.get('created_at') else None, datetime.fromtimestamp(contact_data.get('updated_at')).isoformat() if contact_data.get('updated_at') else None, contact_data.get('id'), json.dumps(contact_data, ensure_ascii=False) )) except Exception as e: logger.error(f"❌ Ошибка сохранения контакта: {e}") def process_task_webhook(webhook_data: dict, webhook_id: int, cursor): """Обработка вебхуков связанных с задачами""" try: event_type = webhook_data.get('event_type', '') tasks_data = webhook_data.get('tasks', {}) if 'add' in tasks_data and tasks_data['add']: for task in tasks_data['add']: create_bot_notification( 'new_task', 'task', task.get('id'), { 'task_id': task.get('id'), 'task_text': task.get('text', 'Без описания'), 'complete_till': task.get('complete_till'), 'webhook_id': webhook_id }, cursor ) except Exception as e: logger.error(f"❌ Ошибка обработки задачи: {e}") def create_bot_notification(event_type: str, entity_type: str, entity_id: str, data: dict, cursor): """Создание уведомления для отправки в Telegram бот""" try: cursor.execute(''' INSERT INTO amocrm_bot_notifications (event_type, entity_type, entity_id, notification_data) VALUES (?, ?, ?, ?) ''', ( event_type, entity_type, str(entity_id), json.dumps(data, ensure_ascii=False) )) except Exception as e: logger.error(f"❌ Ошибка создания уведомления: {e}") async def send_to_telegram_bot(notification_data: dict): """Отправка уведомления в Telegram бота""" try: bot_api_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendMessage" message = format_notification_message(notification_data) async with aiohttp.ClientSession() as session: async with session.post(bot_api_url, json={ 'chat_id': ADMIN_ID, 'text': message, 'parse_mode': 'HTML', 'disable_web_page_preview': True }) as response: if response.status == 200: result = await response.json() logger.info(f"✅ Уведомление отправлено в Telegram") return result else: error_text = await response.text() logger.error(f"❌ Ошибка отправки в Telegram: {error_text}") return None except Exception as e: logger.error(f"❌ Ошибка подключения к Telegram API: {e}") return None def format_notification_message(notification_data: dict) -> str: """Форматирование сообщения для Telegram""" event_type = notification_data.get('event_type', 'unknown') if event_type == 'new_lead': return f"🎯 <b>НОВЫЙ ЛИД В AMOCRM!</b>\n\n" \ f"🆔 <b>ID:</b> {notification_data.get('lead_id')}\n" \ f"📝 <b>Название:</b> {notification_data.get('lead_name')}\n" \ f"💰 <b>Стоимость:</b> {notification_data.get('price', 0)} руб.\n" \ f"⏰ <b>Время:</b> {datetime.now().strftime('%H:%M:%S')}" elif event_type == 'lead_status_changed': return f"🔄 <b>ИЗМЕНЕНИЕ СТАТУСА ЛИДА</b>\n\n" \ f"🆔 <b>ID лида:</b> {notification_data.get('lead_id')}\n" \ f"📊 <b>Новый статус:</b> {notification_data.get('new_status_id')}\n" \ f"⏰ <b>Время:</b> {datetime.now().strftime('%H:%M:%S')}" elif event_type == 'new_contact': return f"👤 <b>НОВЫЙ КОНТАКТ В AMOCRM</b>\n\n" \ f"🆔 <b>ID:</b> {notification_data.get('contact_id')}\n" \ f"📛 <b>Имя:</b> {notification_data.get('contact_name')}\n" \ f"⏰ <b>Время:</b> {datetime.now().strftime('%H:%M:%S')}" elif event_type == 'new_task': return f"📋 <b>НОВАЯ ЗАДАЧА В AMOCRM</b>\n\n" \ f"🆔 <b>ID задачи:</b> {notification_data.get('task_id')}\n" \ f"📝 <b>Описание:</b> {notification_data.get('task_text')[:100]}...\n" \ f"⏰ <b>Время:</b> {datetime.now().strftime('%H:%M:%S')}" else: return f"🔔 <b>СОБЫТИЕ ИЗ AMOCRM</b>\n\n" \ f"📋 <b>Тип:</b> {event_type}\n" \ f"⏰ <b>Время:</b> {datetime.now().strftime('%H:%M:%S')}" async def process_pending_notifications(): """Фоновая задача для обработки ожидающих уведомлений""" try: conn = get_db_connection() cursor = conn.cursor() cursor.execute(''' SELECT id, notification_data FROM amocrm_bot_notifications WHERE status = 'pending' AND sent_to_bot = FALSE ORDER BY created_at ASC LIMIT 10 ''') pending = cursor.fetchall() for notification_id, notification_data_json in pending: try: notification_data = json.loads(notification_data_json) # Отправляем в Telegram result = await send_to_telegram_bot(notification_data) # Обновляем статус cursor.execute(''' UPDATE amocrm_bot_notifications SET status = ?, sent_to_bot = ?, bot_response = ?, processed_at = CURRENT_TIMESTAMP WHERE id = ? ''', ( 'sent' if result else 'failed', True if result else False, json.dumps(result) if result else None, notification_id )) conn.commit() if result: logger.info(f"✅ Уведомление #{notification_id} отправлено") else: logger.warning(f"⚠️ Не удалось отправить уведомление #{notification_id}") except Exception as e: logger.error(f"❌ Ошибка обработки уведомления #{notification_id}: {e}") cursor.execute(''' UPDATE amocrm_bot_notifications SET status = 'error', error_message = ? WHERE id = ? ''', (str(e), notification_id)) conn.commit() conn.close() except Exception as e: logger.error(f"❌ Ошибка в фоновой задаче: {e}") # ===== API ЭНДПОИНТЫ ===== @app.get("/") async def root(): """Корневой эндпоинт""" return { "status": "running", "service": "amocrm-webhook-server", "version": "2.0", "timestamp": datetime.now().isoformat(), "endpoints": { "webhook": "/webhook/amocrm", "health": "/health", "stats": "/stats", "pending_notifications": "/notifications/pending" } } @app.get("/health") async def health_check(): """Проверка работоспособности сервера и базы данных""" try: conn = get_db_connection() cursor = conn.cursor() cursor.execute("SELECT COUNT(*) FROM amocrm_webhook_details") webhook_count = cursor.fetchone()[0] cursor.execute("SELECT COUNT(*) FROM amocrm_bot_notifications WHERE status='pending'") pending_count = cursor.fetchone()[0] conn.close() return { "status": "healthy", "database": "connected", "webhooks_received": webhook_count, "pending_notifications": pending_count, "timestamp": datetime.now().isoformat() } except Exception as e: return { "status": "unhealthy", "error": str(e), "timestamp": datetime.now().isoformat() } @app.get("/stats") async def get_statistics(): """Получение статистики по вебхукам""" try: conn = get_db_connection() cursor = conn.cursor() # Общая статистика cursor.execute("SELECT COUNT(*) FROM amocrm_webhook_details") total_webhooks = cursor.fetchone()[0] cursor.execute("SELECT COUNT(*) FROM amocrm_leads") total_leads = cursor.fetchone()[0] cursor.execute("SELECT COUNT(*) FROM amocrm_contacts") total_contacts = cursor.fetchone()[0] # Статистика по событиям cursor.execute(''' SELECT event_type, COUNT(*) as count FROM amocrm_webhook_details GROUP BY event_type ORDER BY count DESC ''') events_stats = cursor.fetchall() # Статистика по уведомлениям cursor.execute(''' SELECT status, COUNT(*) as count FROM amocrm_bot_notifications GROUP BY status ''') notifications_stats = cursor.fetchall() conn.close() return { "statistics": { "total_webhooks": total_webhooks, "total_leads": total_leads, "total_contacts": total_contacts, "events": [{"type": row[0], "count": row[1]} for row in events_stats], "notifications": [{"status": row[0], "count": row[1]} for row in notifications_stats], "server_time": datetime.now().isoformat() } } except Exception as e: raise HTTPException(status_code=500, detail=f"Ошибка получения статистики: {str(e)}") @app.get("/notifications/pending") async def get_pending_notifications(): """Получение списка ожидающих уведомлений""" try: conn = get_db_connection() cursor = conn.cursor() cursor.execute(''' SELECT id, event_type, entity_type, entity_id, created_at FROM amocrm_bot_notifications WHERE status = 'pending' ORDER BY created_at ASC LIMIT 50 ''') pending = cursor.fetchall() conn.close() return { "pending_notifications": [ { "id": row[0], "event_type": row[1], "entity_type": row[2], "entity_id": row[3], "created_at": row[4] } for row in pending ], "count": len(pending) } except Exception as e: raise HTTPException(status_code=500, detail=f"Ошибка получения уведомлений: {str(e)}") @app.post("/webhook/amocrm") async def handle_amocrm_webhook( request: Request, background_tasks: BackgroundTasks, x_signature: str = Header(None), x_amo_signature: str = Header(None) ): """Основной обработчик вебхуков от amoCRM""" try: # Получаем сырые данные для проверки подписи body_bytes = await request.body() body_str = body_bytes.decode('utf-8') # Проверяем подпись signature = x_signature or x_amo_signature signature_verified = verify_signature(body_bytes, signature) if not signature_verified: logger.warning(f"❌ Неверная подпись вебхука от {request.client.host}") raise HTTPException(status_code=401, detail="Invalid signature") # Парсим JSON try: data = json.loads(body_str) except json.JSONDecodeError as e: logger.error(f"❌ Ошибка парсинга JSON: {e}") raise HTTPException(status_code=400, detail="Invalid JSON format") # Логируем получение вебхука account_id = data.get('account', {}).get('id', 'unknown') event_type = data.get('event_type', 'unknown') logger.info(f"📥 Вебхук от amoCRM | Аккаунт: {account_id} | Событие: {event_type}") # Сохраняем в базу данных headers_dict = dict(request.headers) webhook_id = save_webhook_to_db(data, headers_dict, signature_verified) if not webhook_id: logger.error("❌ Не удалось сохранить вебхук в БД") raise HTTPException(status_code=500, detail="Failed to save webhook") # Добавляем фоновую задачу для обработки уведомлений background_tasks.add_task(process_pending_notifications) # Немедленно обрабатываем важные события if event_type in ['lead_added', 'contact_added']: background_tasks.add_task(process_pending_notifications) return JSONResponse( status_code=200, content={ "status": "success", "message": "Webhook received and processed", "webhook_id": webhook_id, "event_type": event_type, "signature_verified": signature_verified, "timestamp": datetime.now().isoformat() } ) except HTTPException: raise except Exception as e: logger.error(f"❌ Критическая ошибка обработки вебхука: {e}") return JSONResponse( status_code=500, content={ "status": "error", "message": "Internal server error", "timestamp": datetime.now().isoformat() } ) @app.post("/test/webhook") async def test_webhook_endpoint(): """Тестовый эндпоинт для проверки работы (без проверки подписи)""" test_data = { "event_type": "lead_added", "account": { "id": 123456, "subdomain": "test" }, "leads": { "add": [{ "id": 999999, "name": "Тестовый лид", "price": 10000, "status_id": 142, "pipeline_id": 123, "responsible_user_id": 12345, "created_at": int(datetime.now().timestamp()), "updated_at": int(datetime.now().timestamp()) }] }, "timestamp": int(datetime.now().timestamp()) } # Сохраняем тестовый вебхук webhook_id = save_webhook_to_db(test_data, {"test": "true"}, True) # Отправляем тестовое уведомление background_tasks = BackgroundTasks() background_tasks.add_task(process_pending_notifications) return { "status": "test_webhook_processed", "webhook_id": webhook_id, "test_data": test_data, "timestamp": datetime.now().isoformat() } # ===== ЗАПУСК СЕРВЕРА ===== @app.on_event("startup") async def startup_event(): """Действия при запуске сервера""" print("=" * 60) print("🚀 ЗАПУСК AMOCRM WEBHOOK СЕРВЕРА v2.0") print("=" * 60) print(f"📌 Порт: {PORT}") print(f"🔐 Проверка подписи: {'ВКЛЮЧЕНА' if SECRET_KEY else 'ВЫКЛЮЧЕНА'}") print(f"🤖 Интеграция с Telegram ботом: {'ДА' if BOT_TOKEN else 'НЕТ'}") print(f"👑 Администратор: {ADMIN_ID}") print(f"💾 База данных: {DATABASE_PATH}") print(f"🌐 Webhook URL: http://ваш-домен:{PORT}/webhook/amocrm") print(f"🩺 Health check: http://localhost:{PORT}/health") print(f"📊 Статистика: http://localhost:{PORT}/stats") print("=" * 60) # Инициализация базы данных init_database() # Запускаем периодическую обработку уведомлений asyncio.create_task(periodic_notification_processor()) async def periodic_notification_processor(): """Периодическая обработка уведомлений каждые 30 секунд""" while True: try: await process_pending_notifications() except Exception as e: logger.error(f"❌ Ошибка в периодической обработке: {e}") await asyncio.sleep(30) # Проверяем каждые 30 секунд if __name__ == "__main__": # Создаем папку для логов если её нет os.makedirs("amocrm_logs", exist_ok=True) # Запускаем сервер uvicorn.run( app, host="0.0.0.0", port=PORT, log_level="info", access_log=True, reload=False )