Files
Zilant2025/tg_publish.py
T
gitadmin 26d12a4758
ci/woodpecker/push/woodpecker Pipeline was successful
[ВОЛК] коррекция синтаксиса
2025-12-21 19:37:24 +03:00

537 lines
29 KiB
Python

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))
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
)
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"=== Конец отформатированного текста ===")
# Публикация первоначального сообщения без ссылок
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 - пробуем отправить текстовое сообщение
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) # задержка между обращениями к телеграм
# Публикация опроса, если требуется
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:
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}")
# Обновляем запись в базе данных
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}: {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:
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 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())