[ВОЛК] переход tg_publish на httpx_client
ci/woodpecker/push/woodpecker Pipeline failed

This commit is contained in:
2025-12-21 19:26:09 +03:00
parent 5c46c0bcf0
commit 0d411bb5ce
+306 -115
View File
@@ -6,9 +6,8 @@ import time
import html import html
import json import json
import re import re
import httpx
from datetime import datetime, timezone from datetime import datetime, timezone
from telegram import Bot, PollOption
from telegram.error import TelegramError, RetryAfter
from dotenv import load_dotenv from dotenv import load_dotenv
from formatter import get_event_text from formatter import get_event_text
@@ -63,6 +62,13 @@ file_handler = RetryFileHandler(LOG_FILE, encoding='utf-8')
file_handler.setFormatter(formatter) file_handler.setFormatter(formatter)
logger.addHandler(file_handler) 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): def prepare_text(text):
"""Подготовка текстовых полей с экранированием HTML-сущностей""" """Подготовка текстовых полей с экранированием HTML-сущностей"""
if not text: if not text:
@@ -111,9 +117,15 @@ async def publish_to_tg(vk_post_id):
logger.info(f"Запись VK ID {vk_post_id} не найдена или уже опубликована") logger.info(f"Запись VK ID {vk_post_id} не найдена или уже опубликована")
return False 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']) is_org_message = has_org_tag(post['text'])
# Предварительный расчет строки ссылок (с запасом 20 символов для message_id) # Предварительный расчет строки ссылок (с запасом 20 символов для message_id)
@@ -161,126 +173,305 @@ async def publish_to_tg(vk_post_id):
logger.info(f"=== Конец отформатированного текста ===") logger.info(f"=== Конец отформатированного текста ===")
# Публикация первоначального сообщения без ссылок # Публикация первоначального сообщения без ссылок
if has_image and not is_org_message: send_photo_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendPhoto"
try: send_message_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendMessage"
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'<a href="{post["vk_post_url"]}">Оригинал в ВК</a>')
if post['is_event'] and USE_SUBSCRIPTION_BOT: if has_image and not is_org_message:
updated_links.append(f'<a href="https://t.me/{RESPONDER_BOT_NAME}?start=post_{message_id}">🔜 Подписка</a>') try:
photo_payload = {
links_line = " | ".join(updated_links) 'chat_id': CHANNEL_ID,
final_text = f"{text}\n\n{links_line}" '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'<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) # задержка между обращениями к телеграм
# Редактируем сообщение, добавляя ссылки
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) # задержка между обращениями к телеграм 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: try:
await bot.edit_message_caption( # Парсим варианты ответа
chat_id=CHANNEL_ID, options = json.loads(post['poll_options'])
message_id=message_id,
caption=final_text, # Отправляем опрос (в каналах можно отправлять только анонимные опросы)
parse_mode="HTML" send_poll_url = f"https://api.telegram.org/bot{BOT_TOKEN}/sendPoll"
) poll_payload = {
except TelegramError as e: 'chat_id': CHANNEL_ID,
logger.warning(f"Не удалось отредактировать подпись для записи VK ID {vk_post_id}: {e}. Пробуем отредактировать текстовое сообщение.") 'question': post['poll_question'],
await bot.edit_message_text( 'options': options,
chat_id=CHANNEL_ID, 'is_anonymous': True, # В каналах можно отправлять только анонимные опросы
message_id=message_id, 'allows_multiple_answers': post['poll_multiple'],
text=final_text, 'disable_notification': PUBLISH_SILENTLY
parse_mode="HTML", }
disable_web_page_preview=True
) # Добавляем close_date, если указан
else: if post['poll_end_date']:
await bot.edit_message_text( # Преобразуем datetime в Unix timestamp
chat_id=CHANNEL_ID, if isinstance(post['poll_end_date'], datetime):
message_id=message_id, poll_payload['close_date'] = int(post['poll_end_date'].timestamp())
text=final_text, else:
parse_mode="HTML", poll_payload['close_date'] = post['poll_end_date']
disable_web_page_preview=True
) 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()
# Публикация опроса, если требуется except RetryAfterException as e:
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"Получена ошибка FloodWait (429) при публикации записи VK ID {vk_post_id}: {e}")
logger.error(f"Необходимо подождать {e.retry_after} секунд перед следующей попыткой") logger.error(f"Необходимо подождать {e.retry_after} секунд перед следующей попыткой")
raise 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}") logger.error(f"Ошибка Telegram при публикации записи VK ID {vk_post_id}: {e}")
return False return False
except Exception as e: except Exception as e:
@@ -325,7 +516,7 @@ async def publish_to_tg_all():
if post_count < len(posts): if post_count < len(posts):
await asyncio.sleep(2) await asyncio.sleep(2)
except RetryAfter as e: except RetryAfterException as e:
logger.error(f"Прерываем публикацию из-за ошибки FloodWait. Ожидание: {e.retry_after} секунд") logger.error(f"Прерываем публикацию из-за ошибки FloodWait. Ожидание: {e.retry_after} секунд")
break break
except Exception as e: except Exception as e: