From cdb85f61c76ac272a0bc49e2a74bd4a0ba768b51 Mon Sep 17 00:00:00 2001 From: gitadmin Date: Thu, 25 Sep 2025 20:43:09 +0300 Subject: [PATCH] Upload files to "/" --- runner.sh | 113 +++++ tg_mainbot.py | 1206 +++++++++++++++++++++++++++++++++++++++++++++++++ tg_publish.py | 324 +++++++++++++ vk_load.py | 248 ++++++++++ 4 files changed, 1891 insertions(+) create mode 100644 runner.sh create mode 100644 tg_mainbot.py create mode 100644 tg_publish.py create mode 100644 vk_load.py diff --git a/runner.sh b/runner.sh new file mode 100644 index 0000000..3e84426 --- /dev/null +++ b/runner.sh @@ -0,0 +1,113 @@ +#!/bin/bash + +# Директория со скриптами +SCRIPTS_DIR="/opt/testbot" +# Активация виртуального окружения +VENV_ACTIVATE="$SCRIPTS_DIR/venv/bin/activate" + +# Список скриптов для запуска в порядке выполнения +SCRIPTS=("vk_load.py" "zk_load.py" "db_update_shortname.py" "evt_prefetch.py" "tg_publish.py" "evtg_publish.py") + +# Интервалы между скриптами (в секундах) +DELAYS=(2 2 2 2 30) + +# Лог-файл +LOG_FILE="/opt/testbot/runner.log" + +# Функция для логирования +log_message() { + local message="$1" + echo "$(date '+%Y-%m-%d %H:%M:%S') - $message" | tee -a "$LOG_FILE" +} + +# Проверка и настройка лог-файла +setup_logfile() { + # Проверяем, существует ли файл и можно ли в него писать + if [ ! -f "$LOG_FILE" ]; then + # Пытаемся создать файл + if ! touch "$LOG_FILE" 2>/dev/null; then + # Если не получается, пробуем с sudo + if sudo touch "$LOG_FILE" 2>/dev/null; then + sudo chown $(whoami) "$LOG_FILE" + sudo chmod 644 "$LOG_FILE" + log_message "Создан лог-файл с помощью sudo: $LOG_FILE" + else + echo "ОШИБКА: Не удалось создать лог-файл $LOG_FILE!" + exit 1 + fi + fi + fi + + # Проверяем права на запись + if [ ! -w "$LOG_FILE" ]; then + # Пытаемся изменить права + if sudo chown $(whoami) "$LOG_FILE" 2>/dev/null || chmod 644 "$LOG_FILE" 2>/dev/null; then + log_message "Изменены права на лог-файл: $LOG_FILE" + else + echo "ОШИБКА: Нет прав на запись в $LOG_FILE!" + exit 1 + fi + fi +} + +# Настройка лог-файла +setup_logfile + +log_message "=== ЗАПУСК СКРИПТА ===" +log_message "Пользователь: $(whoami)" +log_message "Директория со скриптами: $SCRIPTS_DIR" + +# Проверка существования директории +if [ ! -d "$SCRIPTS_DIR" ]; then + log_message "ОШИБКА: Директория $SCRIPTS_DIR не существует!" + exit 1 +fi + +# Активация виртуального окружения +if [ -f "$VENV_ACTIVATE" ]; then + source "$VENV_ACTIVATE" + log_message "Активировано виртуальное окружение: $VENV_ACTIVATE" +else + log_message "ОШИБКА: Файл активации venv $VENV_ACTIVATE не найден!" + exit 1 +fi + +# Переход в директорию со скриптами +cd "$SCRIPTS_DIR" || { + log_message "ОШИБКА: Не удалось перейти в $SCRIPTS_DIR" + exit 1 +} + +# Запуск скриптов +for i in "${!SCRIPTS[@]}"; do + script="${SCRIPTS[i]}" + + # Проверка существования скрипта + if [ ! -f "$script" ]; then + log_message "ПРЕДУПРЕЖДЕНИЕ: Скрипт $script не найден, пропускаем." + continue + fi + + log_message "Запуск скрипта: $script" + + # Запуск Python-скрипта с записью вывода в лог + if python3 "$script" >> "$LOG_FILE" 2>&1; then + log_message "УСПЕХ: Скрипт $script завершился успешно." + else + exit_code=$? + log_message "ОШИБКА: Скрипт $script завершился с ошибкой (код: $exit_code)!" + log_message "=== ЗАВЕРШЕНИЕ С ОШИБКОЙ ===" + exit $exit_code + fi + + # Пауза после скрипта (кроме последнего) + if [ $i -lt $((${#SCRIPTS[@]} - 1)) ]; then + delay="${DELAYS[i]}" + log_message "Ожидание $delay секунд перед следующим скриптом..." + sleep $delay + fi +done + +log_message "ВСЕ СКРИПТЫ ВЫПОЛНЕНЫ УСПЕШНО!" +log_message "=== ЗАВЕРШЕНИЕ ===" + diff --git a/tg_mainbot.py b/tg_mainbot.py new file mode 100644 index 0000000..68ba9bd --- /dev/null +++ b/tg_mainbot.py @@ -0,0 +1,1206 @@ +import os +import logging +import time +import errno +from dotenv import load_dotenv +import telebot +from telebot.types import InlineKeyboardMarkup, InlineKeyboardButton +import pymysql +from datetime import datetime, timezone +import sys +from urllib.parse import urlparse +import requests +from io import BytesIO +from flask import Flask, request +from formatter import get_post_text, get_event_text + +# Создаем Flask app на верхнем уровне для экспорта +app = Flask(__name__) + +# Загружаем переменные окружения из файла .env +load_dotenv() + +# Настройка логирования +logging.basicConfig( + format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', + level=logging.INFO +) +logger = logging.getLogger(__name__) + +# Получаем данные из переменных окружения +RESPONDER_BOT_TOKEN = os.getenv('RESPONDER_BOT_TOKEN') +CHANNEL_ID = os.getenv('CHANNEL_ID') +WEBHOOK_URL = os.getenv('WEBHOOK_URL') +WEBHOOK_PORT = int(os.getenv('WEBHOOK_PORT', '8443')) +WEBHOOK_SECRET = os.getenv('WEBHOOK_SECRET') +LOG_FILE = os.getenv('LOG_FILE', 'bot.log') +MAX_CAPTION_LENGTH = int(os.getenv('MAX_CAPTION_LENGTH', 1000)) +MAX_TEXT_LENGTH = int(os.getenv('MAX_TEXT_LENGTH', 4000)) + +# Устанавливаем единое ограничение длины текста +UNIFIED_MAX_LENGTH = 800 + +# Извлекаем путь из WEBHOOK_URL +if WEBHOOK_URL: + try: + parsed_url = urlparse(WEBHOOK_URL) + WEBHOOK_PATH = parsed_url.path + if not WEBHOOK_PATH: + WEBHOOK_PATH = '/' + logger.info(f"Извлечен путь вебхука: {WEBHOOK_PATH}") + except Exception as e: + logger.error(f"Ошибка парсинга WEBHOOK_URL: {e}") + WEBHOOK_PATH = '/zil_bot' +else: + WEBHOOK_PATH = '/zil_bot' + logger.warning("WEBHOOK_URL не задан, используем путь по умолчанию") + +# Данные для подключения к MariaDB +MDB_HOST = os.getenv('MDB_HOST') +MDB_USER = os.getenv('MDB_USER') +MDB_PW = os.getenv('MDB_PW') +MDBASE = os.getenv('MDBASE') + +# Проверка обязательных переменных +required_vars = { + "RESPONDER_BOT_TOKEN": RESPONDER_BOT_TOKEN, + "CHANNEL_ID": CHANNEL_ID, + "WEBHOOK_URL": WEBHOOK_URL, + "MDB_HOST": MDB_HOST, + "MDB_USER": MDB_USER, + "MDB_PW": MDB_PW, + "MDBASE": MDBASE +} + +for name, value in required_vars.items(): + if not value: + logger.error(f"Требуется переменная окружения {name}") + exit(1) + +# Создаем экземпляр бота +bot = telebot.TeleBot(RESPONDER_BOT_TOKEN) + +def get_user_name(user): + """Форматирует имя пользователя для записи в БД""" + if user.username: + return f"@{user.username}" + else: + # Используем имя и фамилию, если username отсутствует + name_parts = [] + if user.first_name: + name_parts.append(user.first_name) + if user.last_name: + name_parts.append(user.last_name) + return ' '.join(name_parts) if name_parts else "Неизвестный пользователь" + +def format_user_info(user_id, user_name): + """Форматирует информацию о пользователе для логов""" + return f"{user_name} ({user_id})" + +def log_event(event_type, user_info, post_id=None, message=None): + """Логирует событие с повторными попытками при блокировке файла""" + timestamp = datetime.now(timezone.utc).strftime('[%Y-%m-%d %H:%M:%S]') + log_line = f"{timestamp} [TG_bot] " + + if event_type == "subscribe": + log_line += f"Пользователь {user_info} подписался на событие ID {post_id}" + elif event_type == "unsubscribe": + log_line += f"Пользователь {user_info} отписался от событие ID {post_id}" + elif event_type == "list_request": + log_line += f"Пользователь {user_info} запросил список подписок" + elif event_type == "error": + log_line += f"ОШИБКА: {message} (Пользователь {user_info})" + elif event_type == "security": + log_line += f"SECURITY: {message}" + elif event_type == "system": + log_line += f"SYSTEM: {message}" + elif event_type == "info": + log_line += f"ИНФО: {message} (Пользователь {user_info})" + else: + log_line += f"Неизвестное событие: {event_type} (Пользователь {user_info})" + + # Параметры повторных попыток + intervals = [0.1, 0.2, 0.4, 0.8] # Интервалы между попытками + max_attempts = len(intervals) + 1 # +1 для первой попытки + attempt = 0 + + while attempt < max_attempts: + try: + with open(LOG_FILE, 'a', encoding='utf-8') as log_file: + log_file.write(log_line + '\n') + log_file.flush() # Обеспечить немедленную запись + return + except (IOError, OSError) as e: + # Проверяем, является ли ошибка блокировкой файла + if e.errno in (errno.EAGAIN, errno.EACCES, errno.EWOULDBLOCK, errno.EBUSY): + if attempt < max_attempts - 1: + time.sleep(intervals[attempt]) + attempt += 1 + continue + # Для других ошибок сразу прерываем цикл + break + except Exception as e: + # Обрабатываем все остальные исключения + break + + # Все попытки провалились + error_msg = f"Не удалось записать в лог после {max_attempts} попыток. Ошибка: {e}" + sys.stderr.write(error_msg + '\n') + sys.stderr.write(f"Событие: {log_line}\n") + +def create_db_connection(): + """Создает подключение к базе данных MariaDB через pymysql""" + 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 + ) + return conn + except pymysql.Error as e: + error_msg = f"Ошибка подключения к базе данных: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return None + +def is_user_subscribed(post_id, user_id): + """Проверяет, подписан ли пользователь уже на это событие""" + conn = create_db_connection() + if conn is None: + return False + + try: + with conn.cursor() as cursor: + sql = """ + SELECT COUNT(*) as count + FROM marks + WHERE tg_post_id = %s AND tg_user_id = %s + """ + cursor.execute(sql, (post_id, user_id)) + result = cursor.fetchone() + return result['count'] > 0 + except pymysql.Error as e: + error_msg = f"Ошибка при проверке подписки: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return False + finally: + conn.close() + +def is_user_subscribed_evt(event_id, user_id): + """Проверяет, подписан ли пользователь уже на это событие (таблица events)""" + conn = create_db_connection() + if conn is None: + return False + + try: + with conn.cursor() as cursor: + sql = """ + SELECT COUNT(*) as count + FROM marks_evt + WHERE tg_event_id = %s AND tg_user_id = %s + """ + cursor.execute(sql, (event_id, user_id)) + result = cursor.fetchone() + return result['count'] > 0 + except pymysql.Error as e: + error_msg = f"Ошибка при проверке подписки на событие: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return False + finally: + conn.close() + +def get_user_subscriptions(user_id): + """Возвращает список событий, на которые подписан пользователь, с shortname""" + conn = create_db_connection() + if conn is None: + return [] + + try: + with conn.cursor() as cursor: + sql = """ + SELECT m.tg_post_id, COALESCE(p.shortname, '') as shortname + FROM marks m + LEFT JOIN posts p ON m.tg_post_id = p.tg_message_id + WHERE m.tg_user_id = %s + ORDER BY m.created_at DESC + """ + cursor.execute(sql, (user_id,)) + result = cursor.fetchall() + + return [{'post_id': row['tg_post_id'], 'shortname': row['shortname']} for row in result] + except pymysql.Error as e: + error_msg = f"Ошибка при получении подписок: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return [] + finally: + conn.close() + +def get_user_subscriptions_evt(user_id): + """Возвращает список событий, на которые подписан пользователь, с name""" + conn = create_db_connection() + if conn is None: + return [] + + try: + with conn.cursor() as cursor: + sql = """ + SELECT m.tg_event_id, COALESCE(e.name, '') as name + FROM marks_evt m + LEFT JOIN events e ON m.tg_event_id = e.tg_message_id + WHERE m.tg_user_id = %s + ORDER BY m.created_at DESC + """ + cursor.execute(sql, (user_id,)) + result = cursor.fetchall() + + return [{'event_id': row['tg_event_id'], 'name': row['name']} for row in result] + except pymysql.Error as e: + error_msg = f"Ошибка при получении подписок на события: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return [] + finally: + conn.close() + +def save_mark_to_db(post_id, user_id, user_name): + """Сохраняет отметку пользователя в базе данных""" + conn = create_db_connection() + if conn is None: + return False + + try: + with conn.cursor() as cursor: + # Используем INSERT IGNORE для предотвращения дубликатов + sql = """ + INSERT IGNORE INTO marks (tg_post_id, tg_user_id, tg_user_name, created_at) + VALUES (%s, %s, %s, UTC_TIMESTAMP()) + """ + cursor.execute(sql, (post_id, user_id, user_name)) + conn.commit() + + if cursor.rowcount > 0: + logger.info(f"Добавлена запись: событие={post_id}, пользователь={user_id}, имя={user_name}") + user_info = format_user_info(user_id, user_name) + log_event("subscribe", user_info, post_id) + return True + else: + logger.info(f"Запись уже существует: событие={post_id}, пользователь={user_id}") + # Добавляем запись в лог о попытке повторной подписки + user_info = format_user_info(user_id, user_name) + log_event("info", user_info, post_id, + f"Попытка повторной подписки: событие={post_id}, пользователь={user_info}") + return False + except pymysql.Error as e: + error_msg = f"Ошибка при сохранении в БД: {e}" + logger.error(error_msg) + user_info = format_user_info(user_id, user_name) + log_event("error", user_info, post_id, error_msg) + return False + finally: + conn.close() + +def save_mark_to_db_evt(event_id, user_id, user_name): + """Сохраняет отметку пользователя в базе данных для событий (таблица marks_evt)""" + conn = create_db_connection() + if conn is None: + return False + + try: + with conn.cursor() as cursor: + sql = """ + INSERT IGNORE INTO marks_evt (tg_event_id, tg_user_id, tg_user_name, created_at) + VALUES (%s, %s, %s, UTC_TIMESTAMP()) + """ + cursor.execute(sql, (event_id, user_id, user_name)) + conn.commit() + + if cursor.rowcount > 0: + logger.info(f"Добавлена запись: событие={event_id}, пользователь={user_id}, имя={user_name}") + user_info = format_user_info(user_id, user_name) + log_event("subscribe", user_info, event_id) + return True + else: + logger.info(f"Запись уже существует: событие={event_id}, пользователь={user_id}") + user_info = format_user_info(user_id, user_name) + log_event("info", user_info, event_id, + f"Попытка повторной подписки: событие={event_id}, пользователь={user_info}") + return False + except pymysql.Error as e: + error_msg = f"Ошибка при сохранении в БД (events): {e}" + logger.error(error_msg) + user_info = format_user_info(user_id, user_name) + log_event("error", user_info, event_id, error_msg) + return False + finally: + conn.close() + +def remove_mark_from_db(post_id, user_id, user_name): + """Удаляет отметку пользователя из базы данных""" + conn = create_db_connection() + if conn is None: + return False + + try: + with conn.cursor() as cursor: + sql = """ + DELETE FROM marks + WHERE tg_post_id = %s AND tg_user_id = %s + """ + cursor.execute(sql, (post_id, user_id)) + conn.commit() + + if cursor.rowcount > 0: + logger.info(f"Удалена запись: событие={post_id}, пользователь={user_id}") + user_info = format_user_info(user_id, user_name) + log_event("unsubscribe", user_info, post_id) + return True + else: + logger.info(f"Запись не найдена: событие={post_id}, пользователь={user_id}") + # Добавляем запись в лог о попытке отписки без подписки + user_info = format_user_info(user_id, user_name) + log_event("info", user_info, post_id, + f"Попытка отписки без подписки: событие={post_id}, пользователь={user_info}") + return False + except pymysql.Error as e: + error_msg = f"Ошибка при удалении из БД: {e}" + logger.error(error_msg) + user_info = format_user_info(user_id, user_name) + log_event("error", user_info, post_id, error_msg) + return False + finally: + conn.close() + +def remove_mark_from_db_evt(event_id, user_id, user_name): + """Удаляет отметку пользователя из базы данных для событий (таблица marks_evt)""" + conn = create_db_connection() + if conn is None: + return False + + try: + with conn.cursor() as cursor: + sql = """ + DELETE FROM marks_evt + WHERE tg_event_id = %s AND tg_user_id = %s + """ + cursor.execute(sql, (event_id, user_id)) + conn.commit() + + if cursor.rowcount > 0: + logger.info(f"Удалена запись: событие={event_id}, пользователь={user_id}") + user_info = format_user_info(user_id, user_name) + log_event("unsubscribe", user_info, event_id) + return True + else: + logger.info(f"Запись не найдена: событие={event_id}, пользователь={user_id}") + user_info = format_user_info(user_id, user_name) + log_event("info", user_info, event_id, + f"Попытка отписки без подписки: событие={event_id}, пользователь={user_info}") + return False + except pymysql.Error as e: + error_msg = f"Ошибка при удалении из БД (events): {e}" + logger.error(error_msg) + user_info = format_user_info(user_id, user_name) + log_event("error", user_info, event_id, error_msg) + return False + finally: + conn.close() + +def format_channel_link(post_id=None): + """Форматирует ссылку на событие в канале или на сам канал""" + if not CHANNEL_ID: + return "" + + if post_id: + if CHANNEL_ID.startswith('@'): + return f"https://t.me/{CHANNEL_ID[1:]}/{post_id}" + elif CHANNEL_ID.startswith('-100'): + return f"https://t.me/c/{CHANNEL_ID[4:]}/{post_id}" + else: + return f"https://t.me/c/{CHANNEL_ID}/{post_id}" + else: + if CHANNEL_ID.startswith('@'): + return f"https://t.me/{CHANNEL_ID[1:]}" + elif CHANNEL_ID.startswith('-100'): + return f"https://t.me/c/{CHANNEL_ID[4:]}" + else: + return f"https://t.me/c/{CHANNEL_ID}" + +def create_help_keyboard(): + """Создает клавиатуру для справки""" + keyboard = InlineKeyboardMarkup() + keyboard.row( + InlineKeyboardButton("📜 Подписки", callback_data="cmd_list"), + InlineKeyboardButton("📢 Канал", url=format_channel_link()) + ) + return keyboard + +def create_main_keyboard(): + """Создает основную клавиатуру для /start без параметров""" + keyboard = InlineKeyboardMarkup() + keyboard.row( + InlineKeyboardButton("❓ Справка", callback_data="cmd_help"), + InlineKeyboardButton("📢 Канал", url=format_channel_link()) + ) + return keyboard + +def create_manage_keyboard(post_id, is_subscribed): + """Создает клавиатуру для управления подпиской на посты""" + keyboard = InlineKeyboardMarkup() + + channel_link = format_channel_link(post_id) + if is_subscribed: + keyboard.row( + InlineKeyboardButton("❓ Справка", callback_data="cmd_help"), + InlineKeyboardButton("❌ Отписаться", callback_data=f"unsubscribe_{post_id}"), + InlineKeyboardButton("↪ Назад", url=channel_link) + ) + else: + keyboard.row( + InlineKeyboardButton("❓ Справка", callback_data="cmd_help"), + InlineKeyboardButton("✔ Подписаться", callback_data=f"subscribe_{post_id}"), + InlineKeyboardButton("↪ Назад", url=channel_link) + ) + return keyboard + +def create_manage_keyboard_evt(event_id, is_subscribed): + """Создает клавиатуру для управления подпиской на события""" + keyboard = InlineKeyboardMarkup() + + channel_link = format_channel_link(event_id) + if is_subscribed: + keyboard.row( + InlineKeyboardButton("❓ Справка", callback_data="cmd_help"), + InlineKeyboardButton("❌ Отписаться", callback_data=f"unsubscribe_evt_{event_id}"), + InlineKeyboardButton("↪ Назад", url=channel_link) + ) + else: + keyboard.row( + InlineKeyboardButton("❓ Справка", callback_data="cmd_help"), + InlineKeyboardButton("✔ Подписаться", callback_data=f"subscribe_evt_{event_id}"), + InlineKeyboardButton("↪ Назад", url=channel_link) + ) + + return keyboard + +def get_post_data(tg_message_id): + """Получает данные о событии по его tg_message_id""" + conn = create_db_connection() + if conn is None: + return None + + try: + with conn.cursor() as cursor: + sql = """ + SELECT + vk_post_id, + COALESCE(text, '') AS text, + COALESCE(image_url, '') AS image_url, + COALESCE(vk_post_url, '') AS vk_post_url, + COALESCE(poll_question, '') AS poll_question, + COALESCE(poll_options, '[]') AS poll_options, + poll_multiple, + is_poll, + COALESCE(shortname, '') AS shortname + FROM posts + WHERE tg_message_id = %s + """ + cursor.execute(sql, (tg_message_id,)) + return cursor.fetchone() + except pymysql.Error as e: + error_msg = f"Ошибка при получении данных о событии: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return None + finally: + conn.close() + +def get_event_data(tg_message_id): + """Получает данные о событии из таблица events по его tg_message_id""" + conn = create_db_connection() + if conn is None: + return None + + try: + with conn.cursor() as cursor: + sql = """ + SELECT + about_social_picture, + COALESCE(name, '') AS name, + COALESCE(about, '') AS about + FROM events + WHERE tg_message_id = %s + """ + cursor.execute(sql, (tg_message_id,)) + return cursor.fetchone() + except pymysql.Error as e: + error_msg = f"Ошибка при получении данных о событии (events): {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return None + finally: + conn.close() + +def format_post_message(tg_message_id): + """Форматирует сообщение о событии с помощью модуля formatter""" + post_data = get_post_data(tg_message_id) + if not post_data: + error_msg = f"Событие с tg_message_id={tg_message_id} не найдено в базе." + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return { + 'text': "❌ Событие не найдено в базе данных. Пожалуйста, сообщите администратору.", + 'image_url': None + } + + try: + # Подготавливаем данные для formatter + vk_post_id = post_data['vk_post_id'] or 0 + text = post_data['text'] or "" + image_url = post_data['image_url'] or "" + vk_post_url = post_data['vk_post_url'] or "" + poll_question = post_data['poll_question'] or "" + poll_options = post_data['poll_options'] or "[]" + poll_multiple = bool(post_data['poll_multiple']) if post_data['poll_multiple'] is not None else False + + post_tuple = ( + vk_post_id, + text, + image_url, + vk_post_url, + poll_question, + poll_options, + poll_multiple, + False, # is_event + False # Добавленный элемент со значением False + ) + + # Форматируем событие с единым ограничением длины + formatted = get_post_text( + post_tuple, + max_caption_length=UNIFIED_MAX_LENGTH, + max_text_length=UNIFIED_MAX_LENGTH + ) + + # Используем caption для сообщений с изображениями, text_message для текстовых + if image_url: + base_text = formatted.get('caption') or formatted.get('base_text') or "❌ Не удалось подготовить текст о событии." + else: + base_text = formatted.get('text_message') or formatted.get('base_text') or "❌ Не удалось подготовить текст о событии." + + # Дополнительная обрезка на случай, если форматтер не обрезал текст + if len(base_text) > UNIFIED_MAX_LENGTH: + base_text = base_text[:UNIFIED_MAX_LENGTH - 3] + "..." + + return { + 'text': base_text, + 'image_url': image_url if image_url else None + } + + except Exception as e: + error_msg = f"Ошибка при форматировании текста о событии: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return { + 'text': "❌ Произошла ошибка при форматировании текста о событии. Пожалуйста, попробуйте позже.", + 'image_url': None + } + +def format_event_message(tg_message_id): + """Форматирует сообщение о событии из таблицы events""" + event_data = get_event_data(tg_message_id) + if not event_data: + error_msg = f"Событие с tg_message_id={tg_message_id} не найдено в базе events." + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return { + 'text': "❌ Событие не найдено в базе данных. Пожалуйста, сообщите администратору.", + 'image_url': None + } + + try: + name = event_data['name'] or "" + about = event_data['about'] or "" + image_url = event_data['about_social_picture'] or "" + + # Форматируем текст события + formatted_about = get_event_text(about, UNIFIED_MAX_LENGTH) + + # Формируем итоговый текст + if name and formatted_about: + text = f"{name}\n{formatted_about}" + elif name: + text = f"{name}" + elif formatted_about: + text = formatted_about + else: + text = "❌ Не удалось подготовить текст о событии." + + # Дополнительная обрезка на случай, если текст слишком длинный + if len(text) > UNIFIED_MAX_LENGTH: + text = text[:UNIFIED_MAX_LENGTH - 3] + "..." + + return { + 'text': text, + 'image_url': image_url if image_url else None + } + + except Exception as e: + error_msg = f"Ошибка при форматировании текста о событии (events): {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return { + 'text': "❌ Произошла ошибка при форматировании текста о событии. Пожалуйста, попробуйте позже.", + 'image_url': None + } + +def download_image(image_url): + """Загружает изображение по URL""" + try: + response = requests.get(image_url, timeout=10) + response.raise_for_status() + return BytesIO(response.content) + except Exception as e: + logger.error(f"Ошибка загрузки изображения: {e}") + return None + +def send_post_message(chat_id, text, image_url=None, reply_markup=None): + """Отправляет сообщение с изображением или текстом""" + try: + # Дополнительная обрезка текста + if len(text) > UNIFIED_MAX_LENGTH: + text = text[:UNIFIED_MAX_LENGTH - 3] + "..." + + if image_url: + # Загружаем изображение + image = download_image(image_url) + if image: + # Отправляем фото с подписью + bot.send_photo( + chat_id, + photo=image, + caption=text, + reply_markup=reply_markup, + parse_mode='HTML' + ) + return True + else: + logger.warning(f"Не удалось загрузить изображение: {image_url}") + + # Если изображение отсутствует или не загружено, отправляем текст + bot.send_message( + chat_id, + text, + reply_markup=reply_markup, + parse_mode='HTML' + ) + return True + except Exception as e: + logger.error(f"Ошибка отправки сообщения: {e}") + return False + +def send_management_message(chat_id, post_id, user_id, user_name): + """Отправляет сообщение для управления подпиской на посты""" + try: + # Проверяем статус подписки + is_subscribed = is_user_subscribed(post_id, user_id) + + # Форматируем текст о событии + post_data = format_post_message(post_id) + + # Создаем клавиатуру управления + keyboard = create_manage_keyboard(post_id, is_subscribed) + + # Отправляем сообщение + send_post_message( + chat_id, + post_data['text'], + image_url=post_data['image_url'], + reply_markup=keyboard + ) + except Exception as e: + error_msg = f"Ошибка отправки сообщения управления: {e}" + logger.error(error_msg) + user_info = format_user_info(user_id, user_name) + log_event("error", user_info, post_id, error_msg) + +def send_management_message_evt(chat_id, event_id, user_id, user_name): + """Отправляет сообщение для управления подпиской на события""" + try: + # Проверяем статус подписки + is_subscribed = is_user_subscribed_evt(event_id, user_id) + + # Форматируем текст о событии + event_data = format_event_message(event_id) + + # Создаем клавиатуру управления + keyboard = create_manage_keyboard_evt(event_id, is_subscribed) + + # Отправляем сообщение + send_post_message( + chat_id, + event_data['text'], + image_url=event_data['image_url'], + reply_markup=keyboard + ) + except Exception as e: + error_msg = f"Ошибка отправки сообщения управления (events): {e}" + logger.error(error_msg) + user_info = format_user_info(user_id, user_name) + log_event("error", user_info, event_id, error_msg) + +@bot.message_handler(commands=['start']) +def handle_start(message): + """Обработчик команды /start с параметром""" + try: + # Извлекаем аргументы из команды /start + args = message.text.split() + user = message.from_user + user_id = user.id + user_name = get_user_name(user) + + if len(args) > 1: + # Обработка команды подписки на посты + if args[1].startswith('post_'): + post_id = args[1].split('_')[1] + # Отправляем сообщение управления подпиской + send_management_message(message.chat.id, post_id, user_id, user_name) + return + # Обработка команды подписки на события + elif args[1].startswith('event_'): + event_id = args[1].split('_')[1] + # Отправляем сообщение управления подпиской + send_management_message_evt(message.chat.id, event_id, user_id, user_name) + return + + # Команда /start без параметров + welcome_text = ( + "🌟 Привет! 🌟\n\n" + "Я - информационный бот Зиланткона. " + "Здесь можно подписаться на интересные события будущего Зиланта - и я напомню," + "когда их можно будет включить в свой План Захвата Конвента.\n\n" + ) + + bot.send_message( + message.chat.id, + welcome_text, + reply_markup=create_main_keyboard(), + parse_mode='HTML' + ) + except Exception as e: + error_msg = f"Ошибка в обработке /start: {e}" + logger.error(error_msg) + user_info = format_user_info(message.from_user.id, get_user_name(message.from_user)) if message.from_user else "unknown" + log_event("error", user_info, message=error_msg) + +@bot.message_handler(commands=['home']) +def handle_home(message): + """Обработчик команды /home - переход в основной канал""" + try: + channel_link = format_channel_link() + bot.send_message( + message.chat.id, + "📢 Основной канал Зиланткона\n\n", + reply_markup=InlineKeyboardMarkup().row( + InlineKeyboardButton("📢 Перейти в канал", url=channel_link) + ), + parse_mode='HTML' + ) + except Exception as e: + error_msg = f"Ошибка в обработке /home: {e}" + logger.error(error_msg) + user_info = format_user_info(message.from_user.id, get_user_name(message.from_user)) if message.from_user else "unknown" + log_event("error", user_info, message=error_msg) + +@bot.message_handler(commands=['help']) +def handle_help(message): + """Обработчик команды /help""" + try: + help_text = ( + "ℹ️ Подсказка информ-бота Зиланткона\n\n" + "Я запоминаю, кто чем интересовался, и напоминаю, " + "когда это становится можно включить в План Захвата.\n\n" + "Доступные команды:\n" + "• /start - начать работу с ботом (вы уже здесь!)\n" + "• /help - показать эту справку\n" + "• /home - перейти в основной канал Зиланткона\n" + "• /list - показать текущие подписки\n\n" + "Используйте кнопки ниже для быстрого доступа:" + ) + + bot.send_message( + message.chat.id, + help_text, + reply_markup=create_help_keyboard(), + parse_mode='HTML' + ) + except Exception as e: + error_msg = f"Ошибка в обработке /help: {e}" + logger.error(error_msg) + user_info = format_user_info(message.from_user.id, get_user_name(message.from_user)) if message.from_user else "unknown" + log_event("error", user_info, message=error_msg) + +@bot.message_handler(commands=['list']) +def handle_list(message): + """Обработчик команды /list""" + try: + user_id = message.from_user.id + user_name = get_user_name(message.from_user) + + # Получаем подписки на посты и события + post_subscriptions = get_user_subscriptions(user_id) + event_subscriptions = get_user_subscriptions_evt(user_id) + + # Формируем ответ + response = "📋 Ваши подписки\n\n" + + # Добавляем подписки на посты + if post_subscriptions: + response += "📝 Анонсы:\n" + for sub in post_subscriptions: + post_id = sub['post_id'] + shortname = sub['shortname'] + post_link = format_channel_link(post_id) + + # Используем shortname если есть, иначе ID события + if shortname: + display_name = shortname + else: + # Получаем vk_post_id для отображения + post_data = get_post_data(post_id) + vk_post_id = post_data['vk_post_id'] if post_data else "N/A" + display_name = f"Анонс ВК {vk_post_id}" + + response += f"• {display_name}\n" + response += "\n" + + # Добавляем подписки на события + if event_subscriptions: + response += "🎭 Мероприятия:\n" + for sub in event_subscriptions: + event_id = sub['event_id'] + name = sub['name'] + event_link = format_channel_link(event_id) + + # Используем name если есть, иначе ID события + display_name = name if name else f"Мероприятие {event_id}" + + response += f"• {display_name}\n" + + # Если нет подписок + if not post_subscriptions and not event_subscriptions: + response += "Вы пока не подписаны ни на одно событие.\n\n" + + bot.send_message( + message.chat.id, + response, + reply_markup=create_main_keyboard(), + parse_mode='HTML', + disable_web_page_preview=True + ) + + # Логируем запрос списка подписок + user_info = format_user_info(user_id, user_name) + log_event("list_request", user_info) + + except Exception as e: + error_msg = f"Ошибка в обработке /list: {e}" + logger.error(error_msg) + user_info = format_user_info(message.from_user.id, get_user_name(message.from_user)) if message.from_user else "unknown" + log_event("error", user_info, message=error_msg) + +@bot.callback_query_handler(func=lambda call: True) +def handle_callback(call): + """Обработчик нажатий на кнопки""" + try: + # Обработка командных кнопок + if call.data.startswith("cmd_"): + command = call.data[4:] + if command == "help": + help_text = ( + "ℹ️ Подсказка информ-бота Зиланткона\n\n" + "Я запоминаю, кто чем интересовался, и напоминаю, " + "когда это становится можно включить в План Захвата.\n\n" + "Доступные команды:\n" + "• /start - начать работу с ботом (вы уже здесь!)\n" + "• /help - показать эту справку\n" + "• /home - перейти в основной канал Зиланткона\n" + "• /list - показать текущие подписки\n\n" + "Используйте кнопки ниже для быстрого доступа:" + ) + + bot.send_message( + call.message.chat.id, + help_text, + reply_markup=create_help_keyboard(), + parse_mode='HTML' + ) + elif command == "list": + user = call.from_user + user_id = user.id + user_name = get_user_name(user) + + # Получаем подписки на посты и события + post_subscriptions = get_user_subscriptions(user_id) + event_subscriptions = get_user_subscriptions_evt(user_id) + + # Формируем ответ + response = "📋 Ваши подписки\n\n" + + # Добавляем подписки на посты + if post_subscriptions: + response += "📝 Анонсы:\n" + for sub in post_subscriptions: + post_id = sub['post_id'] + shortname = sub['shortname'] + post_link = format_channel_link(post_id) + + # Используем shortname если есть, иначе ID события + if shortname: + display_name = shortname + else: + # Получаем vk_post_id для отображения + post_data = get_post_data(post_id) + vk_post_id = post_data['vk_post_id'] if post_data else "N/A" + display_name = f"Анонс ВК {vk_post_id}" + + response += f"• {display_name}\n" + response += "\n" + + # Добавляем подписки на события + if event_subscriptions: + response += "🎭 Мероприятия:\n" + for sub in event_subscriptions: + event_id = sub['event_id'] + name = sub['name'] + event_link = format_channel_link(event_id) + + # Используем name если есть, иначе ID события + display_name = name if name else f"Мероприятие {event_id}" + + response += f"• {display_name}\n" + + # Если нет подписок + if not post_subscriptions and not event_subscriptions: + response += "Вы пока не подписаны ни на одно событие.\n\n" + + bot.send_message( + call.message.chat.id, + response, + reply_markup=create_main_keyboard(), + parse_mode='HTML', + disable_web_page_preview=True + ) + + # Логируем запрос списка подписок + user_info = format_user_info(user_id, user_name) + log_event("list_request", user_info) + elif command == "back": + # Возврат к главному меню + welcome_text = ( + "🌟 Привет! 🌟\n\n" + "Я - информационный бот Зиланткона. " + "Здесь можно подписаться на интересные события будущего Зиланта - и я напомню," + "когда их можно будет включить в свой План Захвата Конвента.\n\n" + ) + + bot.send_message( + call.message.chat.id, + welcome_text, + reply_markup=create_main_keyboard(), + parse_mode='HTML' + ) + + # Подтверждаем получение callback + bot.answer_callback_query(call.id) + return + + # Обработка кнопок управления подпиской на посты + if (call.data.startswith("subscribe_") or call.data.startswith("unsubscribe_")) and not call.data.startswith(("subscribe_evt_", "unsubscribe_evt_")): + user = call.from_user + user_id = user.id + user_name = get_user_name(user) + + # Разделяем данные callback + parts = call.data.split('_', 1) + action = parts[0] + post_id = parts[1] + + # Выполняем действие + if action == "subscribe": + save_mark_to_db(post_id, user_id, user_name) + result_text = "✅ Вы успешно подписались!" + else: + remove_mark_from_db(post_id, user_id, user_name) + result_text = "✅ Вы успешно отписались!" + + # Обновляем сообщение с новым статусом + try: + # Получаем текущие данные о событии + post_data = format_post_message(post_id) + text = post_data['text'] + + # Создаем новую клавиатуру + is_subscribed = action == "subscribe" + keyboard = create_manage_keyboard(post_id, is_subscribed) + + # Для сообщений с изображением + if call.message.content_type == 'photo': + # Получаем file_id существующего изображения + file_id = call.message.photo[-1].file_id + + # Редактируем подпись к изображению + bot.edit_message_caption( + chat_id=call.message.chat.id, + message_id=call.message.message_id, + caption=text, + reply_markup=keyboard, + parse_mode='HTML' + ) + else: + # Редактируем текстовое сообщение + bot.edit_message_text( + chat_id=call.message.chat.id, + message_id=call.message.message_id, + text=text, + reply_markup=keyboard, + parse_mode='HTML' + ) + + # Отправляем отдельное сообщение о результате + bot.answer_callback_query(call.id, result_text) + except Exception as e: + logger.error(f"Ошибка обновления сообщения: {e}") + bot.answer_callback_query(call.id, f"{result_text} Но не удалось обновить сообщение.") + return + + # Обработка кнопок управления подпиской на события + if call.data.startswith("subscribe_evt_") or call.data.startswith("unsubscribe_evt_"): + user = call.from_user + user_id = user.id + user_name = get_user_name(user) + parts = call.data.split('_', 2) + action = parts[0] + event_id = parts[2] + + # Выполняем действие + if action == "subscribe": + save_mark_to_db_evt(event_id, user_id, user_name) + result_text = "✅ Вы успешно подписались!" + else: + remove_mark_from_db_evt(event_id, user_id, user_name) + result_text = "✅ Вы успешно отписались!" + + # Обновляем сообщение с новым статусом + try: + # Получаем текущие данные о событии + event_data = format_event_message(event_id) + text = event_data['text'] + + # Создаем новую клавиатуру + is_subscribed = action == "subscribe" + keyboard = create_manage_keyboard_evt(event_id, is_subscribed) + + # Для сообщений с изображением + if call.message.content_type == 'photo': + # Получаем file_id существующего изображения + file_id = call.message.photo[-1].file_id + + # Редактируем подпись к изображению + bot.edit_message_caption( + chat_id=call.message.chat.id, + message_id=call.message.message_id, + caption=text, + reply_markup=keyboard, + parse_mode='HTML' + ) + else: + # Редактируем текстовое сообщение + bot.edit_message_text( + chat_id=call.message.chat.id, + message_id=call.message.message_id, + text=text, + reply_markup=keyboard, + parse_mode='HTML' + ) + + # Отправляем отдельное сообщение о результате + bot.answer_callback_query(call.id, result_text) + except Exception as e: + logger.error(f"Ошибка обновления сообщения (events): {e}") + bot.answer_callback_query(call.id, f"{result_text} Но не удалось обновить сообщение.") + return + + # Подтверждаем получение callback для любых других нажатий + bot.answer_callback_query(call.id) + + except Exception as e: + error_msg = f"Ошибка обработки callback: {e}" + logger.error(error_msg) + user_info = format_user_info(call.from_user.id, get_user_name(call.from_user)) if call.from_user else "unknown" + log_event("error", user_info, message=error_msg) + bot.answer_callback_query(call.id, "Произошла ошибка. Пожалуйста, попробуйте позже.") + +def setup_webhook(): + """Настройка вебхука""" + try: + # Удаляем предыдущий вебхук + bot.remove_webhook() + + # Устанавливаем новый вебхук + bot.set_webhook( + url=WEBHOOK_URL, + secret_token=WEBHOOK_SECRET, + max_connections=40 + ) + logger.info(f"Вебхук установлен: {WEBHOOK_URL}") + logger.info(f"Секретный токен: {'установлен' if WEBHOOK_SECRET else 'не установлен'}") + logger.info(f"Прослушивание порта: {WEBHOOK_PORT}") + logger.info(f"Путь вебхука: {WEBHOOK_PATH}") + + # Логируем успешную настройку вебхука + log_event("system", "system", message=f"Вебхук установлен на {WEBHOOK_URL}") + + except Exception as e: + error_msg = f"Ошибка настройка вебхука: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + exit(1) + +# Используем путь из WEBHOOK_URL +@app.route(WEBHOOK_PATH, methods=['POST']) +def webhook(): + if request.headers.get('X-Telegram-Bot-Api-Secret-Token') != WEBHOOK_SECRET: + logger.warning("Неверный секретный токен!") + log_event("security", "system", message="Попытка доступа с неверным секретным токеном") + return "Unauthorized", 401 + + json_data = request.get_json() + try: + update = telebot.types.Update.de_json(json_data) + bot.process_new_updates([update]) + return "OK", 200 + except Exception as e: + error_msg = f"Ошибка обработки вебхука: {e}" + logger.error(error_msg) + log_event("error", "system", message=error_msg) + return "Internal Server Error", 500 + +if __name__ == '__main__': + # Логируем запуск бота + log_event("system", "system", message="Бот запущен") + + # Настраиваем вебхук + setup_webhook() + + # Логируем запуск Flask + log_event("system", "system", message=f"Flask приложение запущено на порту {WEBHOOK_PORT}, путь: {WEBHOOK_PATH}") + + # Запускаем Flask приложение + app.run(host='0.0.0.0', port=WEBHOOK_PORT) diff --git a/tg_publish.py b/tg_publish.py new file mode 100644 index 0000000..8b724c7 --- /dev/null +++ b/tg_publish.py @@ -0,0 +1,324 @@ +import os +import logging +import asyncio +import pymysql +import time +import html +import json +import re +from datetime import datetime, timezone +from telegram import Bot, PollOption +from telegram.error import TelegramError, RetryAfter +from dotenv import load_dotenv +from formatter import get_event_text + +# Загрузка переменных окружения +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)) + +# Настройка логирования +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) + +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 + + bot = Bot(token=BOT_TOKEN) + + # Проверяем наличие организационного хештега + 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'Оригинал в ВК') + + # Ссылка на подписку (только для событий) + if post['is_event']: + links.append(f'🔜 Подписка') + + # Формируем строку ссылок + 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 has_image and not is_org_message: + try: + message = await bot.send_photo( + chat_id=CHANNEL_ID, + photo=image_url, + caption=text, + parse_mode="HTML", + disable_notification=True + ) + 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=True, + disable_web_page_preview=True + ) + else: + # Если это организационное сообщение или нет изображения, отправляем текстовое сообщение + message = await bot.send_message( + chat_id=CHANNEL_ID, + text=text, + parse_mode="HTML", + disable_notification=True, + 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'Оригинал в ВК') + + if post['is_event']: + updated_links.append(f'🔜 Подписка') + + links_line = " | ".join(updated_links) + final_text = f"{text}\n\n{links_line}" + + await asyncio.sleep(2) # задержка между обращениями к телеграм + + # Редактируем сообщение, добавляя ссылки + if has_image and not is_org_message: + try: + await bot.edit_message_caption( + chat_id=CHANNEL_ID, + message_id=message_id, + caption=final_text, + parse_mode="HTML" + ) + except TelegramError as e: + logger.warning(f"Не удалось отредактировать подпись для записи VK ID {vk_post_id}: {e}. Пробуем отредактировать текстовое сообщение.") + await bot.edit_message_text( + chat_id=CHANNEL_ID, + message_id=message_id, + text=final_text, + parse_mode="HTML", + disable_web_page_preview=True + ) + else: + await bot.edit_message_text( + chat_id=CHANNEL_ID, + message_id=message_id, + text=final_text, + parse_mode="HTML", + disable_web_page_preview=True + ) + + 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']) + 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=True + ) + + 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"Необходимо подождать {e.retry_after} секунд перед следующей попыткой") + raise + except TelegramError as e: + 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 RetryAfter 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()) diff --git a/vk_load.py b/vk_load.py new file mode 100644 index 0000000..623e92d --- /dev/null +++ b/vk_load.py @@ -0,0 +1,248 @@ +import mysql.connector +import requests +import json +import os +import time +import random +from datetime import datetime +from mysql.connector import errorcode +from dotenv import load_dotenv +from ensure_db import ensure_database_structure + +# Загрузка переменных окружения из .env файла +load_dotenv() + +# Получение конфигурации из переменных окружения +VK_GROUP_ID = int(os.getenv('VK_GROUP_ID', -3342146)) +VK_ACCESS_TOKEN = os.getenv('VK_ACCESS_TOKEN', '') +VK_API_VERSION = os.getenv('VK_API_VERSION', '5.131') +IS_EVENT = os.getenv('IS_EVENT', 'false').lower() in ('true', '1', 'yes', 'on') +AUTO_PUBLISH = os.getenv('AUTO_PUBLISH', 'false').lower() in ('true', '1', 'yes', 'on') +LOG_FILE = os.getenv('LOG_FILE', 'vk_loader.log') # Путь к лог-файлу +LOG_PREFIX = "VK_loader" # Уникальный префикс для идентификации скрипта + +# Конфигурация MariaDB из .env +DB_CONFIG = { + 'host': os.getenv('MDB_HOST', 'localhost'), + 'user': os.getenv('MDB_USER', 'root'), + 'password': os.getenv('MDB_PW', ''), + 'database': os.getenv('MDBASE', 'vk_posts') +} + +def log_message(message, max_retries=5, retry_delay=0.1): + """ + Записывает сообщение в лог-файл с обработкой блокировок + и идентификатором скрипта + """ + if not LOG_FILE: + return + + timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S') + # Добавляем идентификатор скрипта к сообщению + log_line = f"[{timestamp}] [{LOG_PREFIX}] {message}\n" + + for attempt in range(max_retries): + try: + with open(LOG_FILE, 'a', encoding='utf-8') as log: + log.write(log_line) + return True + except (IOError, OSError) as e: + if "locked" in str(e).lower() and attempt < max_retries - 1: + # Случайная задержка для уменьшения коллизий + sleep_time = retry_delay * (1 + random.random() * 0.5) + time.sleep(sleep_time) + else: + # Если не удалось записать после всех попыток + print(f"Ошибка записи в лог: {e}") + print(f"Сообщение для лога: {log_line.strip()}") + return False + +def get_vk_posts(): + """Получает последние 10 постов из VK группы""" + url = 'https://api.vk.ru/method/wall.get' + params = { + 'owner_id': VK_GROUP_ID, + 'count': 10, + 'access_token': VK_ACCESS_TOKEN, + 'v': VK_API_VERSION, + 'extended': 1 + } + + try: + response = requests.get(url, params=params, timeout=10) + response.raise_for_status() + data = response.json() + + if 'error' in data: + error_msg = f"VK API Error: {data['error']['error_msg']}" + log_message(error_msg) + raise Exception(error_msg) + + return data['response']['items'] + + except requests.exceptions.RequestException as e: + error_msg = f"Ошибка подключения к VK API: {e}" + log_message(error_msg) + raise Exception(error_msg) + +def process_post(post): + """Извлекает необходимые данные из поста VK""" + post_id = post['id'] + vk_post_url = f"https://vk.ru/wall{VK_GROUP_ID}_{post_id}" + published_at = datetime.utcfromtimestamp(post['date']).strftime('%Y-%m-%d %H:%M:%S') + + # Инициализация полей опроса + is_poll = 0 + poll_question = None + poll_options = None + poll_multiple = False + poll_end_date = None + + # Инициализация поля изображения + image_url = None + + # Обработка вложений + for attachment in post.get('attachments', []): + if attachment['type'] == 'poll': + is_poll = 1 + poll = attachment['poll'] + poll_question = poll['question'] + poll_options = json.dumps([answer['text'] for answer in poll['answers']]) + poll_multiple = bool(poll.get('multiple', False)) + + if poll.get('end_date', 0) > 0: + poll_end_date = datetime.utcfromtimestamp(poll['end_date']).strftime('%Y-%m-%d %H:%M:%S') + + elif attachment['type'] == 'photo' and not image_url: + # Берем только первое изображение + photo = attachment['photo'] + sizes = photo.get('sizes', []) + if sizes: + # Выбираем изображение максимального качества + max_size = max(sizes, key=lambda s: s.get('width', 0) * s.get('height', 0)) + image_url = max_size['url'] + + return { + 'vk_post_id': post_id, + 'text': post['text'], + 'image_url': image_url, + 'vk_post_url': vk_post_url, + 'published_at': published_at, + 'is_poll': is_poll, + 'poll_question': poll_question, + 'poll_options': poll_options, + 'poll_multiple': poll_multiple, + 'poll_end_date': poll_end_date + } + +def save_to_database(posts): + """Сохраняет новые посты в MariaDB с учетом настроек публикации""" + try: + conn = mysql.connector.connect(**DB_CONFIG) + cursor = conn.cursor() + + new_posts_count = 0 + for post in posts: + try: + # Добавляем поля is_event, marked_for_publication, shortname и action_number + cursor.execute(''' + INSERT INTO posts ( + vk_post_id, text, image_url, vk_post_url, published_at, + is_poll, poll_question, poll_options, poll_multiple, poll_end_date, + is_event, marked_for_publication, shortname, action_number + ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) + ON DUPLICATE KEY UPDATE vk_post_id = VALUES(vk_post_id) + ''', ( + post['vk_post_id'], + post['text'], + post['image_url'], + post['vk_post_url'], + post['published_at'], + post['is_poll'], + post['poll_question'], + post['poll_options'], + int(post['poll_multiple']), + post['poll_end_date'], + IS_EVENT, + int(AUTO_PUBLISH), # Значение из .env + None, # shortname - пока не используется, устанавливаем NULL + None # action_number - пока не используется, устанавливаем NULL + )) + if cursor.rowcount > 0: + new_posts_count += 1 + # Логируем добавление нового поста + log_message(f"Добавлен новый пост: ID {post['vk_post_id']}, Дата: {post['published_at']}, Автопубликация: {'Да' if AUTO_PUBLISH else 'Нет'}") + except mysql.connector.Error as err: + error_msg = f"Ошибка при вставке поста {post['vk_post_id']}: {err}" + log_message(error_msg) + continue + + conn.commit() + return new_posts_count + + except mysql.connector.Error as err: + error_msg = f"Ошибка подключения к базе данных: {err}" + log_message(error_msg) + return 0 + finally: + if 'conn' in locals() and conn.is_connected(): + cursor.close() + conn.close() + +def vk_load_10(): + """ + Экспортируемая функция для загрузки 10 последних постов из VK + Возвращает количество добавленных постов + """ + # Стартовая информация + start_msg = f"Старт скрипта | База: {DB_CONFIG['database']}@{DB_CONFIG['host']} | IS_EVENT: {IS_EVENT} | AUTO_PUBLISH: {AUTO_PUBLISH}" + print(start_msg) + log_message(start_msg) + + # Обеспечиваем структуру БД + try: + ensure_database_structure() + except Exception as e: + error_msg = f"Не удалось создать структуру базы данных: {e}" + print(error_msg) + log_message(error_msg) + return 0 + + # Получаем посты из VK + try: + vk_posts = get_vk_posts() + status_msg = f"Получено {len(vk_posts)} постов из VK" + print(status_msg) + log_message(status_msg) + except Exception as e: + error_msg = f"Ошибка при получении данных из VK: {e}" + print(error_msg) + log_message(error_msg) + return 0 + + # Обрабатываем и сохраняем посты + processed_posts = [process_post(post) for post in vk_posts] + new_count = save_to_database(processed_posts) + + # Фиксируем результат + if new_count > 0: + result_msg = f"Добавлено {new_count} новых постов в базу данных" + else: + result_msg = "Новые посты не обнаружены" + + print(result_msg) + log_message(result_msg) + + # Завершение работы + end_msg = "Работа скрипта завершена" + print(end_msg) + log_message(end_msg) + + return new_count + +def main(): + """Основная функция для запуска скрипта из командной строки""" + vk_load_10() + +if __name__ == "__main__": + main()