import os import logging import asyncio import pymysql import time import html import json import re import httpx import traceback import io from datetime import datetime, timezone from dotenv import load_dotenv from formatter import get_event_text # Загрузка переменных окружения 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""" try: if response is not None: logger.error(f"{context} Статус код: {response.status_code}") logger.error(f"{context} Заголовки ответа: {dict(response.headers)}") try: error_data = response.json() logger.error(f"{context} Полный ответ от Telegram API: {json.dumps(error_data, ensure_ascii=False, indent=2)}") except: 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 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 timeout_config = httpx.Timeout( connect=10.0, # Таймаут на подключение read=60.0, # Таймаут на чтение (увеличен для больших изображений) write=30.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)} символов") # Скачиваем изображение 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 # Логируем время начала запроса request_start_time = time.time() logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}") # Отправляем изображение как файл через multipart/form-data # Создаем BytesIO объект для файла image_file = io.BytesIO(image_data) files = { 'photo': ('image.jpg', image_file, 'image/jpeg') } data = { 'chat_id': CHANNEL_ID, 'caption': text, 'parse_mode': 'HTML', 'disable_notification': str(PUBLISH_SILENTLY).lower() } photo_response = await httpx_client.post(send_photo_url, files=files, data=data) request_duration = time.time() - request_start_time logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд") if photo_response.status_code == 200: photo_result = photo_response.json() if photo_result.get('ok'): message_data = photo_result.get('result', {}) message_id = message_data.get('message_id') message_has_image = True # Успешно отправили изображение else: # Ошибка в ответе API - логируем полную информацию log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]") 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: log_telegram_error(photo_response, f"[Отправка фото для VK ID {vk_post_id}]") 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}") except RetryAfterException: 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'}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") 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 try: edit_caption_payload = { 'chat_id': CHANNEL_ID, 'message_id': message_id, 'caption': final_text, 'parse_mode': 'HTML' } try: edit_caption_response = await httpx_client.post(edit_caption_url, json=edit_caption_payload) except httpx.TimeoutException as e: logger.error(f"ТАЙМАУТ при редактировании подписи для записи VK ID {vk_post_id}") logger.error(f"Детали таймаута: {type(e).__name__}: {e}") if hasattr(e, 'request'): logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}") raise # Пробрасываем таймаут наверх if edit_caption_response.status_code == 200: edit_result = edit_caption_response.json() if not edit_result.get('ok'): 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) raise Exception(f"API error: {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) raise Exception(f"HTTP {edit_caption_response.status_code}: {edit_caption_response.text}") except RetryAfterException: 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'}") logger.error(f"Трассировка стека:\n{traceback.format_exc()}") raise # Пробрасываем таймаут наверх except Exception as e: logger.warning(f"Не удалось отредактировать подпись для записи VK ID {vk_post_id}: {e}. Пробуем отредактировать текстовое сообщение.") edit_text_payload = { 'chat_id': CHANNEL_ID, 'message_id': message_id, 'text': final_text, 'parse_mode': 'HTML', 'disable_web_page_preview': True } try: edit_text_response = await httpx_client.post(edit_text_url, json=edit_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 edit_text_response.status_code == 200: edit_result = edit_text_response.json() if not edit_result.get('ok'): log_telegram_error(edit_text_response, f"[Редактирование текста (fallback) для 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) raise Exception(f"API error: {edit_result.get('description', 'Unknown error')}") else: log_telegram_error(edit_text_response, f"[Редактирование текста (fallback) для 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) raise Exception(f"HTTP {edit_text_response.status_code}: {edit_text_response.text}") else: edit_text_payload = { 'chat_id': CHANNEL_ID, 'message_id': message_id, 'text': final_text, 'parse_mode': 'HTML', 'disable_web_page_preview': True } try: edit_text_response = await httpx_client.post(edit_text_url, json=edit_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 edit_text_response.status_code == 200: edit_result = edit_text_response.json() if not edit_result.get('ok'): 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) raise Exception(f"API error: {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) raise Exception(f"HTTP {edit_text_response.status_code}: {edit_text_response.text}") 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())