import os import logging import asyncio import pymysql import time import html import json import re import httpx import requests from requests.adapters import HTTPAdapter from urllib3.util.retry import Retry from urllib3 import Timeout as Urllib3Timeout import traceback import io import sys from datetime import datetime, timezone from dotenv import load_dotenv from formatter import get_event_text # Настройка anyio для правильной работы в отдельном потоке # Устанавливаем правильный backend для anyio (используем asyncio) try: import anyio # Устанавливаем backend для anyio явно if not hasattr(anyio, '_backend'): # Используем asyncio backend os.environ.setdefault('ANYIO_BACKEND', 'asyncio') except ImportError: pass # Загрузка переменных окружения load_dotenv() # Настройки из переменных окружения BOT_TOKEN = os.getenv('POSTER_BOT_TOKEN') RESPONDER_BOT_NAME = os.getenv('RESPONDER_BOT_NAME') CHANNEL_ID = os.getenv('CHANNEL_ID') MDB_HOST = os.getenv('MDB_HOST') MDB_USER = os.getenv('MDB_USER') MDB_PW = os.getenv('MDB_PW') MDBASE = os.getenv('MDBASE') LOG_FILE = os.getenv('LOG_FILE', 'post_publisher.log') ORG_MESSAGE_TAG = os.getenv('ORG_MESSAGE_TAG', '') # Хештег для определения организационных сообщений # Параметры длины сообщений MAX_CAPTION_LENGTH = int(os.getenv('MAX_CAPTION_LENGTH', 1000)) MAX_TEXT_LENGTH = int(os.getenv('MAX_TEXT_LENGTH', 4000)) # Режим публикации без звука PUBLISH_SILENTLY = os.getenv('PUBLISH_SILENTLY', 'false').lower() in ('true', '1', 'yes', 'on') # Использование бота подписки USE_SUBSCRIPTION_BOT = os.getenv('USE_SUBSCRIPTION_BOT', 'true').lower() in ('true', '1', 'yes', 'on') # Режим отладочного логирования данных LOG_DEBUG_DATA = os.getenv('LOG_DEBUG_DATA', 'false').lower() in ('true', '1', 'yes', 'on') # Настройка логирования logger = logging.getLogger('TG_Poster') logger.setLevel(logging.INFO) formatter = logging.Formatter('[%(asctime)s] [%(name)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S') console_handler = logging.StreamHandler() console_handler.setFormatter(formatter) logger.addHandler(console_handler) class RetryFileHandler(logging.FileHandler): def emit(self, record): for _ in range(5): try: super().emit(record) return except (IOError, PermissionError): time.sleep(0.5) print(f"Failed to write to log file after 5 attempts: {record.msg}") file_handler = RetryFileHandler(LOG_FILE, encoding='utf-8') file_handler.setFormatter(formatter) logger.addHandler(file_handler) # Класс для обработки RetryAfter (FloodWait 429) class RetryAfterException(Exception): """Исключение для обработки FloodWait (429) от Telegram API""" def __init__(self, retry_after): self.retry_after = retry_after super().__init__(f"RetryAfter: {retry_after}") def prepare_text(text): """Подготовка текстовых полей с экранированием HTML-сущностей""" if not text: return "" # Экранируем специальные символы HTML text = html.escape(str(text)) return text def has_org_tag(text): """Проверяет, содержит ли текст организационный хештег""" if not text or not ORG_MESSAGE_TAG: return False # Ищем хештег (с учетом возможных пробелов и других символов) pattern = r'#?' + re.escape(ORG_MESSAGE_TAG) + r'\b' return bool(re.search(pattern, text, re.IGNORECASE)) def log_telegram_error(response, context=""): """Логирует полную информацию об ошибке от Telegram API Работает как с httpx.Response, так и с requests.Response """ try: if response is not None: logger.error(f"{context} Статус код: {response.status_code}") logger.error(f"{context} Заголовки ответа: {dict(response.headers)}") try: # Работает для обоих типов ответов (httpx и requests) error_data = response.json() logger.error(f"{context} Полный ответ от Telegram API: {json.dumps(error_data, ensure_ascii=False, indent=2)}") except: # Для requests используем text, для httpx тоже text response_text = getattr(response, 'text', str(response.content[:500]) if hasattr(response, 'content') else 'N/A') logger.error(f"{context} Текст ответа: {response_text}") else: logger.error(f"{context} Ответ от Telegram API отсутствует (None)") except Exception as e: logger.error(f"{context} Ошибка при логировании информации об ошибке: {e}") async def download_image(image_url, httpx_client): """Скачивает изображение по URL и возвращает его содержимое в виде bytes""" try: logger.info(f"Скачивание изображения: {image_url}") download_start_time = time.time() response = await httpx_client.get(image_url, timeout=httpx.Timeout(connect=10.0, read=60.0, write=30.0, pool=10.0)) response.raise_for_status() download_duration = time.time() - download_start_time image_data = response.content logger.info(f"Изображение скачано за {download_duration:.2f} секунд, размер: {len(image_data)} байт") return image_data except Exception as e: logger.error(f"Ошибка при скачивании изображения {image_url}: {e}") raise def _send_photo_sync(send_photo_url, image_data, caption, chat_id, disable_notification, timeout): """ Синхронная функция для отправки фото через requests. Вызывается из отдельного потока для избежания проблем с event loop в gunicorn. Args: send_photo_url: URL для отправки фото в Telegram API image_data: Байты изображения caption: Подпись к фото chat_id: ID чата/канала disable_notification: Отключить уведомления timeout: Кортеж (connect_timeout, read_timeout, write_timeout) или число для общего таймаута Returns: requests.Response объект """ import socket from urllib3.connection import HTTPSConnection, HTTPConnection try: # Используем requests для стабильной работы в gunicorn # requests более надежен для синхронных операций в контексте WSGI image_file = io.BytesIO(image_data) files = { 'photo': ('image.jpg', image_file, 'image/jpeg') } data = { 'chat_id': chat_id, 'caption': caption, 'parse_mode': 'HTML', 'disable_notification': str(disable_notification).lower() } # Извлекаем таймауты из кортежа if isinstance(timeout, tuple): if len(timeout) >= 3: connect_timeout, read_timeout, write_timeout = timeout[0], timeout[1], timeout[2] elif len(timeout) == 2: connect_timeout, read_timeout = timeout[0], timeout[1] write_timeout = read_timeout # Используем read timeout для write else: connect_timeout = read_timeout = write_timeout = timeout[0] else: connect_timeout = read_timeout = write_timeout = timeout # Создаем сессию с кастомным адаптером session = requests.Session() # Создаем кастомный HTTPAdapter с установкой socket timeout # Используем максимальный timeout для всех операций (connect, read, write) class CustomHTTPAdapter(HTTPAdapter): def init_poolmanager(self, *args, **kwargs): # Устанавливаем socket timeout через socket_options socket_options = kwargs.get('socket_options', []) socket_options.append((socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)) kwargs['socket_options'] = socket_options # Устанавливаем timeout для пула соединений # Используем максимальный таймаут для всех операций max_timeout = max(connect_timeout, read_timeout, write_timeout) if 'timeout' not in kwargs: kwargs['timeout'] = Urllib3Timeout( connect=connect_timeout, read=max_timeout ) return super().init_poolmanager(*args, **kwargs) # Монтируем кастомный адаптер adapter = CustomHTTPAdapter(max_retries=Retry(total=0)) session.mount('http://', adapter) session.mount('https://', adapter) # Отправляем запрос с увеличенным таймаутом # Используем максимальный таймаут для всех операций max_timeout = max(connect_timeout, read_timeout, write_timeout) response = session.post( send_photo_url, files=files, data=data, timeout=(connect_timeout, max_timeout), # (connect, read) stream=False # Отключаем потоковую передачу для стабильности ) return response except requests.exceptions.Timeout as e: # Пробрасываем таймаут как есть, чтобы его можно было обработать выше raise except requests.exceptions.RequestException as e: # Пробрасываем другие ошибки requests как есть raise except Exception as e: # Пробрасываем остальные исключения как есть raise finally: # Закрываем сессию if 'session' in locals(): session.close() def _edit_caption_sync(edit_caption_url, chat_id, message_id, caption, timeout): """ Синхронная функция для редактирования caption через requests. Вызывается из отдельного потока для избежания проблем с event loop в gunicorn. Args: edit_caption_url: URL для редактирования caption в Telegram API chat_id: ID чата/канала message_id: ID сообщения для редактирования caption: Новый текст caption timeout: Кортеж (connect_timeout, read_timeout, write_timeout) или (connect_timeout, read_timeout) или число Returns: requests.Response объект """ try: data = { 'chat_id': chat_id, 'message_id': message_id, 'caption': caption, 'parse_mode': 'HTML' } # Преобразуем timeout в формат, который принимает requests # requests принимает либо число, либо кортеж (connect, read) if isinstance(timeout, tuple): if len(timeout) >= 3: # Используем максимальный из read и write timeout connect_timeout, read_timeout, write_timeout = timeout[0], timeout[1], timeout[2] requests_timeout = (connect_timeout, max(read_timeout, write_timeout)) elif len(timeout) == 2: requests_timeout = timeout else: requests_timeout = timeout[0] else: requests_timeout = timeout response = requests.post( edit_caption_url, json=data, timeout=requests_timeout ) return response except requests.exceptions.Timeout as e: raise except requests.exceptions.RequestException as e: raise except Exception as e: raise def _edit_text_sync(edit_text_url, chat_id, message_id, text, timeout): """ Синхронная функция для редактирования текста через requests. Вызывается из отдельного потока для избежания проблем с event loop в gunicorn. Args: edit_text_url: URL для редактирования текста в Telegram API chat_id: ID чата/канала message_id: ID сообщения для редактирования text: Новый текст сообщения timeout: Кортеж (connect_timeout, read_timeout, write_timeout) или (connect_timeout, read_timeout) или число Returns: requests.Response объект """ try: data = { 'chat_id': chat_id, 'message_id': message_id, 'text': text, 'parse_mode': 'HTML', 'disable_web_page_preview': True } # Преобразуем timeout в формат, который принимает requests # requests принимает либо число, либо кортеж (connect, read) if isinstance(timeout, tuple): if len(timeout) >= 3: # Используем максимальный из read и write timeout connect_timeout, read_timeout, write_timeout = timeout[0], timeout[1], timeout[2] requests_timeout = (connect_timeout, max(read_timeout, write_timeout)) elif len(timeout) == 2: requests_timeout = timeout else: requests_timeout = timeout[0] else: requests_timeout = timeout response = requests.post( edit_text_url, json=data, timeout=requests_timeout ) return response except requests.exceptions.Timeout as e: raise except requests.exceptions.RequestException as e: raise except Exception as e: raise async def publish_to_tg(vk_post_id): """Публикация одной записи по VK post ID""" logger.info(f"Запуск публикации для записи VK ID {vk_post_id}") conn = None try: conn = pymysql.connect( host=MDB_HOST, user=MDB_USER, password=MDB_PW, database=MDBASE, charset='utf8mb4', cursorclass=pymysql.cursors.DictCursor ) with conn.cursor() as cursor: cursor.execute(""" SELECT id, vk_post_id, text, image_url, vk_post_url, is_poll, poll_question, poll_options, poll_multiple, poll_end_date, is_event FROM posts WHERE vk_post_id = %s AND marked_for_publication = True AND published_in_tg = False """, (vk_post_id,)) post = cursor.fetchone() if not post: logger.info(f"Запись VK ID {vk_post_id} не найдена или уже опубликована") return False # Создаем httpx клиент с отключенным HTTP/2 # Используем более детальные таймауты: connect, read, write, pool # Таймаут на запись увеличен для больших файлов (до 10 МБ) timeout_config = httpx.Timeout( connect=10.0, # Таймаут на подключение read=120.0, # Таймаут на чтение (увеличен для больших изображений) write=120.0, # Таймаут на запись (увеличен для больших файлов) pool=10.0 # Таймаут на получение соединения из пула ) httpx_client = httpx.AsyncClient( http2=False, timeout=timeout_config, follow_redirects=True ) # Переменная для отслеживания успешности публикации message_id = None publication_successful = False # Флаг, указывающий, было ли отправлено изображение (True) или текстовое сообщение (False) message_has_image = False try: # Проверяем наличие организационного хештега is_org_message = has_org_tag(post['text']) # Предварительный расчет строки ссылок (с запасом 20 символов для message_id) placeholder_message_id = '0' * 20 # Заполнитель для message_id # Формируем ссылки links = [] # Ссылка на оригинал в ВК if post['vk_post_url'] and (post['vk_post_url'].startswith('http://') or post['vk_post_url'].startswith('https://')): links.append(f'Оригинал в ВК') # Ссылка на подписку (только для событий и если включен бот подписки) if post['is_event'] and USE_SUBSCRIPTION_BOT: links.append(f'🔜 Подписка') # Формируем строку ссылок links_line = " | ".join(links) if links else "" links_line_length = len(links_line) # Определяем максимальную длину в зависимости от типа сообщения # Если это организационное сообщение, используем текстовый лимит image_url = post['image_url'] has_image = image_url and image_url.strip() and not is_org_message max_length = MAX_CAPTION_LENGTH if has_image else MAX_TEXT_LENGTH # Вычисляем доступную длину для текста available_length = max_length - links_line_length - 2 # -2 для символов переноса строки if available_length < 0: available_length = 0 # Получаем обработанный текст с учетом доступной длины text = get_event_text(post['text'], available_length) if post['text'] else "" # Логирование отформатированного текста для отладки (только если включено) if LOG_DEBUG_DATA: logger.info(f"Начинаем форматирование текста для записи VK ID {vk_post_id}, доступная длина: {available_length}") logger.info(f"=== Отформатированный текст для записи VK ID {vk_post_id} ===") logger.info(f"Длина текста: {len(text)}") logger.info(f"Текст (repr): {repr(text)}") logger.info(f"Текст (содержимое):") # Выводим текст построчно для лучшей читаемости в логе for i, line in enumerate(text.split('\n'), 1): logger.info(f" Строка {i}: {repr(line)}") logger.info(f"=== Конец отформатированного текста ===") # Публикация первоначального сообщения без ссылок send_photo_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendPhoto" send_message_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendMessage" if has_image and not is_org_message: try: # Логируем информацию об изображении перед отправкой logger.info(f"Отправка изображения для записи VK ID {vk_post_id}") logger.info(f"URL изображения: {image_url}") logger.info(f"Длина URL изображения: {len(image_url)} символов") logger.info(f"Длина подписи: {len(text)} символов") # Скачиваем изображение image_data = None try: image_data = await download_image(image_url, httpx_client) except Exception as download_error: logger.error(f"Ошибка при скачивании изображения для записи VK ID {vk_post_id}: {download_error}") raise # Пробрасываем ошибку, чтобы перейти к fallback # Вычисляем динамический таймаут на запись в зависимости от размера файла # Базовый таймаут 60 секунд + 10 секунд на каждый МБ (минимум 120 секунд) file_size_mb = len(image_data) / (1024 * 1024) write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0)) logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд") # Используем requests для отправки файлов через отдельный поток # Это необходимо для стабильной работы с gunicorn, который использует синхронные воркеры # requests более надежен для синхронных операций в контексте WSGI logger.info(f"Используем requests для отправки файла размером {file_size_mb:.2f} МБ через отдельный поток") # Логируем время начала запроса request_start_time = time.time() logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}") # Вызываем синхронную функцию через asyncio.to_thread() # Это позволяет избежать проблем с event loop в gunicorn try: # Используем кортеж для таймаута: (connect, read, write) # write_timeout применяется к операциям записи в сокет timeout_tuple = (10.0, write_timeout, write_timeout) # (connect, read, write) photo_response = await asyncio.to_thread( _send_photo_sync, send_photo_url, image_data, text, CHANNEL_ID, PUBLISH_SILENTLY, timeout_tuple ) except Exception as sync_error: # Обработка ошибок от синхронной функции request_duration = time.time() - request_start_time error_msg = str(sync_error) # Проверяем, является ли это таймаутом is_timeout = ( 'Timeout' in error_msg or isinstance(sync_error, requests.exceptions.Timeout) or 'timeout' in error_msg.lower() ) if is_timeout: timeout_type = type(sync_error).__name__ logger.error(f"ТАЙМАУТ ({timeout_type}) при отправке изображения (requests) для записи VK ID {vk_post_id}") logger.error(f"Время до таймаута: {request_duration:.2f} секунд") logger.error(f"Тип исключения: {timeout_type}") logger.error(f"Сообщение об ошибке: {error_msg}") logger.error(f"Полная информация об исключении: {repr(sync_error)}") logger.error(f"URL изображения: {image_url}") logger.error(f"Размер файла: {file_size_mb:.2f} МБ") logger.error(f"Запрос, вызвавший таймаут: POST {send_photo_url}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise # Пробрасываем таймаут наверх else: # Другие ошибки logger.error(f"Ошибка при отправке изображения (requests) для записи VK ID {vk_post_id}: {error_msg}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise request_duration = time.time() - request_start_time logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд") # Обработка ответа от requests if photo_response.status_code == 200: try: photo_result = photo_response.json() except Exception as json_error: logger.error(f"Ошибка при парсинге JSON ответа для записи VK ID {vk_post_id}: {json_error}") logger.error(f"Текст ответа: {photo_response.text[:500]}") raise Exception(f"Ошибка парсинга JSON ответа: {json_error}") if photo_result.get('ok'): message_data = photo_result.get('result', {}) message_id = message_data.get('message_id') message_has_image = True # Успешно отправили изображение else: # Ошибка в ответе API - логируем полную информацию logger.error(f"[Отправка фото для VK ID {vk_post_id}] Статус код: {photo_response.status_code}") logger.error(f"[Отправка фото для VK ID {vk_post_id}] Полный ответ от Telegram API: {json.dumps(photo_result, ensure_ascii=False, indent=2)}") error_code = photo_result.get('error_code') if error_code == 429: retry_after = photo_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) raise Exception(f"API error: {photo_result.get('description', 'Unknown error')}") else: logger.error(f"[Отправка фото для VK ID {vk_post_id}] Статус код: {photo_response.status_code}") logger.error(f"[Отправка фото для VK ID {vk_post_id}] Текст ответа: {photo_response.text[:500]}") if photo_response.status_code == 429: retry_after = 60 try: error_data = photo_response.json() retry_after = error_data.get('parameters', {}).get('retry_after', 60) except: pass raise RetryAfterException(retry_after) raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text[:500]}") except RetryAfterException: raise except Exception as e: # Проверяем, является ли это таймаутом # Может быть от requests или httpx is_timeout = ( isinstance(e, (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout)) or isinstance(e, requests.exceptions.Timeout) or 'Timeout' in str(e) or 'timeout' in str(e).lower() ) if is_timeout: timeout_type = type(e).__name__ request_duration = time.time() - request_start_time if 'request_start_time' in locals() else 0 logger.error(f"ТАЙМАУТ ({timeout_type}) при отправке изображения для записи VK ID {vk_post_id}") if request_duration > 0: logger.error(f"Время до таймаута: {request_duration:.2f} секунд") logger.error(f"Тип исключения: {timeout_type}") logger.error(f"Сообщение об ошибке: {str(e)}") logger.error(f"Полная информация об исключении: {repr(e)}") logger.error(f"URL изображения: {image_url}") if image_data is not None: logger.error(f"Размер файла: {len(image_data) / (1024 * 1024):.2f} МБ") logger.error(f"Запрос, вызвавший таймаут: POST {send_photo_url}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise # Пробрасываем таймаут наверх # Если это не таймаут, пробрасываем дальше для обработки как обычная ошибка raise except Exception as e: logger.warning(f"Не удалось отправить изображение для записи VK ID {vk_post_id}: {e}. Отправляем текстовое сообщение.") message_has_image = False # Отправляем текстовое сообщение вместо изображения text_payload = { 'chat_id': CHANNEL_ID, 'text': text, 'parse_mode': 'HTML', 'disable_notification': PUBLISH_SILENTLY, 'disable_web_page_preview': True } try: text_response = await httpx_client.post(send_message_url, json=text_payload) except httpx.TimeoutException as timeout_e: logger.error(f"ТАЙМАУТ при отправке текста (fallback) для записи VK ID {vk_post_id}") logger.error(f"Тип исключения: {type(timeout_e).__name__}") logger.error(f"Сообщение об ошибке: {str(timeout_e)}") logger.error(f"Полная информация об исключении: {repr(timeout_e)}") if hasattr(timeout_e, 'request'): logger.error(f"Запрос, вызвавший таймаут: {timeout_e.request.method} {timeout_e.request.url if hasattr(timeout_e.request, 'url') else 'N/A'}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise # Пробрасываем таймаут наверх if text_response.status_code == 200: text_result = text_response.json() if text_result.get('ok'): message_data = text_result.get('result', {}) message_id = message_data.get('message_id') else: log_telegram_error(text_response, f"[Отправка текста (fallback) для VK ID {vk_post_id}]") error_code = text_result.get('error_code') if error_code == 429: retry_after = text_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) raise Exception(f"API error: {text_result.get('description', 'Unknown error')}") else: log_telegram_error(text_response, f"[Отправка текста (fallback) для VK ID {vk_post_id}]") if text_response.status_code == 429: retry_after = 60 try: error_data = text_response.json() retry_after = error_data.get('parameters', {}).get('retry_after', 60) except: pass raise RetryAfterException(retry_after) raise Exception(f"HTTP {text_response.status_code}: {text_response.text}") else: # Если это организационное сообщение или нет изображения, отправляем текстовое сообщение message_has_image = False # Отправляем текстовое сообщение text_payload = { 'chat_id': CHANNEL_ID, 'text': text, 'parse_mode': 'HTML', 'disable_notification': PUBLISH_SILENTLY, 'disable_web_page_preview': True } try: text_response = await httpx_client.post(send_message_url, json=text_payload) except httpx.TimeoutException as e: logger.error(f"ТАЙМАУТ при отправке текста для записи VK ID {vk_post_id}") logger.error(f"Тип исключения: {type(e).__name__}") logger.error(f"Сообщение об ошибке: {str(e)}") logger.error(f"Полная информация об исключении: {repr(e)}") if hasattr(e, 'request'): logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise # Пробрасываем таймаут наверх if text_response.status_code == 200: text_result = text_response.json() if text_result.get('ok'): message_data = text_result.get('result', {}) message_id = message_data.get('message_id') else: log_telegram_error(text_response, f"[Отправка текста для VK ID {vk_post_id}]") error_code = text_result.get('error_code') if error_code == 429: retry_after = text_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) raise Exception(f"API error: {text_result.get('description', 'Unknown error')}") else: log_telegram_error(text_response, f"[Отправка текста для VK ID {vk_post_id}]") if text_response.status_code == 429: retry_after = 60 try: error_data = text_response.json() retry_after = error_data.get('parameters', {}).get('retry_after', 60) except: pass raise RetryAfterException(retry_after) raise Exception(f"HTTP {text_response.status_code}: {text_response.text}") # Формируем финальный текст с ссылками if links: # Обновляем ссылку подписки с реальным message_id updated_links = [] if post['vk_post_url'] and (post['vk_post_url'].startswith('http://') or post['vk_post_url'].startswith('https://')): updated_links.append(f'Оригинал в ВК') if post['is_event'] and USE_SUBSCRIPTION_BOT: updated_links.append(f'🔜 Подписка') links_line = " | ".join(updated_links) final_text = f"{text}\n\n{links_line}" await asyncio.sleep(2) # задержка между обращениями к телеграм # Редактируем сообщение, добавляя ссылки edit_caption_url = f"https://api.telegram.org/bot{BOT_TOKEN}/editMessageCaption" edit_text_url = f"https://api.telegram.org/bot{BOT_TOKEN}/editMessageText" # Используем message_has_image для определения типа сообщения # Если было отправлено изображение, редактируем caption, иначе - текст if message_has_image: # Сообщение с изображением - редактируем caption через requests edit_caption_success = False try: # Используем requests через asyncio.to_thread() для стабильной работы в gunicorn # Увеличиваем таймаут для редактирования caption edit_timeout = (10.0, 60.0, 60.0) # (connect, read, write) edit_caption_response = await asyncio.to_thread( _edit_caption_sync, edit_caption_url, CHANNEL_ID, message_id, final_text, edit_timeout ) if edit_caption_response.status_code == 200: try: edit_result = edit_caption_response.json() except Exception as json_error: logger.error(f"Ошибка при парсинге JSON ответа при редактировании caption для записи VK ID {vk_post_id}: {json_error}") logger.error(f"Текст ответа: {edit_caption_response.text[:500]}") raise Exception(f"Ошибка парсинга JSON ответа: {json_error}") if edit_result.get('ok'): edit_caption_success = True logger.info(f"Caption успешно отредактирован для записи VK ID {vk_post_id}") else: # Ошибка в ответе API log_telegram_error(edit_caption_response, f"[Редактирование подписи для VK ID {vk_post_id}]") error_code = edit_result.get('error_code') if error_code == 429: retry_after = edit_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) # Не критичная ошибка - сообщение уже опубликовано logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: {edit_result.get('description', 'Unknown error')}. Сообщение опубликовано, но без ссылок.") else: log_telegram_error(edit_caption_response, f"[Редактирование подписи для VK ID {vk_post_id}]") if edit_caption_response.status_code == 429: retry_after = 60 try: error_data = edit_caption_response.json() retry_after = error_data.get('parameters', {}).get('retry_after', 60) except: pass raise RetryAfterException(retry_after) # Не критичная ошибка - сообщение уже опубликовано logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: HTTP {edit_caption_response.status_code}. Сообщение опубликовано, но без ссылок.") except RetryAfterException: raise except Exception as e: # Проверяем, является ли это таймаутом is_timeout = ( isinstance(e, requests.exceptions.Timeout) or 'Timeout' in str(e) or 'timeout' in str(e).lower() or 'disconnected' in str(e).lower() ) if is_timeout: logger.warning(f"Таймаут при редактировании caption для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.") else: logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.") # Не пробрасываем исключение - сообщение уже опубликовано, просто без ссылок # Это считается частичным успехом else: # Текстовое сообщение - редактируем текст через requests edit_text_success = False try: # Используем requests через asyncio.to_thread() для стабильной работы в gunicorn # Увеличиваем таймаут для редактирования текста edit_timeout = (10.0, 60.0, 60.0) # (connect, read, write) edit_text_response = await asyncio.to_thread( _edit_text_sync, edit_text_url, CHANNEL_ID, message_id, final_text, edit_timeout ) if edit_text_response.status_code == 200: try: edit_result = edit_text_response.json() except Exception as json_error: logger.error(f"Ошибка при парсинге JSON ответа при редактировании текста для записи VK ID {vk_post_id}: {json_error}") logger.error(f"Текст ответа: {edit_text_response.text[:500]}") raise Exception(f"Ошибка парсинга JSON ответа: {json_error}") if edit_result.get('ok'): edit_text_success = True logger.info(f"Текст успешно отредактирован для записи VK ID {vk_post_id}") else: # Ошибка в ответе API log_telegram_error(edit_text_response, f"[Редактирование текста для VK ID {vk_post_id}]") error_code = edit_result.get('error_code') if error_code == 429: retry_after = edit_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) # Не критичная ошибка - сообщение уже опубликовано logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: {edit_result.get('description', 'Unknown error')}. Сообщение опубликовано, но без ссылок.") else: log_telegram_error(edit_text_response, f"[Редактирование текста для VK ID {vk_post_id}]") if edit_text_response.status_code == 429: retry_after = 60 try: error_data = edit_text_response.json() retry_after = error_data.get('parameters', {}).get('retry_after', 60) except: pass raise RetryAfterException(retry_after) # Не критичная ошибка - сообщение уже опубликовано logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: HTTP {edit_text_response.status_code}. Сообщение опубликовано, но без ссылок.") except RetryAfterException: raise except Exception as e: # Проверяем, является ли это таймаутом is_timeout = ( isinstance(e, requests.exceptions.Timeout) or 'Timeout' in str(e) or 'timeout' in str(e).lower() or 'disconnected' in str(e).lower() ) if is_timeout: logger.warning(f"Таймаут при редактировании текста для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.") else: logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.") # Не пробрасываем исключение - сообщение уже опубликовано, просто без ссылок # Это считается частичным успехом await asyncio.sleep(2) # задержка между обращениями к телеграм # Публикация опроса, если требуется poll_message_id = None if post['is_poll'] and post['poll_question'] and post['poll_options']: try: # Парсим варианты ответа options = json.loads(post['poll_options']) # Отправляем опрос (в каналах можно отправлять только анонимные опросы) send_poll_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendPoll" poll_payload = { 'chat_id': CHANNEL_ID, 'question': post['poll_question'], 'options': options, 'is_anonymous': True, # В каналах можно отправлять только анонимные опросы 'allows_multiple_answers': post['poll_multiple'], 'disable_notification': PUBLISH_SILENTLY } # Добавляем close_date, если указан if post['poll_end_date']: # Преобразуем datetime в Unix timestamp if isinstance(post['poll_end_date'], datetime): poll_payload['close_date'] = int(post['poll_end_date'].timestamp()) else: poll_payload['close_date'] = post['poll_end_date'] try: poll_response = await httpx_client.post(send_poll_url, json=poll_payload) except httpx.TimeoutException as e: logger.error(f"ТАЙМАУТ при отправке опроса для записи VK ID {vk_post_id}") logger.error(f"Тип исключения: {type(e).__name__}") logger.error(f"Сообщение об ошибке: {str(e)}") logger.error(f"Полная информация об исключении: {repr(e)}") if hasattr(e, 'request'): logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise # Пробрасываем таймаут наверх if poll_response.status_code == 200: poll_result = poll_response.json() if poll_result.get('ok'): poll_message_data = poll_result.get('result', {}) poll_message_id = poll_message_data.get('message_id') logger.info(f"Опубликован опрос для записи VK ID {vk_post_id}") else: log_telegram_error(poll_response, f"[Отправка опроса для VK ID {vk_post_id}]") error_code = poll_result.get('error_code') if error_code == 429: retry_after = poll_result.get('parameters', {}).get('retry_after', 60) raise RetryAfterException(retry_after) raise Exception(f"API error: {poll_result.get('description', 'Unknown error')}") else: log_telegram_error(poll_response, f"[Отправка опроса для VK ID {vk_post_id}]") if poll_response.status_code == 429: retry_after = 60 try: error_data = poll_response.json() retry_after = error_data.get('parameters', {}).get('retry_after', 60) except: pass raise RetryAfterException(retry_after) raise Exception(f"HTTP {poll_response.status_code}: {poll_response.text}") except RetryAfterException: raise except httpx.TimeoutException as e: logger.error(f"ТАЙМАУТ при публикации опроса для записи VK ID {vk_post_id}") logger.error(f"Детали таймаута: {type(e).__name__}: {e}") raise # Пробрасываем таймаут наверх except Exception as e: logger.error(f"Ошибка при публикации опроса для записи VK ID {vk_post_id}: {e}") # Проверяем, что message_id был получен (публикация успешна) if message_id is None: logger.error(f"Не удалось получить message_id для записи VK ID {vk_post_id}. Публикация не завершена.") return False # Обновляем запись в базе данных только если публикация успешна publication_successful = True current_time_utc = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S") cursor.execute(""" UPDATE posts SET published_in_tg = True, tg_publication_date = %s, tg_message_id = %s, tg_poll_id = %s WHERE vk_post_id = %s """, (current_time_utc, message_id, poll_message_id, vk_post_id)) conn.commit() logger.info(f"Запись VK ID {vk_post_id} успешно опубликована") return True finally: await httpx_client.aclose() except RetryAfterException as e: logger.error(f"Получена ошибка FloodWait (429) при публикации записи VK ID {vk_post_id}: {e}") logger.error(f"Необходимо подождать {e.retry_after} секунд перед следующей попыткой") raise except httpx.TimeoutException as e: logger.error(f"ТАЙМАУТ при публикации записи VK ID {vk_post_id}") logger.error(f"Тип исключения: {type(e).__name__}") logger.error(f"Сообщение об ошибке: {str(e)}") logger.error(f"Полная информация об исключении: {repr(e)}") if hasattr(e, 'request'): logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}") # Пытаемся получить информацию о таймауте из исключения if hasattr(e, 'timeout'): logger.error(f"Настройки таймаута: {e.timeout}") # Проверяем, запущены ли мы из gunicorn import sys if 'gunicorn' in sys.modules: logger.warning("Обнаружен gunicorn - возможен конфликт с event loop при запуске асинхронного кода в отдельном потоке") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") # Убеждаемся, что БД не обновляется при таймауте return False except Exception as e: error_msg = str(e) if error_msg.startswith("RetryAfter:") or (hasattr(e, 'retry_after')): retry_after = getattr(e, 'retry_after', int(error_msg.split(":")[1]) if ":" in error_msg else 60) logger.error(f"Получена ошибка FloodWait (429) при публикации записи VK ID {vk_post_id}") logger.error(f"Необходимо подождать {retry_after} секунд перед следующей попыткой") raise RetryAfterException(retry_after) logger.error(f"Ошибка Telegram при публикации записи VK ID {vk_post_id}: {type(e).__name__}: {e}") logger.error(f"Детали ошибки: {repr(e)}") # Убеждаемся, что БД не обновляется при ошибке return False finally: if conn: conn.close() async def publish_to_tg_all(): """Публикация всех неопубликованных записей (максимум 10 за один вызов)""" logger.info("Запуск скрипта публикации всех записей") conn = None try: conn = pymysql.connect( host=MDB_HOST, user=MDB_USER, password=MDB_PW, database=MDBASE, charset='utf8mb4', cursorclass=pymysql.cursors.DictCursor ) with conn.cursor() as cursor: cursor.execute(""" SELECT vk_post_id FROM posts WHERE marked_for_publication = True AND published_in_tg = False LIMIT 5 """) posts = cursor.fetchall() post_count = 0 for post in posts: try: success = await publish_to_tg(post['vk_post_id']) if success: post_count += 1 # Задержка между публикациями разных записей if post_count < len(posts): await asyncio.sleep(2) except RetryAfterException as e: logger.error(f"Прерываем публикацию из-за ошибки FloodWait. Ожидание: {e.retry_after} секунд") break except Exception as e: logger.error(f"Ошибка при публикации записи VK ID {post['vk_post_id']}: {e}") # Продолжаем публикацию следующих записей, несмотря на ошибку logger.info(f"Опубликовано записей в этом запуске: {post_count}") except Exception as e: logger.error(f"Ошибка при публикации: {e}") finally: if conn: conn.close() logger.info("Завершение работы скрипта") if __name__ == "__main__": asyncio.run(publish_to_tg_all())