Files
Zilant2025/tg_publish.py
T
gitadmin 7608afa77b
ci/woodpecker/push/woodpecker Pipeline was successful
[ВОЛК] фикс ConnectionCls
2026-01-10 18:16:47 +03:00

1002 lines
61 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import os
import logging
import asyncio
import pymysql
import time
import html
import json
import re
import httpx
import requests
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
from urllib3 import Timeout as Urllib3Timeout
import traceback
import io
import sys
from datetime import datetime, timezone
from dotenv import load_dotenv
from formatter import get_event_text
# Настройка 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
def _send_photo_sync(send_photo_url, image_data, caption, chat_id, disable_notification, timeout):
"""
Синхронная функция для отправки фото через requests.
Вызывается из отдельного потока для избежания проблем с event loop в gunicorn.
Args:
send_photo_url: URL для отправки фото в Telegram API
image_data: Байты изображения
caption: Подпись к фото
chat_id: ID чата/канала
disable_notification: Отключить уведомления
timeout: Кортеж (connect_timeout, read_timeout, write_timeout) или число для общего таймаута
Returns:
requests.Response объект
"""
import socket
from urllib3.connection import HTTPSConnection, HTTPConnection
try:
# Используем requests для стабильной работы в gunicorn
# requests более надежен для синхронных операций в контексте WSGI
image_file = io.BytesIO(image_data)
files = {
'photo': ('image.jpg', image_file, 'image/jpeg')
}
data = {
'chat_id': chat_id,
'caption': caption,
'parse_mode': 'HTML',
'disable_notification': str(disable_notification).lower()
}
# Извлекаем таймауты из кортежа
if isinstance(timeout, tuple):
if len(timeout) >= 3:
connect_timeout, read_timeout, write_timeout = timeout[0], timeout[1], timeout[2]
elif len(timeout) == 2:
connect_timeout, read_timeout = timeout[0], timeout[1]
write_timeout = read_timeout # Используем read timeout для write
else:
connect_timeout = read_timeout = write_timeout = timeout[0]
else:
connect_timeout = read_timeout = write_timeout = timeout
# Создаем сессию с кастомным адаптером
session = requests.Session()
# Создаем кастомный HTTPAdapter с установкой socket timeout
# Используем максимальный timeout для всех операций (connect, read, write)
class CustomHTTPAdapter(HTTPAdapter):
def init_poolmanager(self, *args, **kwargs):
# Устанавливаем socket timeout через socket_options
socket_options = kwargs.get('socket_options', [])
socket_options.append((socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1))
kwargs['socket_options'] = socket_options
# Устанавливаем timeout для пула соединений
# Используем максимальный таймаут для всех операций
max_timeout = max(connect_timeout, read_timeout, write_timeout)
if 'timeout' not in kwargs:
kwargs['timeout'] = Urllib3Timeout(
connect=connect_timeout,
read=max_timeout
)
return super().init_poolmanager(*args, **kwargs)
# Монтируем кастомный адаптер
adapter = CustomHTTPAdapter(max_retries=Retry(total=0))
session.mount('http://', adapter)
session.mount('https://', adapter)
# Отправляем запрос с увеличенным таймаутом
# Используем максимальный таймаут для всех операций
max_timeout = max(connect_timeout, read_timeout, write_timeout)
response = session.post(
send_photo_url,
files=files,
data=data,
timeout=(connect_timeout, max_timeout), # (connect, read)
stream=False # Отключаем потоковую передачу для стабильности
)
return response
except requests.exceptions.Timeout as e:
# Пробрасываем таймаут как есть, чтобы его можно было обработать выше
raise
except requests.exceptions.RequestException as e:
# Пробрасываем другие ошибки requests как есть
raise
except Exception as e:
# Пробрасываем остальные исключения как есть
raise
finally:
# Закрываем сессию
if 'session' in locals():
session.close()
def _edit_caption_sync(edit_caption_url, chat_id, message_id, caption, timeout):
"""
Синхронная функция для редактирования caption через requests.
Вызывается из отдельного потока для избежания проблем с event loop в gunicorn.
Args:
edit_caption_url: URL для редактирования caption в Telegram API
chat_id: ID чата/канала
message_id: ID сообщения для редактирования
caption: Новый текст caption
timeout: Кортеж (connect_timeout, read_timeout) или число для общего таймаута
Returns:
requests.Response объект
"""
try:
data = {
'chat_id': chat_id,
'message_id': message_id,
'caption': caption,
'parse_mode': 'HTML'
}
response = requests.post(
edit_caption_url,
json=data,
timeout=timeout
)
return response
except requests.exceptions.Timeout as e:
raise
except requests.exceptions.RequestException as e:
raise
except Exception as e:
raise
def _edit_text_sync(edit_text_url, chat_id, message_id, text, timeout):
"""
Синхронная функция для редактирования текста через requests.
Вызывается из отдельного потока для избежания проблем с event loop в gunicorn.
Args:
edit_text_url: URL для редактирования текста в Telegram API
chat_id: ID чата/канала
message_id: ID сообщения для редактирования
text: Новый текст сообщения
timeout: Кортеж (connect_timeout, read_timeout) или число для общего таймаута
Returns:
requests.Response объект
"""
try:
data = {
'chat_id': chat_id,
'message_id': message_id,
'text': text,
'parse_mode': 'HTML',
'disable_web_page_preview': True
}
response = requests.post(
edit_text_url,
json=data,
timeout=timeout
)
return response
except requests.exceptions.Timeout as e:
raise
except requests.exceptions.RequestException as e:
raise
except Exception as 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
# Создаем httpx клиент с отключенным HTTP/2
# Используем более детальные таймауты: connect, read, write, pool
# Таймаут на запись увеличен для больших файлов (до 10 МБ)
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
)
# Переменная для отслеживания успешности публикации
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"=== Конец отформатированного текста ===")
# Публикация первоначального сообщения без ссылок
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:
# Логируем информацию об изображении перед отправкой
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
# Вычисляем динамический таймаут на запись в зависимости от размера файла
# Базовый таймаут 60 секунд + 10 секунд на каждый МБ (минимум 120 секунд)
file_size_mb = len(image_data) / (1024 * 1024)
write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0))
logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд")
# Используем requests для отправки файлов через отдельный поток
# Это необходимо для стабильной работы с gunicorn, который использует синхронные воркеры
# requests более надежен для синхронных операций в контексте WSGI
logger.info(f"Используем requests для отправки файла размером {file_size_mb:.2f} МБ через отдельный поток")
# Логируем время начала запроса
request_start_time = time.time()
logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}")
# Вызываем синхронную функцию через asyncio.to_thread()
# Это позволяет избежать проблем с event loop в gunicorn
try:
# Используем кортеж для таймаута: (connect, read, write)
# write_timeout применяется к операциям записи в сокет
timeout_tuple = (10.0, write_timeout, write_timeout) # (connect, read, write)
photo_response = await asyncio.to_thread(
_send_photo_sync,
send_photo_url,
image_data,
text,
CHANNEL_ID,
PUBLISH_SILENTLY,
timeout_tuple
)
except Exception as sync_error:
# Обработка ошибок от синхронной функции
request_duration = time.time() - request_start_time
error_msg = str(sync_error)
# Проверяем, является ли это таймаутом
is_timeout = (
'Timeout' in error_msg or
isinstance(sync_error, requests.exceptions.Timeout) or
'timeout' in error_msg.lower()
)
if is_timeout:
timeout_type = type(sync_error).__name__
logger.error(f"ТАЙМАУТ ({timeout_type}) при отправке изображения (requests) для записи VK ID {vk_post_id}")
logger.error(f"Время до таймаута: {request_duration:.2f} секунд")
logger.error(f"Тип исключения: {timeout_type}")
logger.error(f"Сообщение об ошибке: {error_msg}")
logger.error(f"Полная информация об исключении: {repr(sync_error)}")
logger.error(f"URL изображения: {image_url}")
logger.error(f"Размер файла: {file_size_mb:.2f} МБ")
logger.error(f"Запрос, вызвавший таймаут: POST {send_photo_url}")
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
raise # Пробрасываем таймаут наверх
else:
# Другие ошибки
logger.error(f"Ошибка при отправке изображения (requests) для записи VK ID {vk_post_id}: {error_msg}")
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
raise
request_duration = time.time() - request_start_time
logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд")
# Обработка ответа от requests
if photo_response.status_code == 200:
try:
photo_result = photo_response.json()
except Exception as json_error:
logger.error(f"Ошибка при парсинге JSON ответа для записи VK ID {vk_post_id}: {json_error}")
logger.error(f"Текст ответа: {photo_response.text[:500]}")
raise Exception(f"Ошибка парсинга JSON ответа: {json_error}")
if photo_result.get('ok'):
message_data = photo_result.get('result', {})
message_id = message_data.get('message_id')
message_has_image = True # Успешно отправили изображение
else:
# Ошибка в ответе API - логируем полную информацию
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Статус код: {photo_response.status_code}")
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Полный ответ от Telegram API: {json.dumps(photo_result, ensure_ascii=False, indent=2)}")
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:
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Статус код: {photo_response.status_code}")
logger.error(f"[Отправка фото для VK ID {vk_post_id}] Текст ответа: {photo_response.text[:500]}")
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[:500]}")
except RetryAfterException:
raise
except Exception as e:
# Проверяем, является ли это таймаутом
# Может быть от requests или httpx
is_timeout = (
isinstance(e, (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout)) or
isinstance(e, requests.exceptions.Timeout) or
'Timeout' in str(e) or
'timeout' in str(e).lower()
)
if is_timeout:
timeout_type = type(e).__name__
request_duration = time.time() - request_start_time if 'request_start_time' in locals() else 0
logger.error(f"ТАЙМАУТ ({timeout_type}) при отправке изображения для записи VK ID {vk_post_id}")
if request_duration > 0:
logger.error(f"Время до таймаута: {request_duration:.2f} секунд")
logger.error(f"Тип исключения: {timeout_type}")
logger.error(f"Сообщение об ошибке: {str(e)}")
logger.error(f"Полная информация об исключении: {repr(e)}")
logger.error(f"URL изображения: {image_url}")
if image_data is not None:
logger.error(f"Размер файла: {len(image_data) / (1024 * 1024):.2f} МБ")
logger.error(f"Запрос, вызвавший таймаут: POST {send_photo_url}")
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
raise # Пробрасываем таймаут наверх
# Если это не таймаут, пробрасываем дальше для обработки как обычная ошибка
raise
except Exception as e:
logger.warning(f"Не удалось отправить изображение для записи VK ID {vk_post_id}: {e}. Отправляем текстовое сообщение.")
message_has_image = False # Отправляем текстовое сообщение вместо изображения
text_payload = {
'chat_id': CHANNEL_ID,
'text': text,
'parse_mode': 'HTML',
'disable_notification': PUBLISH_SILENTLY,
'disable_web_page_preview': True
}
try:
text_response = await httpx_client.post(send_message_url, json=text_payload)
except httpx.TimeoutException as timeout_e:
logger.error(f"ТАЙМАУТ при отправке текста (fallback) для записи VK ID {vk_post_id}")
logger.error(f"Тип исключения: {type(timeout_e).__name__}")
logger.error(f"Сообщение об ошибке: {str(timeout_e)}")
logger.error(f"Полная информация об исключении: {repr(timeout_e)}")
if hasattr(timeout_e, 'request'):
logger.error(f"Запрос, вызвавший таймаут: {timeout_e.request.method} {timeout_e.request.url if hasattr(timeout_e.request, 'url') else 'N/A'}")
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
raise # Пробрасываем таймаут наверх
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:
log_telegram_error(text_response, f"[Отправка текста (fallback) для VK ID {vk_post_id}]")
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:
log_telegram_error(text_response, f"[Отправка текста (fallback) для VK ID {vk_post_id}]")
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:
# Если это организационное сообщение или нет изображения, отправляем текстовое сообщение
message_has_image = False # Отправляем текстовое сообщение
text_payload = {
'chat_id': CHANNEL_ID,
'text': text,
'parse_mode': 'HTML',
'disable_notification': PUBLISH_SILENTLY,
'disable_web_page_preview': True
}
try:
text_response = await httpx_client.post(send_message_url, json=text_payload)
except httpx.TimeoutException 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)}")
if hasattr(e, 'request'):
logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}")
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
raise # Пробрасываем таймаут наверх
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:
log_telegram_error(text_response, f"[Отправка текста для VK ID {vk_post_id}]")
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:
log_telegram_error(text_response, f"[Отправка текста для VK ID {vk_post_id}]")
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"
# Используем message_has_image для определения типа сообщения
# Если было отправлено изображение, редактируем caption, иначе - текст
if message_has_image:
# Сообщение с изображением - редактируем caption через requests
edit_caption_success = False
try:
# Используем requests через asyncio.to_thread() для стабильной работы в gunicorn
# Увеличиваем таймаут для редактирования caption
edit_timeout = (10.0, 60.0, 60.0) # (connect, read, write)
edit_caption_response = await asyncio.to_thread(
_edit_caption_sync,
edit_caption_url,
CHANNEL_ID,
message_id,
final_text,
edit_timeout
)
if edit_caption_response.status_code == 200:
try:
edit_result = edit_caption_response.json()
except Exception as json_error:
logger.error(f"Ошибка при парсинге JSON ответа при редактировании caption для записи VK ID {vk_post_id}: {json_error}")
logger.error(f"Текст ответа: {edit_caption_response.text[:500]}")
raise Exception(f"Ошибка парсинга JSON ответа: {json_error}")
if edit_result.get('ok'):
edit_caption_success = True
logger.info(f"Caption успешно отредактирован для записи VK ID {vk_post_id}")
else:
# Ошибка в ответе API
log_telegram_error(edit_caption_response, f"[Редактирование подписи для VK ID {vk_post_id}]")
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)
# Не критичная ошибка - сообщение уже опубликовано
logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: {edit_result.get('description', 'Unknown error')}. Сообщение опубликовано, но без ссылок.")
else:
log_telegram_error(edit_caption_response, f"[Редактирование подписи для VK ID {vk_post_id}]")
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)
# Не критичная ошибка - сообщение уже опубликовано
logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: HTTP {edit_caption_response.status_code}. Сообщение опубликовано, но без ссылок.")
except RetryAfterException:
raise
except Exception as e:
# Проверяем, является ли это таймаутом
is_timeout = (
isinstance(e, requests.exceptions.Timeout) or
'Timeout' in str(e) or
'timeout' in str(e).lower() or
'disconnected' in str(e).lower()
)
if is_timeout:
logger.warning(f"Таймаут при редактировании caption для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.")
else:
logger.warning(f"Не удалось отредактировать caption для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.")
# Не пробрасываем исключение - сообщение уже опубликовано, просто без ссылок
# Это считается частичным успехом
else:
# Текстовое сообщение - редактируем текст через requests
edit_text_success = False
try:
# Используем requests через asyncio.to_thread() для стабильной работы в gunicorn
# Увеличиваем таймаут для редактирования текста
edit_timeout = (10.0, 60.0, 60.0) # (connect, read, write)
edit_text_response = await asyncio.to_thread(
_edit_text_sync,
edit_text_url,
CHANNEL_ID,
message_id,
final_text,
edit_timeout
)
if edit_text_response.status_code == 200:
try:
edit_result = edit_text_response.json()
except Exception as json_error:
logger.error(f"Ошибка при парсинге JSON ответа при редактировании текста для записи VK ID {vk_post_id}: {json_error}")
logger.error(f"Текст ответа: {edit_text_response.text[:500]}")
raise Exception(f"Ошибка парсинга JSON ответа: {json_error}")
if edit_result.get('ok'):
edit_text_success = True
logger.info(f"Текст успешно отредактирован для записи VK ID {vk_post_id}")
else:
# Ошибка в ответе API
log_telegram_error(edit_text_response, f"[Редактирование текста для VK ID {vk_post_id}]")
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)
# Не критичная ошибка - сообщение уже опубликовано
logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: {edit_result.get('description', 'Unknown error')}. Сообщение опубликовано, но без ссылок.")
else:
log_telegram_error(edit_text_response, f"[Редактирование текста для VK ID {vk_post_id}]")
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)
# Не критичная ошибка - сообщение уже опубликовано
logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: HTTP {edit_text_response.status_code}. Сообщение опубликовано, но без ссылок.")
except RetryAfterException:
raise
except Exception as e:
# Проверяем, является ли это таймаутом
is_timeout = (
isinstance(e, requests.exceptions.Timeout) or
'Timeout' in str(e) or
'timeout' in str(e).lower() or
'disconnected' in str(e).lower()
)
if is_timeout:
logger.warning(f"Таймаут при редактировании текста для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.")
else:
logger.warning(f"Не удалось отредактировать текст для записи VK ID {vk_post_id}: {e}. Сообщение опубликовано, но без ссылок.")
# Не пробрасываем исключение - сообщение уже опубликовано, просто без ссылок
# Это считается частичным успехом
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']
try:
poll_response = await httpx_client.post(send_poll_url, json=poll_payload)
except httpx.TimeoutException 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)}")
if hasattr(e, 'request'):
logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}")
logger.error(f"Трассировка стека:\n{traceback.format_exc()}")
raise # Пробрасываем таймаут наверх
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:
log_telegram_error(poll_response, f"[Отправка опроса для VK ID {vk_post_id}]")
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:
log_telegram_error(poll_response, f"[Отправка опроса для VK ID {vk_post_id}]")
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 httpx.TimeoutException as e:
logger.error(f"ТАЙМАУТ при публикации опроса для записи VK ID {vk_post_id}")
logger.error(f"Детали таймаута: {type(e).__name__}: {e}")
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 httpx.TimeoutException 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)}")
if hasattr(e, 'request'):
logger.error(f"Запрос, вызвавший таймаут: {e.request.method} {e.request.url if hasattr(e.request, 'url') else 'N/A'}")
# Пытаемся получить информацию о таймауте из исключения
if hasattr(e, 'timeout'):
logger.error(f"Настройки таймаута: {e.timeout}")
# Проверяем, запущены ли мы из gunicorn
import sys
if 'gunicorn' in sys.modules:
logger.warning("Обнаружен gunicorn - возможен конфликт с event loop при запуске асинхронного кода в отдельном потоке")
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())