660 lines
42 KiB
Python
660 lines
42 KiB
Python
import os
|
||
import logging
|
||
import asyncio
|
||
import pymysql
|
||
import time
|
||
import html
|
||
import json
|
||
import re
|
||
import httpx
|
||
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)
|
||
try:
|
||
import anyio
|
||
# Устанавливаем backend для anyio явно
|
||
if not hasattr(anyio, '_backend'):
|
||
# Используем asyncio backend
|
||
os.environ.setdefault('ANYIO_BACKEND', 'asyncio')
|
||
except ImportError:
|
||
pass
|
||
|
||
# Загрузка переменных окружения
|
||
load_dotenv()
|
||
|
||
# Настройки из переменных окружения
|
||
BOT_TOKEN = os.getenv('POSTER_BOT_TOKEN')
|
||
RESPONDER_BOT_NAME = os.getenv('RESPONDER_BOT_NAME')
|
||
CHANNEL_ID = os.getenv('CHANNEL_ID')
|
||
MDB_HOST = os.getenv('MDB_HOST')
|
||
MDB_USER = os.getenv('MDB_USER')
|
||
MDB_PW = os.getenv('MDB_PW')
|
||
MDBASE = os.getenv('MDBASE')
|
||
LOG_FILE = os.getenv('LOG_FILE', 'post_publisher.log')
|
||
ORG_MESSAGE_TAG = os.getenv('ORG_MESSAGE_TAG', '') # Хештег для определения организационных сообщений
|
||
|
||
# Параметры длины сообщений
|
||
MAX_CAPTION_LENGTH = int(os.getenv('MAX_CAPTION_LENGTH', 1000))
|
||
MAX_TEXT_LENGTH = int(os.getenv('MAX_TEXT_LENGTH', 4000))
|
||
|
||
# Режим публикации без звука
|
||
PUBLISH_SILENTLY = os.getenv('PUBLISH_SILENTLY', 'false').lower() in ('true', '1', 'yes', 'on')
|
||
|
||
# Использование бота подписки
|
||
USE_SUBSCRIPTION_BOT = os.getenv('USE_SUBSCRIPTION_BOT', 'true').lower() in ('true', '1', 'yes', 'on')
|
||
|
||
# Режим отладочного логирования данных
|
||
LOG_DEBUG_DATA = os.getenv('LOG_DEBUG_DATA', 'false').lower() in ('true', '1', 'yes', 'on')
|
||
|
||
# Настройка логирования
|
||
logger = logging.getLogger('TG_Poster')
|
||
logger.setLevel(logging.INFO)
|
||
|
||
formatter = logging.Formatter('[%(asctime)s] [%(name)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')
|
||
|
||
console_handler = logging.StreamHandler()
|
||
console_handler.setFormatter(formatter)
|
||
logger.addHandler(console_handler)
|
||
|
||
class RetryFileHandler(logging.FileHandler):
|
||
def emit(self, record):
|
||
for _ in range(5):
|
||
try:
|
||
super().emit(record)
|
||
return
|
||
except (IOError, PermissionError):
|
||
time.sleep(0.5)
|
||
print(f"Failed to write to log file after 5 attempts: {record.msg}")
|
||
|
||
file_handler = RetryFileHandler(LOG_FILE, encoding='utf-8')
|
||
file_handler.setFormatter(formatter)
|
||
logger.addHandler(file_handler)
|
||
|
||
# Класс для обработки RetryAfter (FloodWait 429)
|
||
class RetryAfterException(Exception):
|
||
"""Исключение для обработки FloodWait (429) от Telegram API"""
|
||
def __init__(self, retry_after):
|
||
self.retry_after = retry_after
|
||
super().__init__(f"RetryAfter: {retry_after}")
|
||
|
||
def prepare_text(text):
|
||
"""Подготовка текстовых полей с экранированием HTML-сущностей"""
|
||
if not text:
|
||
return ""
|
||
|
||
# Экранируем специальные символы HTML
|
||
text = html.escape(str(text))
|
||
|
||
return text
|
||
|
||
def has_org_tag(text):
|
||
"""Проверяет, содержит ли текст организационный хештег"""
|
||
if not text or not ORG_MESSAGE_TAG:
|
||
return False
|
||
|
||
# Ищем хештег (с учетом возможных пробелов и других символов)
|
||
pattern = r'#?' + re.escape(ORG_MESSAGE_TAG) + r'\b'
|
||
return bool(re.search(pattern, text, re.IGNORECASE))
|
||
|
||
def log_telegram_error(response, context=""):
|
||
"""Логирует полную информацию об ошибке от Telegram API
|
||
Работает как с httpx.Response, так и с requests.Response
|
||
"""
|
||
try:
|
||
if response is not None:
|
||
logger.error(f"{context} Статус код: {response.status_code}")
|
||
logger.error(f"{context} Заголовки ответа: {dict(response.headers)}")
|
||
try:
|
||
# Работает для обоих типов ответов (httpx и requests)
|
||
error_data = response.json()
|
||
logger.error(f"{context} Полный ответ от Telegram API: {json.dumps(error_data, ensure_ascii=False, indent=2)}")
|
||
except:
|
||
# Для requests используем text, для httpx тоже text
|
||
response_text = getattr(response, 'text', str(response.content[:500]) if hasattr(response, 'content') else 'N/A')
|
||
logger.error(f"{context} Текст ответа: {response_text}")
|
||
else:
|
||
logger.error(f"{context} Ответ от Telegram API отсутствует (None)")
|
||
except Exception as e:
|
||
logger.error(f"{context} Ошибка при логировании информации об ошибке: {e}")
|
||
|
||
async def download_image(image_url, httpx_client):
|
||
"""Скачивает изображение по URL и возвращает его содержимое в виде bytes"""
|
||
try:
|
||
logger.info(f"Скачивание изображения: {image_url}")
|
||
download_start_time = time.time()
|
||
|
||
response = await httpx_client.get(image_url, timeout=httpx.Timeout(connect=10.0, read=60.0, write=30.0, pool=10.0))
|
||
response.raise_for_status()
|
||
|
||
download_duration = time.time() - download_start_time
|
||
image_data = response.content
|
||
logger.info(f"Изображение скачано за {download_duration:.2f} секунд, размер: {len(image_data)} байт")
|
||
|
||
return image_data
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при скачивании изображения {image_url}: {e}")
|
||
raise
|
||
|
||
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
|
||
|
||
# Настройка таймаутов для 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, # Таймаут на чтение (увеличен для больших изображений)
|
||
write=120.0, # Таймаут на запись (увеличен для больших файлов)
|
||
pool=10.0 # Таймаут на получение соединения из пула
|
||
)
|
||
httpx_client = httpx.AsyncClient(
|
||
http2=False,
|
||
timeout=timeout_config,
|
||
follow_redirects=True
|
||
)
|
||
|
||
# Параметры для повторных попыток при таймаутах
|
||
max_retries = 10
|
||
retry_interval = 5
|
||
|
||
# Переменная для отслеживания успешности публикации
|
||
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'<a href="{post["vk_post_url"]}">Оригинал в ВК</a>')
|
||
|
||
# Ссылка на подписку (только для событий и если включен бот подписки)
|
||
if post['is_event'] and USE_SUBSCRIPTION_BOT:
|
||
links.append(f'<a href="https://t.me/{RESPONDER_BOT_NAME}?start=post_{placeholder_message_id}">🔜 Подписка</a>')
|
||
|
||
# Формируем строку ссылок
|
||
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"=== Конец отформатированного текста ===")
|
||
|
||
# Публикация первоначального сообщения без ссылок
|
||
if has_image and not is_org_message:
|
||
try:
|
||
# Логируем информацию об изображении перед отправкой
|
||
logger.info(f"Отправка изображения для записи VK ID {vk_post_id}")
|
||
logger.info(f"URL изображения: {image_url}")
|
||
logger.info(f"Длина URL изображения: {len(image_url)} символов")
|
||
logger.info(f"Длина подписи: {len(text)} символов")
|
||
|
||
# Скачиваем изображение
|
||
image_data = None
|
||
try:
|
||
image_data = await download_image(image_url, httpx_client)
|
||
except Exception as download_error:
|
||
logger.error(f"Ошибка при скачивании изображения для записи VK ID {vk_post_id}: {download_error}")
|
||
raise # Пробрасываем ошибку, чтобы перейти к fallback
|
||
|
||
file_size_mb = len(image_data) / (1024 * 1024)
|
||
logger.info(f"Размер файла: {file_size_mb:.2f} МБ")
|
||
|
||
# Повторные попытки отправки фото при таймауте (до 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
|
||
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("Не удалось опубликовать изображение после всех попыток")
|
||
|
||
except RetryAfterException:
|
||
raise
|
||
except Exception as e:
|
||
# Если не удалось отправить изображение, пробуем отправить текстовое сообщение
|
||
logger.warning(f"Не удалось отправить изображение для записи VK ID {vk_post_id}: {e}. Отправляем текстовое сообщение.")
|
||
message_has_image = False # Отправляем текстовое сообщение вместо изображения
|
||
|
||
# Повторные попытки отправки текста (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)
|
||
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 # Отправляем текстовое сообщение
|
||
|
||
# Повторные попытки отправки текста при таймауте
|
||
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)
|
||
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:
|
||
# Обновляем ссылку подписки с реальным 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'<a href="{post["vk_post_url"]}">Оригинал в ВК</a>')
|
||
|
||
if post['is_event'] and USE_SUBSCRIPTION_BOT:
|
||
updated_links.append(f'<a href="https://t.me/{RESPONDER_BOT_NAME}?start=post_{message_id}">🔜 Подписка</a>')
|
||
|
||
links_line = " | ".join(updated_links)
|
||
final_text = f"{text}\n\n{links_line}"
|
||
|
||
await asyncio.sleep(2) # задержка между обращениями к телеграм
|
||
|
||
# Редактируем сообщение, добавляя ссылки
|
||
# Используем message_has_image для определения типа сообщения
|
||
# Если было отправлено изображение, редактируем caption, иначе - текст
|
||
if message_has_image:
|
||
# Сообщение с изображением - редактируем 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:
|
||
logger.warning(f"Таймаут при редактировании caption после {max_retries} попыток для записи VK ID {vk_post_id}. Сообщение опубликовано, но без ссылок.")
|
||
# Не критичная ошибка - сообщение уже опубликовано
|
||
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}: {type(e).__name__}: {e}. Сообщение опубликовано, но без ссылок.")
|
||
break
|
||
except Exception as e:
|
||
# Не критичная ошибка - сообщение уже опубликовано
|
||
logger.warning(f"Неожиданная ошибка при редактировании caption для записи VK ID {vk_post_id}: {type(e).__name__}: {e}. Сообщение опубликовано, но без ссылок.")
|
||
break
|
||
else:
|
||
# Текстовое сообщение - редактируем текст с повторными попытками
|
||
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:
|
||
logger.warning(f"Таймаут при редактировании текста после {max_retries} попыток для записи VK ID {vk_post_id}. Сообщение опубликовано, но без ссылок.")
|
||
# Не критичная ошибка - сообщение уже опубликовано
|
||
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}: {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) # задержка между обращениями к телеграм
|
||
|
||
# Публикация опроса, если требуется
|
||
poll_message_id = None
|
||
if post['is_poll'] and post['poll_question'] and post['poll_options']:
|
||
try:
|
||
# Парсим варианты ответа
|
||
options = json.loads(post['poll_options'])
|
||
|
||
# Подготавливаем параметры для опроса
|
||
poll_kwargs = {
|
||
'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_kwargs['close_date'] = int(post['poll_end_date'].timestamp())
|
||
else:
|
||
poll_kwargs['close_date'] = post['poll_end_date']
|
||
|
||
# Повторные попытки отправки опроса при таймауте
|
||
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}")
|
||
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
|
||
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 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 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)}")
|
||
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())
|