refactor(utils): move modules to utils package
This commit is contained in:
@@ -0,0 +1,84 @@
|
||||
import logging
|
||||
import asyncio
|
||||
|
||||
from LOLZTEAM.Client import Forum
|
||||
from typing import Dict, Any, Optional, Union, List
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Lolz:
|
||||
def __init__(self, token: str):
|
||||
try:
|
||||
self.client = Forum(token=token, timeout=15)
|
||||
self.client.settings.logger.enable()
|
||||
|
||||
except:
|
||||
raise
|
||||
|
||||
async def get_post(self, post_id: Union[str, int]) -> Dict[str, Any]:
|
||||
response = await self.client.posts.get(post_id=post_id)
|
||||
return (response.json()).get("post", {})
|
||||
|
||||
async def get_thread_posts(
|
||||
self, thread_id: Union[str, int], start_page: int = 1
|
||||
) -> List[Dict[str, Any]]:
|
||||
all_posts = []
|
||||
page = start_page
|
||||
|
||||
while True:
|
||||
response = await self.client.posts.list(thread_id=thread_id, page=page)
|
||||
posts = (response.json()).get("posts", [])
|
||||
|
||||
if len(posts) == 0:
|
||||
break
|
||||
|
||||
all_posts.extend(posts)
|
||||
logger.info(
|
||||
f"Получено {len(posts)} постов из темы {thread_id} на странице {page}."
|
||||
)
|
||||
|
||||
page += 1
|
||||
await asyncio.sleep(1)
|
||||
|
||||
if all_posts:
|
||||
logger.info(f"Всего получено {len(all_posts)} постов из темы {thread_id}.")
|
||||
else:
|
||||
logger.info(f"Постов в теме {thread_id} не найдено.")
|
||||
|
||||
return all_posts
|
||||
|
||||
async def get_all_thread_posts(
|
||||
self, thread_id: Union[str, int]
|
||||
) -> List[Dict[str, Any]]:
|
||||
return await self.get_thread_posts(thread_id=thread_id, start_page=1)
|
||||
|
||||
async def get_post_comments(self, post_id: int) -> List[Dict[str, Any]]:
|
||||
response = await self.client.posts.comments.list(post_id=post_id)
|
||||
comments = (response.json()).get("comments", [])
|
||||
|
||||
return comments
|
||||
|
||||
async def has_comments(
|
||||
self, post_id: Optional[int], post: Dict[str, Any] = None
|
||||
) -> bool:
|
||||
if not post:
|
||||
if not post_id:
|
||||
post = {}
|
||||
else:
|
||||
post = await self.get_post(post_id=post_id)
|
||||
|
||||
return post.get("post_comment_count", 0) > 0
|
||||
|
||||
async def create_comment(self, post_id: int, comment_body: str) -> bool:
|
||||
logger.info(f"Публикую комментарий к посту {post_id}...")
|
||||
response = await self.client.posts.comments.create(
|
||||
post_id=post_id, comment_body=comment_body
|
||||
)
|
||||
|
||||
if response.status_code == 200:
|
||||
logger.info(f"Комментарий к посту {post_id} успешно опубликован.")
|
||||
return True
|
||||
else:
|
||||
logger.error(f"Не удалось опубликовать комментарий к посту {post_id}.")
|
||||
return False
|
||||
@@ -0,0 +1,32 @@
|
||||
from typing import Set, List
|
||||
|
||||
|
||||
class ProcessedPostsManager:
|
||||
def __init__(self, file_path: str = "processed.txt"):
|
||||
self.file_path = file_path
|
||||
self.processed_posts: Set[int] = self._load()
|
||||
|
||||
def _load(self) -> Set[int]:
|
||||
try:
|
||||
with open(self.file_path, "r", encoding="utf-8") as processed:
|
||||
return set(processed.readlines())
|
||||
|
||||
except:
|
||||
return set()
|
||||
|
||||
def _save(self):
|
||||
with open(self.file_path, "w", encoding="utf-8") as processed:
|
||||
processed.writelines(lines=self.processed_posts)
|
||||
|
||||
def is_processed(self, post_id: int) -> bool:
|
||||
return post_id in self.processed_posts
|
||||
|
||||
def mark_processed(self, post_id: int):
|
||||
self.processed_posts.add(post_id)
|
||||
self._save()
|
||||
|
||||
def add_existing_posts(self, post_ids: List[int]):
|
||||
for post_id in post_ids:
|
||||
self.processed_posts.add(post_id)
|
||||
|
||||
self._save()
|
||||
@@ -0,0 +1,272 @@
|
||||
import re
|
||||
import os
|
||||
import random
|
||||
import logging
|
||||
import asyncio
|
||||
|
||||
from pyrogram.client import Client
|
||||
from pyrogram.errors import FloodWait
|
||||
|
||||
from typing import List, Optional, Dict, Any
|
||||
|
||||
from lolz import Lolz
|
||||
from config import Config
|
||||
from misc import ProcessedPostsManager
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class TelegramLinkExtractor:
|
||||
@staticmethod
|
||||
def extract(text: str) -> List[str]:
|
||||
patterns = [
|
||||
r"https?://(?:www\.)?(?:t\.me|telegram\.me)/([a-zA-Z0-9_]+(?:/\d+)?)",
|
||||
r"\[MEDIA=telegram\]([a-zA-Z0-9_]+(?:/\d+)?)\[/MEDIA\]",
|
||||
r'data-telegram-post="([a-zA-Z0-9_]+/\d+)"',
|
||||
]
|
||||
|
||||
all_matches = {
|
||||
f"https://t.me/{match}"
|
||||
for p in patterns
|
||||
for match in re.findall(p, text, re.I)
|
||||
}
|
||||
|
||||
return list(all_matches)
|
||||
|
||||
@staticmethod
|
||||
def parse(link: str) -> Optional[tuple[str, Optional[int]]]:
|
||||
match = re.search(r"t\.me/([^/]+)(?:/(\d+))?", link)
|
||||
|
||||
if match:
|
||||
channel = match.group(1)
|
||||
message_id = int(match.group(2)) if match.group(2) else None
|
||||
return channel, message_id
|
||||
|
||||
return None
|
||||
|
||||
|
||||
class TelegramStarsBot:
|
||||
SESSION_NAME = "stars_bot"
|
||||
|
||||
def __init__(self, config: Config):
|
||||
self.config = config
|
||||
|
||||
self.start_page: int = 1
|
||||
self.lolz_api = Lolz(config.lolz_token)
|
||||
self.processed_manager = ProcessedPostsManager(config.processed_posts_file)
|
||||
|
||||
self.client: Optional[Client] = None
|
||||
|
||||
async def parse_existing_posts(self):
|
||||
logger.info("Начинаю парсинг всех существующих постов в теме...")
|
||||
all_posts = await self.lolz_api.get_all_thread_posts(
|
||||
thread_id=self.config.forum_thread_id
|
||||
)
|
||||
|
||||
if all_posts:
|
||||
post_ids = 0
|
||||
for post in all_posts:
|
||||
post_id = post.get("post_id")
|
||||
has_comments = self.lolz_api.has_comments(post=post)
|
||||
|
||||
if has_comments:
|
||||
self.processed_manager.mark_processed(post_id)
|
||||
post_ids += 1
|
||||
|
||||
logger.info(f"Добавлено {len(post_ids)} постов в список обработанных.")
|
||||
|
||||
else:
|
||||
logger.info("Не найдено постов для добавления в обработанные.")
|
||||
|
||||
async def send_star(self, channel: str, message_id: Optional[int] = None) -> bool:
|
||||
if not hasattr(self.client, "send_paid_reaction"):
|
||||
logger.error("Платные реакции недоступны. Отправка 'звезд' невозможна.")
|
||||
return False
|
||||
|
||||
try:
|
||||
for attempt in range(self.config.max_retries):
|
||||
try:
|
||||
if message_id is None:
|
||||
history = await self.client.get_chat_history(
|
||||
chat_id=channel, limit=1
|
||||
)
|
||||
|
||||
if len(history) == 0:
|
||||
logger.error(
|
||||
f"Не удалось найти сообщения в канале https://t.me/{channel}"
|
||||
)
|
||||
return False
|
||||
|
||||
message = history[0]
|
||||
message_id = message.id
|
||||
|
||||
await self.client.send_paid_reaction(
|
||||
chat_id=channel,
|
||||
message_id=message_id,
|
||||
amount=self.config.stars_count,
|
||||
)
|
||||
logger.info(
|
||||
f"Отправлено {self.config.stars_count} звезд в https://t.me/{channel}/{message_id}"
|
||||
)
|
||||
return True
|
||||
|
||||
except FloodWait as e:
|
||||
logger.warning(f"FloodWait: необходимо подождать {e.x + 2} секунд.")
|
||||
await asyncio.sleep(e.x + 2)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"Попытка {attempt + 1} отправки звезд не удалась: {e}"
|
||||
)
|
||||
|
||||
if attempt < self.config.max_retries - 1:
|
||||
await asyncio.sleep(3 * (attempt + 1))
|
||||
|
||||
return False
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Критическая ошибка при отправке звезд: {e}")
|
||||
return False
|
||||
|
||||
async def _process_single_post(self, post: Dict[str, Any]):
|
||||
post_id = post.get("post_id")
|
||||
poster_user_id = post.get("poster_user_id")
|
||||
|
||||
if not post_id or self.processed_manager.is_processed(post_id):
|
||||
return
|
||||
|
||||
logger.info(f"Найден новый пост для обработки: ID {post_id}")
|
||||
|
||||
if self.config.skip_posts_with_comments:
|
||||
if await self.lolz_api.has_comments(post_id):
|
||||
logger.info(
|
||||
f"Пост {post_id} уже имеет комментарии. Пропускаю обработку."
|
||||
)
|
||||
self.processed_manager.mark_processed(post_id)
|
||||
return
|
||||
|
||||
post_content = post.get("post_body_html") or post.get("post_body")
|
||||
if not post_content:
|
||||
logger.warning(f"У поста {post_id} отсутствует содержимое. Пропускаем.")
|
||||
self.processed_manager.mark_processed(post_id)
|
||||
return
|
||||
|
||||
links = TelegramLinkExtractor.extract(post_content)
|
||||
if not links:
|
||||
logger.info(f"В посте {post_id} не найдено ссылок Telegram.")
|
||||
self.processed_manager.mark_processed(post_id)
|
||||
return
|
||||
|
||||
successful_reactions = 0
|
||||
for link in links:
|
||||
parsed_link = TelegramLinkExtractor.parse(link)
|
||||
if parsed_link:
|
||||
channel, message_id = parsed_link
|
||||
if await self.send_star(channel, message_id):
|
||||
successful_reactions += 1
|
||||
await asyncio.sleep(1)
|
||||
|
||||
if successful_reactions > 0 and self.config.enable_reply:
|
||||
await asyncio.sleep(self.config.api_delay)
|
||||
|
||||
# Теперь можно отправлять больше 1 сообщения
|
||||
# см. README.md
|
||||
if poster_user_id:
|
||||
replies = self.config.reply_templates
|
||||
|
||||
for reply_options in replies:
|
||||
await asyncio.sleep(self.config.api_delay)
|
||||
|
||||
reply_message = f"[userids={poster_user_id};align=left]{random.choice(reply_options)}[/userids]"
|
||||
await self.lolz_api.create_comment(post_id, reply_message)
|
||||
|
||||
logger.info(f"Пост {post_id} полностью обработан.")
|
||||
self.processed_manager.mark_processed(post_id)
|
||||
|
||||
async def _main_loop(self):
|
||||
while True:
|
||||
try:
|
||||
logger.info(
|
||||
f"Проверка новых постов в теме {self.config.forum_thread_id} начиная со страницы {self.start_page}..."
|
||||
)
|
||||
posts = await self.lolz_api.get_thread_posts(
|
||||
self.config.forum_thread_id, self.start_page
|
||||
)
|
||||
|
||||
if posts:
|
||||
for post in posts:
|
||||
await self._process_single_post(post)
|
||||
else:
|
||||
logger.info("Новых постов для обработки не найдено.")
|
||||
|
||||
logger.info(f"Ожидание {self.config.check_interval} секунд...")
|
||||
await asyncio.sleep(self.config.check_interval)
|
||||
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Получен сигнал прерывания (Ctrl+C).")
|
||||
break
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"Критическая ошибка в главном цикле: {e}")
|
||||
await asyncio.sleep(self.config.check_interval)
|
||||
|
||||
async def start(self):
|
||||
is_first_login = not os.path.exists(f"{self.SESSION_NAME}.session")
|
||||
|
||||
if is_first_login:
|
||||
logger.info("Сессия Telegram не найдена. Запускаю процесс входа...")
|
||||
else:
|
||||
choice = self.config.defualt_choice
|
||||
|
||||
if choice == "2":
|
||||
logger.info("Выбран режим парсинга существующих постов.")
|
||||
self.client = Client(
|
||||
name=self.SESSION_NAME,
|
||||
api_id=self.config.api_id,
|
||||
api_hash=self.config.api_hash,
|
||||
)
|
||||
|
||||
await self.client.start()
|
||||
await self.parse_existing_posts()
|
||||
await self.client.stop()
|
||||
|
||||
logger.info(
|
||||
"Парсинг завершен. Все существующие посты добавлены в обработанные."
|
||||
)
|
||||
|
||||
return
|
||||
|
||||
self.start_page = self.config.start_page
|
||||
|
||||
self.client = Client(
|
||||
name=self.SESSION_NAME,
|
||||
api_id=self.config.api_id,
|
||||
api_hash=self.config.api_hash,
|
||||
)
|
||||
|
||||
try:
|
||||
await self.client.start()
|
||||
logger.info("Клиент Telegram успешно запущен.")
|
||||
|
||||
if is_first_login:
|
||||
logger.info("Аккаунт Telegram успешно подключен.")
|
||||
logger.info("Пожалуйста, перезапустите скрипт для начала работы.")
|
||||
return
|
||||
|
||||
logger.info("=" * 40)
|
||||
logger.info(
|
||||
f"Комментарии на форуме: {'ВКЛЮЧЕНЫ' if self.config.enable_reply else 'ВЫКЛЮЧЕНЫ'}"
|
||||
)
|
||||
logger.info(
|
||||
f"Пропуск постов с комментариями: {'ВКЛЮЧЕН' if self.config.skip_posts_with_comments else 'ВЫКЛЮЧЕН'}"
|
||||
)
|
||||
logger.info(f"Проверка начинается со страницы: {self.start_page}")
|
||||
logger.info("Бот в работе. Для остановки нажмите Ctrl+C.")
|
||||
logger.info("=" * 40)
|
||||
await self._main_loop()
|
||||
|
||||
finally:
|
||||
if self.client and self.client.is_connected:
|
||||
await self.client.stop()
|
||||
|
||||
logger.info("Бот остановлен.")
|
||||
Reference in New Issue
Block a user