From a8f61f724faca3fa801c030e60b70ebf589d78d4 Mon Sep 17 00:00:00 2001 From: gitadmin Date: Sat, 10 Jan 2026 22:11:15 +0300 Subject: [PATCH] =?UTF-8?q?[=D0=92=D0=9E=D0=9B=D0=9A]=20retrying=20logic?= =?UTF-8?q?=20=D0=B8=20=D0=B4=D1=80=D1=83=D0=B3=D0=B0=D1=8F=20=D0=B1=D0=B8?= =?UTF-8?q?=D0=B1=D0=BB=D0=B8=D0=BE=D1=82=D0=B5=D0=BA=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tg_publish.py | 977 ++++++++++++-------------------------------------- 1 file changed, 222 insertions(+), 755 deletions(-) diff --git a/tg_publish.py b/tg_publish.py index 4408bd6..5893e67 100644 --- a/tg_publish.py +++ b/tg_publish.py @@ -7,16 +7,15 @@ 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 +from telegram import Bot +from telegram.error import TelegramError, TimedOut +from telegram.request import HTTPXRequest # Настройка anyio для правильной работы в отдельном потоке # Устанавливаем правильный backend для anyio (используем asyncio) @@ -145,372 +144,6 @@ async def download_image(image_url, httpx_client): 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 - ) - - pool = super().init_poolmanager(*args, **kwargs) - - # Устанавливаем socket timeout для всех соединений в пуле - # Это критично для операций записи в сокет - # Используем несколько подходов для надежности - try: - if hasattr(pool, 'ConnectionCls') and pool.ConnectionCls: - original_connect = pool.ConnectionCls.connect - def connect_with_socket_timeout(self): - result = original_connect(self) - # Устанавливаем socket timeout после подключения - # Это применяется ко всем операциям с сокетом (read и write) - if hasattr(self, 'sock') and self.sock: - self.sock.settimeout(write_timeout) - return result - pool.ConnectionCls.connect = connect_with_socket_timeout - except Exception: - # Если не удалось установить через ConnectionCls, используем другой подход - pass - - return pool - - def _new_conn(self): - """Создаем новое соединение с установкой socket timeout""" - conn = super()._new_conn() - # Устанавливаем socket timeout для операций записи - try: - if hasattr(conn, 'sock') and conn.sock: - conn.sock.settimeout(write_timeout) - elif hasattr(conn, 'connect'): - # Если соединение еще не создано, перехватываем метод connect - original_connect = conn.connect - def connect_with_timeout(): - result = original_connect() - if hasattr(conn, 'sock') and conn.sock: - conn.sock.settimeout(write_timeout) - return result - conn.connect = connect_with_timeout - except Exception: - # Игнорируем ошибки при установке timeout - pass - return conn - - # Монтируем кастомный адаптер - 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 объект - """ - import socket - - try: - data = { - 'chat_id': chat_id, - 'message_id': message_id, - 'caption': caption, - 'parse_mode': 'HTML' - } - - # Извлекаем таймауты - 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 - else: - connect_timeout = read_timeout = write_timeout = timeout[0] - else: - connect_timeout = read_timeout = write_timeout = timeout - - # Преобразуем timeout в формат, который принимает requests - requests_timeout = (connect_timeout, max(read_timeout, write_timeout)) - - # Создаем сессию с кастомным адаптером для установки socket timeout - session = requests.Session() - - class CustomHTTPAdapter(HTTPAdapter): - def init_poolmanager(self, *args, **kwargs): - socket_options = kwargs.get('socket_options', []) - socket_options.append((socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)) - kwargs['socket_options'] = socket_options - - max_timeout = max(connect_timeout, read_timeout, write_timeout) - if 'timeout' not in kwargs: - kwargs['timeout'] = Urllib3Timeout( - connect=connect_timeout, - read=max_timeout - ) - - pool = super().init_poolmanager(*args, **kwargs) - - # Устанавливаем socket timeout для всех соединений в пуле - try: - if hasattr(pool, 'ConnectionCls') and pool.ConnectionCls: - original_connect = pool.ConnectionCls.connect - def connect_with_socket_timeout(self): - result = original_connect(self) - if hasattr(self, 'sock') and self.sock: - self.sock.settimeout(write_timeout) - return result - pool.ConnectionCls.connect = connect_with_socket_timeout - except Exception: - pass - - return pool - - def _new_conn(self): - """Создаем новое соединение с установкой socket timeout""" - conn = super()._new_conn() - try: - if hasattr(conn, 'sock') and conn.sock: - conn.sock.settimeout(write_timeout) - elif hasattr(conn, 'connect'): - original_connect = conn.connect - def connect_with_timeout(): - result = original_connect() - if hasattr(conn, 'sock') and conn.sock: - conn.sock.settimeout(write_timeout) - return result - conn.connect = connect_with_timeout - except Exception: - pass - return conn - - adapter = CustomHTTPAdapter(max_retries=Retry(total=0)) - session.mount('http://', adapter) - session.mount('https://', adapter) - - response = session.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 - finally: - if 'session' in locals(): - session.close() - -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 объект - """ - import socket - - try: - data = { - 'chat_id': chat_id, - 'message_id': message_id, - 'text': text, - 'parse_mode': 'HTML', - 'disable_web_page_preview': True - } - - # Извлекаем таймауты - 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 - else: - connect_timeout = read_timeout = write_timeout = timeout[0] - else: - connect_timeout = read_timeout = write_timeout = timeout - - # Преобразуем timeout в формат, который принимает requests - requests_timeout = (connect_timeout, max(read_timeout, write_timeout)) - - # Создаем сессию с кастомным адаптером для установки socket timeout - session = requests.Session() - - class CustomHTTPAdapter(HTTPAdapter): - def init_poolmanager(self, *args, **kwargs): - socket_options = kwargs.get('socket_options', []) - socket_options.append((socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)) - kwargs['socket_options'] = socket_options - - max_timeout = max(connect_timeout, read_timeout, write_timeout) - if 'timeout' not in kwargs: - kwargs['timeout'] = Urllib3Timeout( - connect=connect_timeout, - read=max_timeout - ) - - pool = super().init_poolmanager(*args, **kwargs) - - # Устанавливаем socket timeout для всех соединений в пуле - try: - if hasattr(pool, 'ConnectionCls') and pool.ConnectionCls: - original_connect = pool.ConnectionCls.connect - def connect_with_socket_timeout(self): - result = original_connect(self) - if hasattr(self, 'sock') and self.sock: - self.sock.settimeout(write_timeout) - return result - pool.ConnectionCls.connect = connect_with_socket_timeout - except Exception: - pass - - return pool - - def _new_conn(self): - """Создаем новое соединение с установкой socket timeout""" - conn = super()._new_conn() - try: - if hasattr(conn, 'sock') and conn.sock: - conn.sock.settimeout(write_timeout) - elif hasattr(conn, 'connect'): - original_connect = conn.connect - def connect_with_timeout(): - result = original_connect() - if hasattr(conn, 'sock') and conn.sock: - conn.sock.settimeout(write_timeout) - return result - conn.connect = connect_with_timeout - except Exception: - pass - return conn - - adapter = CustomHTTPAdapter(max_retries=Retry(total=0)) - session.mount('http://', adapter) - session.mount('https://', adapter) - - response = session.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 - finally: - if 'session' in locals(): - session.close() - async def publish_to_tg(vk_post_id): """Публикация одной записи по VK post ID""" logger.info(f"Запуск публикации для записи VK ID {vk_post_id}") @@ -540,9 +173,19 @@ async def publish_to_tg(vk_post_id): logger.info(f"Запись VK ID {vk_post_id} не найдена или уже опубликована") return False - # Создаем httpx клиент с отключенным HTTP/2 - # Используем более детальные таймауты: connect, read, write, pool - # Таймаут на запись увеличен для больших файлов (до 10 МБ) + # Настройка таймаутов для HTTPXRequest (аналогично test_tg_poster.py) + request = HTTPXRequest( + read_timeout=120.0, # Таймаут на чтение (увеличен для больших изображений) + write_timeout=120.0, # Таймаут на запись (увеличен для больших файлов) + connect_timeout=30.0, # Таймаут на подключение (увеличен) + pool_timeout=10.0, # Таймаут на получение соединения из пула + media_write_timeout=120.0 # Таймаут для медиа + ) + + # Создание экземпляра бота с кастомным Request + bot = Bot(token=BOT_TOKEN, request=request) + + # Создаем httpx клиент только для скачивания изображений timeout_config = httpx.Timeout( connect=10.0, # Таймаут на подключение read=120.0, # Таймаут на чтение (увеличен для больших изображений) @@ -555,6 +198,10 @@ async def publish_to_tg(vk_post_id): follow_redirects=True ) + # Параметры для повторных попыток при таймаутах + max_retries = 10 + retry_interval = 5 + # Переменная для отслеживания успешности публикации message_id = None publication_successful = False @@ -610,9 +257,6 @@ async def publish_to_tg(vk_post_id): 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: # Логируем информацию об изображении перед отправкой @@ -629,224 +273,135 @@ async def publish_to_tg(vk_post_id): 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} секунд") + logger.info(f"Размер файла: {file_size_mb:.2f} МБ") - # Используем 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}") + # Повторные попытки отправки фото при таймауте (до 10 раз) + message = None + for attempt in range(1, max_retries + 1): + try: + image_file = io.BytesIO(image_data) + message = await bot.send_photo( + chat_id=CHANNEL_ID, + photo=image_file, + caption=text, + parse_mode='HTML', + disable_notification=PUBLISH_SILENTLY + ) + message_id = message.message_id + message_has_image = True + logger.info(f"Сообщение с изображением успешно опубликовано. Message ID: {message_id}") + break # Успешная публикация, выходим из цикла + except TimedOut as e: + if attempt < max_retries: + logger.warning(f"Таймаут при публикации изображения (попытка {attempt}/{max_retries}) для записи VK ID {vk_post_id}. Повтор через {retry_interval} секунд...") + logger.warning(f"Детали таймаута: {type(e).__name__}: {e}") + await asyncio.sleep(retry_interval) + else: + logger.error(f"Таймаут при публикации изображения после {max_retries} попыток для записи VK ID {vk_post_id}") + logger.error(f"Детали ошибки: {type(e).__name__}: {e}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise # Пробрасываем таймаут наверх + except TelegramError as e: + # Проверяем, является ли это ошибкой FloodWait (429) + if hasattr(e, 'retry_after') or '429' in str(e): + retry_after = getattr(e, 'retry_after', 60) + raise RetryAfterException(retry_after) + logger.error(f"Ошибка Telegram при публикации изображения для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + logger.error(f"Детали ошибки: {repr(e)}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") - raise # Пробрасываем таймаут наверх - else: - # Другие ошибки - logger.error(f"Ошибка при отправке изображения (requests) для записи VK ID {vk_post_id}: {error_msg}") + raise + except Exception as e: + logger.error(f"Неожиданная ошибка при публикации изображения для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + logger.error(f"Детали ошибки: {repr(e)}") 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 message_id is None: + raise Exception("Не удалось опубликовать изображение после всех попыток") - 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) + # Повторные попытки отправки текста (fallback) при таймауте + for attempt in range(1, max_retries + 1): + try: + message = await bot.send_message( + chat_id=CHANNEL_ID, + text=text, + parse_mode='HTML', + disable_notification=PUBLISH_SILENTLY, + disable_web_page_preview=True + ) + message_id = message.message_id + logger.info(f"Текстовое сообщение (fallback) успешно опубликовано. Message ID: {message_id}") + break # Успешная публикация, выходим из цикла + except TimedOut as e: + if attempt < max_retries: + logger.warning(f"Таймаут при отправке текста (fallback) (попытка {attempt}/{max_retries}) для записи VK ID {vk_post_id}. Повтор через {retry_interval} секунд...") + await asyncio.sleep(retry_interval) + else: + logger.error(f"Таймаут при отправке текста (fallback) после {max_retries} попыток для записи VK ID {vk_post_id}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise + except TelegramError as e: + if hasattr(e, 'retry_after') or '429' in str(e): + retry_after = getattr(e, '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}") + logger.error(f"Ошибка Telegram при отправке текста (fallback) для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + raise + except Exception as e: + logger.error(f"Неожиданная ошибка при отправке текста (fallback) для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + raise 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) + # Повторные попытки отправки текста при таймауте + message = None + for attempt in range(1, max_retries + 1): + try: + message = await bot.send_message( + chat_id=CHANNEL_ID, + text=text, + parse_mode='HTML', + disable_notification=PUBLISH_SILENTLY, + disable_web_page_preview=True + ) + message_id = message.message_id + logger.info(f"Текстовое сообщение успешно опубликовано. Message ID: {message_id}") + break # Успешная публикация, выходим из цикла + except TimedOut as e: + if attempt < max_retries: + logger.warning(f"Таймаут при публикации текста (попытка {attempt}/{max_retries}) для записи VK ID {vk_post_id}. Повтор через {retry_interval} секунд...") + logger.warning(f"Детали таймаута: {type(e).__name__}: {e}") + await asyncio.sleep(retry_interval) + else: + logger.error(f"Таймаут при публикации текста после {max_retries} попыток для записи VK ID {vk_post_id}") + logger.error(f"Детали ошибки: {type(e).__name__}: {e}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise # Пробрасываем таймаут наверх + except TelegramError as e: + # Проверяем, является ли это ошибкой FloodWait (429) + if hasattr(e, 'retry_after') or '429' in str(e): + retry_after = getattr(e, '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}") + logger.error(f"Ошибка Telegram при публикации текста для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + logger.error(f"Детали ошибки: {repr(e)}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise + except Exception as e: + logger.error(f"Неожиданная ошибка при публикации текста для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + logger.error(f"Детали ошибки: {repr(e)}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise + + if message_id is None: + raise Exception("Не удалось опубликовать текстовое сообщение после всех попыток") # Формируем финальный текст с ссылками if links: @@ -864,145 +419,75 @@ async def publish_to_tg(vk_post_id): 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}") + # Сообщение с изображением - редактируем caption с повторными попытками + edit_success = False + for attempt in range(1, max_retries + 1): + try: + await bot.edit_message_caption( + chat_id=CHANNEL_ID, + message_id=message_id, + caption=final_text, + parse_mode='HTML' + ) + edit_success = True + logger.info(f"Caption успешно отредактирован для записи VK ID {vk_post_id}") + break # Успешное редактирование, выходим из цикла + except TimedOut as e: + if attempt < max_retries: + logger.warning(f"Таймаут при редактировании caption (попытка {attempt}/{max_retries}) для записи VK ID {vk_post_id}. Повтор через {retry_interval} секунд...") + await asyncio.sleep(retry_interval) 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 после {max_retries} попыток для записи VK ID {vk_post_id}. Сообщение опубликовано, но без ссылок.") # Не критичная ошибка - сообщение уже опубликовано - 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 + break + except TelegramError as e: + if hasattr(e, 'retry_after') or '429' in str(e): + retry_after = getattr(e, 'retry_after', 60) 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}. Сообщение опубликовано, но без ссылок.") - - # Не пробрасываем исключение - сообщение уже опубликовано, просто без ссылок - # Это считается частичным успехом + logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: {type(e).__name__}: {e}. Сообщение опубликовано, но без ссылок.") + break + except Exception as e: + # Не критичная ошибка - сообщение уже опубликовано + logger.warning(f"Неожиданная ошибка при редактировании caption для записи VK ID {vk_post_id}: {type(e).__name__}: {e}. Сообщение опубликовано, но без ссылок.") + break 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}") + # Текстовое сообщение - редактируем текст с повторными попытками + edit_success = False + for attempt in range(1, max_retries + 1): + try: + await bot.edit_message_text( + chat_id=CHANNEL_ID, + message_id=message_id, + text=final_text, + parse_mode='HTML', + disable_web_page_preview=True + ) + edit_success = True + logger.info(f"Текст успешно отредактирован для записи VK ID {vk_post_id}") + break # Успешное редактирование, выходим из цикла + except TimedOut as e: + if attempt < max_retries: + logger.warning(f"Таймаут при редактировании текста (попытка {attempt}/{max_retries}) для записи VK ID {vk_post_id}. Повтор через {retry_interval} секунд...") + await asyncio.sleep(retry_interval) 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"Таймаут при редактировании текста после {max_retries} попыток для записи VK ID {vk_post_id}. Сообщение опубликовано, но без ссылок.") # Не критичная ошибка - сообщение уже опубликовано - 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 + break + except TelegramError as e: + if hasattr(e, 'retry_after') or '429' in str(e): + retry_after = getattr(e, 'retry_after', 60) 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}. Сообщение опубликовано, но без ссылок.") - - # Не пробрасываем исключение - сообщение уже опубликовано, просто без ссылок - # Это считается частичным успехом + logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: {type(e).__name__}: {e}. Сообщение опубликовано, но без ссылок.") + break + except Exception as e: + # Не критичная ошибка - сообщение уже опубликовано + logger.warning(f"Неожиданная ошибка при редактировании текста для записи VK ID {vk_post_id}: {type(e).__name__}: {e}. Сообщение опубликовано, но без ссылок.") + break await asyncio.sleep(2) # задержка между обращениями к телеграм @@ -1013,9 +498,8 @@ async def publish_to_tg(vk_post_id): # Парсим варианты ответа options = json.loads(post['poll_options']) - # Отправляем опрос (в каналах можно отправлять только анонимные опросы) - send_poll_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendPoll" - poll_payload = { + # Подготавливаем параметры для опроса + poll_kwargs = { 'chat_id': CHANNEL_ID, 'question': post['poll_question'], 'options': options, @@ -1028,53 +512,45 @@ async def publish_to_tg(vk_post_id): 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()) + poll_kwargs['close_date'] = int(post['poll_end_date'].timestamp()) else: - poll_payload['close_date'] = post['poll_end_date'] + poll_kwargs['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') + # Повторные попытки отправки опроса при таймауте + poll_message = None + for attempt in range(1, max_retries + 1): + try: + poll_message = await bot.send_poll(**poll_kwargs) + poll_message_id = poll_message.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) + break # Успешная публикация, выходим из цикла + except TimedOut as e: + if attempt < max_retries: + logger.warning(f"Таймаут при отправке опроса (попытка {attempt}/{max_retries}) для записи VK ID {vk_post_id}. Повтор через {retry_interval} секунд...") + logger.warning(f"Детали таймаута: {type(e).__name__}: {e}") + await asyncio.sleep(retry_interval) + else: + logger.error(f"Таймаут при отправке опроса после {max_retries} попыток для записи VK ID {vk_post_id}") + logger.error(f"Детали ошибки: {type(e).__name__}: {e}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise # Пробрасываем таймаут наверх + except TelegramError as e: + # Проверяем, является ли это ошибкой FloodWait (429) + if hasattr(e, 'retry_after') or '429' in str(e): + retry_after = getattr(e, '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}") + logger.error(f"Ошибка Telegram при публикации опроса для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + logger.error(f"Детали ошибки: {repr(e)}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise + except Exception as e: + logger.error(f"Неожиданная ошибка при публикации опроса для записи VK ID {vk_post_id}: {type(e).__name__}: {e}") + logger.error(f"Детали ошибки: {repr(e)}") + logger.error(f"Трассировка стека:\n{traceback.format_exc()}") + raise 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}") @@ -1105,20 +581,11 @@ async def publish_to_tg(vk_post_id): logger.error(f"Получена ошибка FloodWait (429) при публикации записи VK ID {vk_post_id}: {e}") logger.error(f"Необходимо подождать {e.retry_after} секунд перед следующей попыткой") raise - except httpx.TimeoutException as e: + except TimedOut 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