import os import logging import asyncio import pymysql import time import html import json import re import httpx 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 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 httpx_client = httpx.AsyncClient( http2=False, timeout=20.0, follow_redirects=True ) # Переменная для отслеживания успешности публикации message_id = None publication_successful = 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: photo_payload = { 'chat_id': CHANNEL_ID, 'photo': image_url, 'caption': text, 'parse_mode': 'HTML', 'disable_notification': PUBLISH_SILENTLY } photo_response = await httpx_client.post(send_photo_url, json=photo_payload) 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') 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__}: {e}") raise # Пробрасываем таймаут наверх except Exception as e: logger.warning(f"Не удалось отправить изображение для записи VK ID {vk_post_id}: {e}. Отправляем текстовое сообщение.") text_payload = { 'chat_id': CHANNEL_ID, 'text': text, 'parse_mode': 'HTML', 'disable_notification': PUBLISH_SILENTLY, 'disable_web_page_preview': True } text_response = await httpx_client.post(send_message_url, json=text_payload) 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: # Если это организационное сообщение или нет изображения, отправляем текстовое сообщение text_payload = { 'chat_id': CHANNEL_ID, 'text': text, 'parse_mode': 'HTML', 'disable_notification': PUBLISH_SILENTLY, 'disable_web_page_preview': True } text_response = await httpx_client.post(send_message_url, json=text_payload) 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" if has_image and not is_org_message: try: edit_caption_payload = { 'chat_id': CHANNEL_ID, 'message_id': message_id, 'caption': final_text, 'parse_mode': 'HTML' } edit_caption_response = await httpx_client.post(edit_caption_url, json=edit_caption_payload) 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__}: {e}") 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 } edit_text_response = await httpx_client.post(edit_text_url, json=edit_text_payload) 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 } edit_text_response = await httpx_client.post(edit_text_url, json=edit_text_payload) 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'] poll_response = await httpx_client.post(send_poll_url, json=poll_payload) 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__}: {e}") if hasattr(e, 'request'): logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}") # Убеждаемся, что БД не обновляется при таймауте 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())