import os import logging import asyncio import pymysql import time import html import json import re from datetime import datetime, timezone from telegram import Bot, PollOption from telegram.error import TelegramError, RetryAfter 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') # Настройка логирования 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) 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)) 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 bot = Bot(token=BOT_TOKEN) # Проверяем наличие организационного хештега 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 has_image and not is_org_message: try: message = await bot.send_photo( chat_id=CHANNEL_ID, photo=image_url, caption=text, parse_mode="HTML", disable_notification=PUBLISH_SILENTLY ) except TelegramError as e: logger.warning(f"Не удалось отправить изображение для записи VK ID {vk_post_id}: {e}. Отправляем текстовое сообщение.") message = await bot.send_message( chat_id=CHANNEL_ID, text=text, parse_mode="HTML", disable_notification=PUBLISH_SILENTLY, disable_web_page_preview=True ) else: # Если это организационное сообщение или нет изображения, отправляем текстовое сообщение message = await bot.send_message( chat_id=CHANNEL_ID, text=text, parse_mode="HTML", disable_notification=PUBLISH_SILENTLY, disable_web_page_preview=True ) # Получаем ID сообщения для использования в ссылке подписки message_id = message.message_id # Формируем финальный текст с ссылками 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) # задержка между обращениями к телеграм # Редактируем сообщение, добавляя ссылки if has_image and not is_org_message: try: await bot.edit_message_caption( chat_id=CHANNEL_ID, message_id=message_id, caption=final_text, parse_mode="HTML" ) except TelegramError as e: logger.warning(f"Не удалось отредактировать подпись для записи VK ID {vk_post_id}: {e}. Пробуем отредактировать текстовое сообщение.") await bot.edit_message_text( chat_id=CHANNEL_ID, message_id=message_id, text=final_text, parse_mode="HTML", disable_web_page_preview=True ) else: await bot.edit_message_text( chat_id=CHANNEL_ID, message_id=message_id, text=final_text, parse_mode="HTML", disable_web_page_preview=True ) 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']) poll_options = [PollOption(option, 0) for option in options] # Отправляем опрос (в каналах можно отправлять только анонимные опросы) poll_message = await bot.send_poll( chat_id=CHANNEL_ID, question=post['poll_question'], options=poll_options, is_anonymous=True, # В каналах можно отправлять только анонимные опросы allows_multiple_answers=post['poll_multiple'], close_date=post['poll_end_date'] if post['poll_end_date'] else None, disable_notification=PUBLISH_SILENTLY ) poll_message_id = poll_message.message_id logger.info(f"Опубликован опрос для записи VK ID {vk_post_id}") except Exception as e: logger.error(f"Ошибка при публикации опроса для записи VK ID {vk_post_id}: {e}") # Обновляем запись в базе данных 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 except RetryAfter as e: logger.error(f"Получена ошибка FloodWait (429) при публикации записи VK ID {vk_post_id}: {e}") logger.error(f"Необходимо подождать {e.retry_after} секунд перед следующей попыткой") raise except TelegramError as e: logger.error(f"Ошибка Telegram при публикации записи VK ID {vk_post_id}: {e}") return False except Exception as e: logger.error(f"Неожиданная ошибка при публикации записи VK ID {vk_post_id}: {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 RetryAfter 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())