833 lines
44 KiB
Python
833 lines
44 KiB
Python
import os
|
||
import logging
|
||
import asyncio
|
||
import pymysql
|
||
import time
|
||
import html
|
||
import httpx
|
||
from datetime import datetime, timezone, timedelta
|
||
from urllib.parse import urlparse, unquote
|
||
from dotenv import load_dotenv
|
||
from formatter import get_event_text
|
||
from telegram_relay import (
|
||
bot_api_method_url,
|
||
EDIT_LINKS_INITIAL_DELAY_SEC,
|
||
EDIT_LINKS_MAX_ATTEMPTS,
|
||
EDIT_LINKS_RETRY_INTERVAL_SEC,
|
||
published_text_has_links,
|
||
)
|
||
from season_links import subscription_start_link
|
||
|
||
# Загрузка переменных окружения
|
||
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')
|
||
DESC_PREFIX = os.getenv('DESC_PREFIX')
|
||
PZK_PREFIX = os.getenv('PZK_PREFIX')
|
||
LOG_FILE = os.getenv('LOG_FILE', 'evtg_publisher.log')
|
||
EVENT_POST_DELAY = int(os.getenv('EVENT_POST_DELAY', 0)) # Задержка в минутах
|
||
WORKMODE = os.getenv('WORKMODE', 'ZILANT') # Режим работы: ZILANT или VOLK
|
||
|
||
# Параметры длины сообщений
|
||
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')
|
||
|
||
# Настройка логирования
|
||
logger = logging.getLogger('TG_EV_post')
|
||
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 format_tags(tags_str):
|
||
"""Форматирование тегов с добавлением #"""
|
||
if not tags_str:
|
||
return ""
|
||
tags = tags_str.split()
|
||
formatted_tags = []
|
||
for tag in tags:
|
||
if len(tag) > 1 and not tag.startswith('#'):
|
||
formatted_tags.append(f"#{tag}")
|
||
else:
|
||
formatted_tags.append(tag)
|
||
return " ".join(formatted_tags)
|
||
|
||
def prepare_text(text):
|
||
"""Подготовка текстовых полей с экранированием HTML-сущностей"""
|
||
if not text:
|
||
return ""
|
||
|
||
# Экранируем специальные символы HTML
|
||
text = html.escape(str(text))
|
||
|
||
return text
|
||
|
||
|
||
def _filename_from_image_url(image_url):
|
||
path = unquote(urlparse(image_url).path)
|
||
name = os.path.basename(path) or "image.jpg"
|
||
if "." not in name:
|
||
name = f"{name}.jpg"
|
||
return name
|
||
|
||
|
||
async def download_image(image_url, httpx_client):
|
||
"""Скачивает картинку самим сервером (как tg_publish), не отдаёт URL Telegram.
|
||
|
||
sendPhoto с photo=http://... Telegram качает сам и часто отвечает
|
||
«failed to get HTTP URL content» (http, недоступность с DC Telegram).
|
||
"""
|
||
logger.info(f"Скачивание изображения: {image_url}")
|
||
started = 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()
|
||
image_data = response.content
|
||
content_type = (response.headers.get("content-type") or "image/jpeg").split(";")[0].strip()
|
||
if not content_type.startswith("image/"):
|
||
content_type = "image/jpeg"
|
||
logger.info(
|
||
f"Изображение скачано за {time.time() - started:.2f} с, "
|
||
f"размер: {len(image_data)} байт, type={content_type}"
|
||
)
|
||
return image_data, content_type, _filename_from_image_url(image_url)
|
||
|
||
|
||
async def _edit_event_add_links(httpx_client, has_image, message_id, final_text, entity_label):
|
||
"""Edit caption/text со ссылками; повтор до появления href= в ответе API."""
|
||
await asyncio.sleep(EDIT_LINKS_INITIAL_DELAY_SEC)
|
||
method = "editMessageCaption" if has_image else "editMessageText"
|
||
field = "caption" if has_image else "text"
|
||
api_url = bot_api_method_url(BOT_TOKEN, method)
|
||
|
||
logger.info(
|
||
f"Попытка редактирования {entity_label}: message_id={message_id}, "
|
||
f"final_text_length={len(final_text)}"
|
||
)
|
||
|
||
for attempt in range(1, EDIT_LINKS_MAX_ATTEMPTS + 1):
|
||
payload = {
|
||
'chat_id': CHANNEL_ID,
|
||
'message_id': message_id,
|
||
field: final_text,
|
||
'parse_mode': 'HTML',
|
||
}
|
||
if not has_image:
|
||
payload['disable_web_page_preview'] = True
|
||
|
||
try:
|
||
response = await httpx_client.post(api_url, json=payload)
|
||
except httpx.TimeoutException as e:
|
||
if attempt >= EDIT_LINKS_MAX_ATTEMPTS:
|
||
logger.error(
|
||
f"Таймаут edit ссылок для {entity_label} после "
|
||
f"{EDIT_LINKS_MAX_ATTEMPTS} попыток: {e}"
|
||
)
|
||
return False
|
||
logger.warning(
|
||
f"Таймаут edit (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) для {entity_label}"
|
||
)
|
||
await asyncio.sleep(EDIT_LINKS_RETRY_INTERVAL_SEC)
|
||
continue
|
||
|
||
if response.status_code == 200:
|
||
data = response.json()
|
||
if data.get('ok'):
|
||
result_text = (data.get('result') or {}).get(field) or ""
|
||
if published_text_has_links(result_text):
|
||
logger.info(
|
||
f"Ссылки добавлены для {entity_label} "
|
||
f"(попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS})"
|
||
)
|
||
return True
|
||
logger.warning(
|
||
f"edit ok, но href= не найден (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label}"
|
||
)
|
||
else:
|
||
error_code = data.get('error_code')
|
||
error_description = (data.get('description') or '').lower()
|
||
if error_code == 429:
|
||
retry_after = data.get('parameters', {}).get('retry_after', 60)
|
||
raise Exception(f"RetryAfter:{retry_after}")
|
||
if error_code == 400 and "not modified" in error_description:
|
||
logger.warning(
|
||
f"message is not modified (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label} — повтор"
|
||
)
|
||
else:
|
||
logger.warning(
|
||
f"API error edit (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label}: {data.get('description')}"
|
||
)
|
||
elif response.status_code == 429:
|
||
retry_after = 60
|
||
try:
|
||
retry_after = response.json().get('parameters', {}).get('retry_after', 60)
|
||
except Exception:
|
||
pass
|
||
raise Exception(f"RetryAfter:{retry_after}")
|
||
elif response.status_code == 400:
|
||
try:
|
||
error_description = (response.json().get('description') or '').lower()
|
||
if "not modified" in error_description:
|
||
logger.warning(
|
||
f"message is not modified HTTP 400 (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label} — повтор"
|
||
)
|
||
else:
|
||
logger.warning(
|
||
f"HTTP 400 edit (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label}: {response.text[:300]}"
|
||
)
|
||
except Exception:
|
||
logger.warning(
|
||
f"HTTP 400 edit (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label}: {response.text[:300]}"
|
||
)
|
||
else:
|
||
logger.warning(
|
||
f"HTTP {response.status_code} edit (попытка {attempt}/{EDIT_LINKS_MAX_ATTEMPTS}) "
|
||
f"для {entity_label}: {response.text[:300]}"
|
||
)
|
||
|
||
if attempt < EDIT_LINKS_MAX_ATTEMPTS:
|
||
await asyncio.sleep(EDIT_LINKS_RETRY_INTERVAL_SEC)
|
||
|
||
logger.error(
|
||
f"Не удалось добавить строку ссылок для {entity_label} "
|
||
f"после {EDIT_LINKS_MAX_ATTEMPTS} попыток"
|
||
)
|
||
return False
|
||
|
||
|
||
async def tg_post_event(httpx_client, event_data):
|
||
"""Публикация одного события"""
|
||
try:
|
||
# Форматирование текста
|
||
tags = format_tags(event_data['tags'])
|
||
unit_name_raw = event_data.get('unit_name') or ""
|
||
name = prepare_text(event_data['name'] or "")
|
||
|
||
# Формируем название площадки с возможной ссылкой
|
||
unit_name = ""
|
||
if unit_name_raw:
|
||
unit_tg_id = event_data.get('unit_tg_id')
|
||
unit_title = event_data.get('unit_title') or unit_name_raw
|
||
unit_title_escaped = prepare_text(unit_title)
|
||
|
||
if unit_tg_id:
|
||
# Если площадка опубликована, делаем ссылку на сообщение
|
||
# Используем логику из format_channel_link из tg_mainbot.py
|
||
if CHANNEL_ID.startswith('@'):
|
||
# Публичный канал
|
||
channel_link = f"https://t.me/{CHANNEL_ID[1:]}/{unit_tg_id}"
|
||
elif CHANNEL_ID.startswith('-100'):
|
||
# Приватный канал с префиксом -100
|
||
channel_link = f"https://t.me/c/{CHANNEL_ID[4:]}/{unit_tg_id}"
|
||
else:
|
||
# Приватный канал без префикса или другой формат
|
||
channel_link = f"https://t.me/c/{CHANNEL_ID}/{unit_tg_id}"
|
||
|
||
# Формируем ссылку на сообщение
|
||
unit_name = f'<a href="{channel_link}">{unit_title_escaped}</a>'
|
||
else:
|
||
# Если площадка не опубликована, просто текст
|
||
unit_name = unit_title_escaped
|
||
|
||
# Получаем и подготавливаем описание события
|
||
about_raw = event_data['about']
|
||
|
||
# Формируем базовый текст (теги + название площадки + название события)
|
||
caption_parts = []
|
||
if tags:
|
||
caption_parts.append(tags)
|
||
if unit_name:
|
||
caption_parts.append(unit_name)
|
||
if name:
|
||
caption_parts.append(f"<b>{name}</b>")
|
||
base_text = "\n".join(caption_parts)
|
||
base_text_length = len(base_text)
|
||
|
||
# Предварительный расчет строки ссылок (с запасом 20 символов для message_id)
|
||
number = event_data['number']
|
||
placeholder_message_id = '0' * 20 # Заполнитель для message_id
|
||
|
||
# Формируем ссылки в зависимости от USE_SUBSCRIPTION_BOT и наличия DESC_PREFIX
|
||
links_parts = []
|
||
if DESC_PREFIX:
|
||
links_parts.append(f'🔎 <a href="{DESC_PREFIX}{number}/">Инфо</a>')
|
||
if USE_SUBSCRIPTION_BOT:
|
||
links_parts.append(f'📌 <a href="{PZK_PREFIX}{number}/">Иду</a>')
|
||
links_parts.append(
|
||
f'🔜 <a href="{subscription_start_link(RESPONDER_BOT_NAME, "event", placeholder_message_id)}">Подписка</a>'
|
||
)
|
||
|
||
if links_parts:
|
||
links_line_placeholder = ' | '.join(links_parts)
|
||
links_line_length = len(links_line_placeholder)
|
||
else:
|
||
links_line_placeholder = ""
|
||
links_line_length = 0
|
||
|
||
# Определяем максимальную длину в зависимости от типа сообщения
|
||
image_url = event_data['about_social_picture']
|
||
has_image = image_url and image_url.strip()
|
||
max_length = MAX_CAPTION_LENGTH if has_image else MAX_TEXT_LENGTH
|
||
|
||
# Вычисляем доступную длину для текста события
|
||
available_length = max_length - base_text_length - links_line_length - 3 # -3 для символов переноса строки
|
||
if available_length < 0:
|
||
available_length = 0
|
||
|
||
# Получаем обработанный текст с учетом доступной длины
|
||
about = get_event_text(about_raw, available_length) if about_raw else ""
|
||
|
||
# Формируем первоначальный текст
|
||
initial_text = f"{base_text}\n\n{about}" if about else base_text
|
||
|
||
# Логируем детальную информацию о данных для публикации
|
||
logger.info(f"=== Детали публикации события номер {event_data['number']} ===")
|
||
logger.info(f"Длина initial_text: {len(initial_text)} символов")
|
||
logger.info(f"Длина base_text: {base_text_length} символов")
|
||
logger.info(f"Длина about: {len(about)} символов")
|
||
logger.info(f"Длина links_line_placeholder: {links_line_length} символов")
|
||
logger.info(f"Максимальная длина (max_length): {max_length} символов")
|
||
logger.info(f"Доступная длина для about: {available_length} символов")
|
||
logger.info(f"Есть изображение: {has_image}")
|
||
if has_image:
|
||
logger.info(f"URL изображения: {image_url}")
|
||
logger.info(f"Длина URL изображения: {len(image_url)} символов")
|
||
logger.info(f"Текст для публикации: {initial_text[:200]}..." if len(initial_text) > 200 else f"Текст для публикации: {initial_text}")
|
||
|
||
# Публикация первоначального сообщения без ссылок
|
||
send_photo_url = bot_api_method_url(BOT_TOKEN, "sendPhoto")
|
||
send_message_url = bot_api_method_url(BOT_TOKEN, "sendMessage")
|
||
message_is_photo = False
|
||
|
||
if has_image:
|
||
try:
|
||
logger.info(
|
||
f"Попытка отправки фото для события {event_data['number']}: "
|
||
f"chat_id={CHANNEL_ID}, caption_length={len(initial_text)}, "
|
||
f"image_url={image_url}"
|
||
)
|
||
image_data, content_type, filename = await download_image(image_url, httpx_client)
|
||
photo_data = {
|
||
'chat_id': CHANNEL_ID,
|
||
'caption': initial_text,
|
||
'parse_mode': 'HTML',
|
||
'disable_notification': 'true' if PUBLISH_SILENTLY else 'false',
|
||
}
|
||
photo_files = {
|
||
'photo': (filename, image_data, content_type),
|
||
}
|
||
photo_response = await httpx_client.post(
|
||
send_photo_url, data=photo_data, files=photo_files
|
||
)
|
||
|
||
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')
|
||
message_is_photo = True
|
||
logger.info(f"Фото успешно отправлено для события {event_data['number']}, message_id={message_id}")
|
||
else:
|
||
# Ошибка в ответе API - пробуем отправить текстовое сообщение
|
||
error_code = photo_result.get('error_code')
|
||
if error_code == 429:
|
||
# FloodWait - пробрасываем как RetryAfter
|
||
retry_after = photo_result.get('parameters', {}).get('retry_after', 60)
|
||
raise Exception(f"RetryAfter:{retry_after}")
|
||
raise Exception(f"API error: {photo_result.get('description', 'Unknown error')}")
|
||
else:
|
||
if photo_response.status_code == 429:
|
||
# FloodWait
|
||
retry_after = 60
|
||
try:
|
||
error_data = photo_response.json()
|
||
retry_after = error_data.get('parameters', {}).get('retry_after', 60)
|
||
except:
|
||
pass
|
||
raise Exception(f"RetryAfter:{retry_after}")
|
||
raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text}")
|
||
|
||
except httpx.TimeoutException as e:
|
||
logger.error(f"ТАЙМАУТ при отправке изображения для события номер {event_data['number']}")
|
||
logger.error(f"Полная диагностическая информация об ошибке TimedOut:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e)}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
logger.error(f" Параметры запроса:")
|
||
logger.error(f" chat_id: {CHANNEL_ID}")
|
||
logger.error(f" image_url: {image_url}")
|
||
logger.error(f" caption_length: {len(initial_text)}")
|
||
logger.error(f" parse_mode: HTML")
|
||
raise
|
||
except Exception as e:
|
||
error_msg = str(e)
|
||
if error_msg.startswith("RetryAfter:"):
|
||
retry_after = int(error_msg.split(":")[1])
|
||
raise Exception(f"RetryAfter:{retry_after}")
|
||
logger.warning(f"Не удалось отправить изображение для события номер {event_data['number']}: {e}. Отправляем текстовое сообщение.")
|
||
logger.error(f"Полная информация об ошибке Telegram при отправке изображения для события {event_data['number']}:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e) if hasattr(e, '__dict__') else 'N/A'}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
logger.info(f"Попытка отправки текстового сообщения вместо фото для события {event_data['number']}")
|
||
try:
|
||
text_payload = {
|
||
'chat_id': CHANNEL_ID,
|
||
'text': initial_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')
|
||
logger.info(f"Текстовое сообщение успешно отправлено для события {event_data['number']}, message_id={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 Exception(f"RetryAfter:{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 Exception(f"RetryAfter:{retry_after}")
|
||
raise Exception(f"HTTP {text_response.status_code}: {text_response.text}")
|
||
except httpx.TimeoutException as e:
|
||
logger.error(f"ТАЙМАУТ при отправке текстового сообщения для события номер {event_data['number']}")
|
||
logger.error(f"Полная диагностическая информация об ошибке TimedOut:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e)}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
logger.error(f" Параметры запроса:")
|
||
logger.error(f" chat_id: {CHANNEL_ID}")
|
||
logger.error(f" text_length: {len(initial_text)}")
|
||
logger.error(f" parse_mode: HTML")
|
||
raise
|
||
except Exception as e:
|
||
error_msg = str(e)
|
||
if error_msg.startswith("RetryAfter:"):
|
||
raise
|
||
raise
|
||
else:
|
||
logger.info(f"Попытка отправки текстового сообщения для события {event_data['number']}: chat_id={CHANNEL_ID}, text_length={len(initial_text)}")
|
||
try:
|
||
text_payload = {
|
||
'chat_id': CHANNEL_ID,
|
||
'text': initial_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')
|
||
logger.info(f"Текстовое сообщение успешно отправлено для события {event_data['number']}, message_id={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 Exception(f"RetryAfter:{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 Exception(f"RetryAfter:{retry_after}")
|
||
raise Exception(f"HTTP {text_response.status_code}: {text_response.text}")
|
||
except httpx.TimeoutException as e:
|
||
logger.error(f"ТАЙМАУТ при отправке текстового сообщения для события номер {event_data['number']}")
|
||
logger.error(f"Полная диагностическая информация об ошибке TimedOut:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e)}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
logger.error(f" Параметры запроса:")
|
||
logger.error(f" chat_id: {CHANNEL_ID}")
|
||
logger.error(f" text_length: {len(initial_text)}")
|
||
logger.error(f" parse_mode: HTML")
|
||
raise
|
||
|
||
# Формируем финальный текст с ссылками
|
||
links_parts = []
|
||
if DESC_PREFIX:
|
||
links_parts.append(f'🔎 <a href="{DESC_PREFIX}{number}/">Инфо</a>')
|
||
if USE_SUBSCRIPTION_BOT:
|
||
links_parts.append(f'📌 <a href="{PZK_PREFIX}{number}/">Иду</a>')
|
||
links_parts.append(
|
||
f'🔜 <a href="{subscription_start_link(RESPONDER_BOT_NAME, "event", message_id)}">Подписка</a>'
|
||
)
|
||
|
||
if links_parts:
|
||
links_line = ' | '.join(links_parts)
|
||
final_text = f"{initial_text}\n\n{links_line}"
|
||
else:
|
||
final_text = initial_text
|
||
|
||
# Проверяем, изменился ли текст. Если нет - не редактируем сообщение
|
||
if final_text == initial_text:
|
||
logger.info(f"Финальный текст идентичен исходному для события {event_data['number']}, редактирование не требуется")
|
||
return message_id
|
||
|
||
entity_label = f"события {event_data['number']}"
|
||
links_ok = await _edit_event_add_links(
|
||
httpx_client, message_is_photo, message_id, final_text, entity_label
|
||
)
|
||
if not links_ok and message_is_photo:
|
||
links_ok = await _edit_event_add_links(
|
||
httpx_client, False, message_id, final_text, entity_label
|
||
)
|
||
if not links_ok:
|
||
logger.error(
|
||
f"Событие {event_data['number']}: сообщение отправлено (id={message_id}), "
|
||
f"но ссылки не добавлены"
|
||
)
|
||
return None
|
||
|
||
return message_id # Возвращаем только message_id
|
||
|
||
except httpx.TimeoutException as e:
|
||
logger.error(f"ТАЙМАУТ при публикации события номер {event_data['number']}")
|
||
logger.error(f"Полная диагностическая информация об ошибке TimedOut:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e)}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
logger.error(f" Контекст события:")
|
||
logger.error(f" number: {event_data.get('number')}")
|
||
logger.error(f" has_image: {bool(event_data.get('about_social_picture'))}")
|
||
logger.error(f" image_url_length: {len(event_data.get('about_social_picture', '') or '')}")
|
||
logger.error(f" initial_text_length: {len(initial_text) if 'initial_text' in locals() else 'N/A'}")
|
||
logger.error(f" final_text_length: {len(final_text) if 'final_text' in locals() else 'N/A'}")
|
||
return None
|
||
except Exception as e:
|
||
error_msg = str(e)
|
||
if error_msg.startswith("RetryAfter:"):
|
||
# Обработка ошибки FloodWait (429)
|
||
retry_after = int(error_msg.split(":")[1])
|
||
logger.error(f"Получена ошибка FloodWait (429) при публикации события номер {event_data['number']}")
|
||
logger.error(f"Необходимо подождать {retry_after} секунд перед следующей попыткой")
|
||
logger.error(f"Полная информация об ошибке RetryAfter для события {event_data['number']}:")
|
||
logger.error(f" Тип ошибки: RetryAfter")
|
||
logger.error(f" Сообщение: {error_msg}")
|
||
logger.error(f" retry_after: {retry_after}")
|
||
raise RetryAfterException(retry_after)
|
||
logger.error(f"Ошибка Telegram при публикации события номер {event_data['number']}: {e}")
|
||
logger.error(f"Полная информация об ошибке Telegram для события {event_data['number']}:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e) if hasattr(e, '__dict__') else 'N/A'}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
return None
|
||
except Exception as e:
|
||
logger.error(f"Неожиданная ошибка при публикации события номер {event_data['number']}: {e}")
|
||
return None
|
||
|
||
async def tg_post_event_by_id(id_event):
|
||
"""Публикация события по ID"""
|
||
logger.info(f"Запуск публикации для события ID {id_event}")
|
||
conn = None
|
||
try:
|
||
conn = pymysql.connect(
|
||
host=MDB_HOST,
|
||
user=MDB_USER,
|
||
password=MDB_PW,
|
||
database=MDBASE,
|
||
charset='utf8mb4',
|
||
cursorclass=pymysql.cursors.DictCursor,
|
||
init_command="SET time_zone='+00:00'" # Устанавливаем часовой пояс UTC
|
||
)
|
||
|
||
# Проверка времени изменения только в режиме ZILANT
|
||
with conn.cursor() as cursor:
|
||
if WORKMODE == "ZILANT":
|
||
# Вычисляем время, до которого должно быть изменено событие для публикации
|
||
modified_before = datetime.now(timezone.utc) - timedelta(minutes=EVENT_POST_DELAY)
|
||
|
||
cursor.execute("""
|
||
SELECT id_event, number, tags, unit_name, name, about, about_social_picture, tg_message_id
|
||
FROM events
|
||
WHERE id_event = %s
|
||
AND marked_to_publication = True
|
||
AND is_posted_tg = False
|
||
AND modified <= %s
|
||
""", (id_event, modified_before))
|
||
event = cursor.fetchone()
|
||
|
||
if not event:
|
||
logger.info(f"Событие ID {id_event} не найдено, уже опубликовано или изменено менее чем {EVENT_POST_DELAY} минут назад")
|
||
return
|
||
else:
|
||
# В режиме VOLK проверка времени не выполняется
|
||
cursor.execute("""
|
||
SELECT id_event, number, tags, unit_name, name, about, about_social_picture, tg_message_id
|
||
FROM events
|
||
WHERE id_event = %s
|
||
AND marked_to_publication = True
|
||
AND is_posted_tg = False
|
||
""", (id_event,))
|
||
event = cursor.fetchone()
|
||
|
||
if not event:
|
||
logger.info(f"Событие ID {id_event} не найдено или уже опубликовано")
|
||
return
|
||
|
||
# Получаем информацию о категории (площадке), если указана unit_name
|
||
unit_tg_id = None
|
||
unit_title = None
|
||
if event.get('unit_name'):
|
||
cursor.execute("""
|
||
SELECT TITLE, tg_id
|
||
FROM categories
|
||
WHERE ID = %s
|
||
""", (event['unit_name'],))
|
||
category = cursor.fetchone()
|
||
if category:
|
||
unit_title = category.get('TITLE')
|
||
unit_tg_id = category.get('tg_id')
|
||
|
||
# Добавляем информацию о категории в данные события
|
||
event['unit_tg_id'] = unit_tg_id
|
||
event['unit_title'] = unit_title
|
||
|
||
old_message_id = event.get('tg_message_id') # Сохраняем старый ID сообщения
|
||
# Создаем httpx клиент с отключенным HTTP/2
|
||
httpx_client = httpx.AsyncClient(
|
||
http2=False,
|
||
timeout=20.0,
|
||
follow_redirects=True
|
||
)
|
||
tg_message_id = None
|
||
try:
|
||
tg_message_id = await tg_post_event(httpx_client, event) # Получаем новый message_id
|
||
finally:
|
||
await httpx_client.aclose()
|
||
|
||
if tg_message_id:
|
||
# Используем UTC время для записи в базу данных
|
||
current_time_utc = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
|
||
|
||
# Обновляем базу данных в отдельном блоке try-except для гарантии сохранения
|
||
try:
|
||
# Обновляем подписки если был предыдущий message_id
|
||
if old_message_id:
|
||
cursor.execute("""
|
||
UPDATE marks_evt
|
||
SET tg_event_id = %s
|
||
WHERE tg_event_id = %s
|
||
""", (tg_message_id, old_message_id))
|
||
logger.info(f"Обновлены подписки для события {id_event}: {cursor.rowcount} записей")
|
||
|
||
# Обновляем данные о публикации события
|
||
cursor.execute("""
|
||
UPDATE events
|
||
SET is_posted_tg = True,
|
||
tg_message_id = %s,
|
||
tg_posted_date = %s
|
||
WHERE id_event = %s
|
||
""", (tg_message_id, current_time_utc, id_event))
|
||
conn.commit()
|
||
logger.info(f"Событие номер {event['number']} успешно опубликовано. Данные сохранены в БД: message_id={tg_message_id}, date={current_time_utc}")
|
||
except Exception as db_error:
|
||
logger.error(f"Ошибка при сохранении данных о публикации события {id_event} в БД: {db_error}")
|
||
# Пытаемся выполнить commit даже при ошибке
|
||
try:
|
||
conn.rollback()
|
||
# Повторная попытка обновления
|
||
cursor.execute("""
|
||
UPDATE events
|
||
SET is_posted_tg = True,
|
||
tg_message_id = %s,
|
||
tg_posted_date = %s
|
||
WHERE id_event = %s
|
||
""", (tg_message_id, current_time_utc, id_event))
|
||
conn.commit()
|
||
logger.info(f"Данные о публикации события {id_event} успешно сохранены после повторной попытки")
|
||
except Exception as retry_error:
|
||
logger.error(f"Критическая ошибка при сохранении данных о публикации события {id_event}: {retry_error}")
|
||
raise
|
||
else:
|
||
logger.error(f"Не удалось опубликовать событие ID {id_event}")
|
||
|
||
except RetryAfterException as e:
|
||
# Обработка RetryAfter (FloodWait)
|
||
logger.error(f"Полная информация об ошибке RetryAfter в tg_post_event_by_id для события ID {id_event}:")
|
||
logger.error(f" Тип ошибки: RetryAfter")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" retry_after: {e.retry_after}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
raise
|
||
except Exception as e:
|
||
# Проверяем, не RetryAfter ли это в строковом виде
|
||
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"Полная информация об ошибке RetryAfter в tg_post_event_by_id для события ID {id_event}:")
|
||
logger.error(f" Тип ошибки: RetryAfter")
|
||
logger.error(f" Сообщение: {error_msg}")
|
||
logger.error(f" retry_after: {retry_after}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
raise RetryAfterException(retry_after)
|
||
# Если это не RetryAfter, обрабатываем как обычную ошибку
|
||
logger.error(f"Ошибка при публикации: {e}")
|
||
logger.error(f"Полная информация об ошибке в tg_post_event_by_id для события ID {id_event}:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e) if hasattr(e, '__dict__') else 'N/A'}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
finally:
|
||
if conn:
|
||
conn.close()
|
||
|
||
async def tg_post_all_events():
|
||
"""Публикация всех неопубликованных событий (максимум 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,
|
||
init_command="SET time_zone='+00:00'" # Устанавливаем часовой пояс UTC
|
||
)
|
||
|
||
with conn.cursor() as cursor:
|
||
# Проверка времени изменения только в режиме ZILANT
|
||
if WORKMODE == "ZILANT":
|
||
# Вычисляем время, до которого должно быть изменено событие для публикации
|
||
modified_before = datetime.now(timezone.utc) - timedelta(minutes=EVENT_POST_DELAY)
|
||
|
||
cursor.execute("""
|
||
SELECT id_event
|
||
FROM events
|
||
WHERE marked_to_publication = True
|
||
AND is_posted_tg = False
|
||
AND modified <= %s
|
||
LIMIT 5
|
||
""", (modified_before,))
|
||
else:
|
||
# В режиме VOLK проверка времени не выполняется
|
||
cursor.execute("""
|
||
SELECT id_event
|
||
FROM events
|
||
WHERE marked_to_publication = True
|
||
AND is_posted_tg = False
|
||
LIMIT 5
|
||
""")
|
||
events = cursor.fetchall()
|
||
|
||
event_count = 0
|
||
for event in events:
|
||
try:
|
||
await tg_post_event_by_id(event['id_event'])
|
||
event_count += 1
|
||
|
||
# Задержка между публикациями разных событий
|
||
if event_count < len(events):
|
||
await asyncio.sleep(2)
|
||
|
||
except RetryAfterException as e:
|
||
logger.error(f"Прерываем публикацию из-за ошибки FloodWait. Ожидание: {e.retry_after} секунд")
|
||
logger.error(f"Полная информация об ошибке RetryAfter в tg_post_all_events для события ID {event['id_event']}:")
|
||
logger.error(f" Тип ошибки: RetryAfter")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" retry_after: {e.retry_after}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
break
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при публикации события ID {event['id_event']}: {e}")
|
||
logger.error(f"Полная информация об ошибке в tg_post_all_events для события ID {event['id_event']}:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e) if hasattr(e, '__dict__') else 'N/A'}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
# Продолжаем публикацию следующих событий, несмотря на ошибку
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при публикации события ID {event['id_event']}: {e}")
|
||
logger.error(f"Полная информация об ошибке в tg_post_all_events для события ID {event['id_event']}:")
|
||
logger.error(f" Тип ошибки: {type(e).__name__}")
|
||
logger.error(f" Сообщение: {str(e)}")
|
||
logger.error(f" Все атрибуты ошибки: {vars(e)}")
|
||
logger.error(f" Полное представление: {repr(e)}")
|
||
# Продолжаем публикацию следующих событий, несмотря на ошибку
|
||
|
||
logger.info(f"Опубликовано событий в этом запуске: {event_count}")
|
||
|
||
except Exception as e:
|
||
logger.error(f"Ошибка при публикации: {e}")
|
||
finally:
|
||
if conn:
|
||
conn.close()
|
||
logger.info("Завершение работы скрипта")
|
||
|
||
if __name__ == "__main__":
|
||
asyncio.run(tg_post_all_events())
|