diff --git a/tg_publish.py b/tg_publish.py index 44f9443..b71b018 100644 --- a/tg_publish.py +++ b/tg_publish.py @@ -6,9 +6,8 @@ import time import html import json import re +import httpx 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 @@ -63,6 +62,13 @@ 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: @@ -111,9 +117,15 @@ async def publish_to_tg(vk_post_id): logger.info(f"Запись VK ID {vk_post_id} не найдена или уже опубликована") return False - bot = Bot(token=BOT_TOKEN) + # Создаем httpx клиент с отключенным HTTP/2 + httpx_client = httpx.AsyncClient( + http2=False, + timeout=20.0, + follow_redirects=True + ) - # Проверяем наличие организационного хештега + try: + # Проверяем наличие организационного хештега is_org_message = has_org_tag(post['text']) # Предварительный расчет строки ссылок (с запасом 20 символов для message_id) @@ -161,126 +173,305 @@ async def publish_to_tg(vk_post_id): logger.info(f"=== Конец отформатированного текста ===") # Публикация первоначального сообщения без ссылок - 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'Оригинал в ВК') + 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 post['is_event'] and USE_SUBSCRIPTION_BOT: - updated_links.append(f'🔜 Подписка') - - links_line = " | ".join(updated_links) - final_text = f"{text}\n\n{links_line}" + 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 - пробуем отправить текстовое сообщение + 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: + 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 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: + 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: + 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: + 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: + 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'): + 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: + 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 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'): + 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: + 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'): + 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: + 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) # задержка между обращениями к телеграм - # Редактируем сообщение, добавляя ссылки - if has_image and not is_org_message: + # Публикация опроса, если требуется + poll_message_id = None + if post['is_poll'] and post['poll_question'] and post['poll_options']: 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 - ) + # Парсим варианты ответа + 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: + 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: + 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 Exception as e: + logger.error(f"Ошибка при публикации опроса для записи VK ID {vk_post_id}: {e}") - await asyncio.sleep(2) # задержка между обращениями к телеграм + # Обновляем запись в базе данных + 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() - # Публикация опроса, если требуется - 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: + except RetryAfterException as e: logger.error(f"Получена ошибка FloodWait (429) при публикации записи VK ID {vk_post_id}: {e}") logger.error(f"Необходимо подождать {e.retry_after} секунд перед следующей попыткой") raise - except TelegramError as e: + except httpx.TimeoutException as e: + logger.error(f"ТАЙМАУТ при публикации записи VK ID {vk_post_id}: {e}") + 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}: {e}") return False except Exception as e: @@ -325,7 +516,7 @@ async def publish_to_tg_all(): if post_count < len(posts): await asyncio.sleep(2) - except RetryAfter as e: + except RetryAfterException as e: logger.error(f"Прерываем публикацию из-за ошибки FloodWait. Ожидание: {e.retry_after} секунд") break except Exception as e: