151 lines
6.5 KiB
Python
151 lines
6.5 KiB
Python
import os
|
|
import logging
|
|
from typing import Optional, Callable, Awaitable
|
|
from telethon import TelegramClient, events
|
|
from telethon.errors import SessionPasswordNeededError
|
|
from telethon.tl.types import MessageMediaPhoto, MessageMediaDocument
|
|
from db.repository import Repository
|
|
from core.dedup import compute_content_hash, compute_file_hash
|
|
from services.ai_processor import AIProcessor
|
|
from core.metrics import COLLECTED_POSTS_TOTAL
|
|
from core.proxy import get_telegram_proxy
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
SESSION_DIR = os.getenv("SESSION_DIR", "/app/sessions" if os.path.exists("/app") else "/projects/telegram-bots/copykar/sessions")
|
|
MEDIA_DIR = os.getenv("MEDIA_DIR", "/app/data/media" if os.path.exists("/app") else "/projects/telegram-bots/copykar/data/media")
|
|
|
|
class CollectorService:
|
|
def __init__(
|
|
self,
|
|
repo: Repository,
|
|
ai_processor: AIProcessor,
|
|
api_id: Optional[int] = None,
|
|
api_hash: Optional[str] = None,
|
|
phone: Optional[str] = None,
|
|
session_name: Optional[str] = None,
|
|
):
|
|
self.repo = repo
|
|
self.ai_processor = ai_processor
|
|
self.api_id = api_id or int(os.getenv("API_ID", "0"))
|
|
self.api_hash = api_hash or os.getenv("API_HASH", "")
|
|
self.phone = phone or os.getenv("PHONE")
|
|
self.session_name = session_name or os.path.join(SESSION_DIR, "collector.session")
|
|
os.makedirs(os.path.dirname(self.session_name), exist_ok=True)
|
|
os.makedirs(MEDIA_DIR, exist_ok=True)
|
|
self.client = TelegramClient(self.session_name, self.api_id, self.api_hash, proxy=get_telegram_proxy())
|
|
self.phone_code_hash: Optional[str] = None
|
|
self._handlers_registered = False
|
|
|
|
async def start(self, notify_fn: Optional[Callable[[str], Awaitable[None]]] = None):
|
|
logger.info("Initializing Collector Userbot client...")
|
|
await self.client.connect()
|
|
|
|
if await self.client.is_user_authorized():
|
|
me = await self.client.get_me()
|
|
logger.info(f"Collector Userbot is authorized as: {me.first_name} (@{me.username})")
|
|
self._register_handlers()
|
|
return True
|
|
|
|
logger.warning("Collector Userbot is not authorized. Requesting login code...")
|
|
if self.phone and notify_fn:
|
|
try:
|
|
sent = await self.client.send_code_request(self.phone)
|
|
self.phone_code_hash = sent.phone_code_hash
|
|
await notify_fn(
|
|
f"🔐 <b>Collector Userbot Login Required</b>\n\n"
|
|
f"A login code was sent to phone <code>{self.phone}</code>.\n\n"
|
|
f"Please reply with: <code>/code <your_code></code>\n"
|
|
f"(Or <code>/password <2fa_password></code> if 2FA is enabled)."
|
|
)
|
|
except Exception as e:
|
|
logger.error(f"Failed to send login code request: {e}")
|
|
await notify_fn(f"❌ Failed to request login code: {e}")
|
|
return False
|
|
|
|
async def submit_code(self, code: str) -> str:
|
|
if not self.phone or not self.phone_code_hash:
|
|
# Re-request code
|
|
sent = await self.client.send_code_request(self.phone)
|
|
self.phone_code_hash = sent.phone_code_hash
|
|
|
|
try:
|
|
await self.client.sign_in(phone=self.phone, code=code, phone_code_hash=self.phone_code_hash)
|
|
me = await self.client.get_me()
|
|
self._register_handlers()
|
|
return f"✅ Logged in successfully as <b>{me.first_name}</b> (@{me.username or 'none'}). Collector is now active!"
|
|
except SessionPasswordNeededError:
|
|
return "🔐 <b>Two-Factor Authentication (2FA) is enabled.</b> Please send: <code>/password <your_2fa_password></code>"
|
|
except Exception as e:
|
|
return f"❌ Login failed: {e}"
|
|
|
|
async def submit_password(self, password: str) -> str:
|
|
try:
|
|
await self.client.sign_in(password=password)
|
|
me = await self.client.get_me()
|
|
self._register_handlers()
|
|
return f"✅ 2FA Verified! Logged in as <b>{me.first_name}</b> (@{me.username or 'none'}). Collector is now active!"
|
|
except Exception as e:
|
|
return f"❌ 2FA verification failed: {e}"
|
|
|
|
def _register_handlers(self):
|
|
if self._handlers_registered:
|
|
return
|
|
|
|
@self.client.on(events.NewMessage)
|
|
async def on_new_message(event: events.NewMessage.Event):
|
|
await self._handle_message(event)
|
|
|
|
self._handlers_registered = True
|
|
logger.info("Collector real-time event handlers registered.")
|
|
|
|
async def _handle_message(self, event: events.NewMessage.Event):
|
|
try:
|
|
chat_id = event.chat_id
|
|
source = await self.repo.get_source_by_channel_id(chat_id)
|
|
if not source or not source.is_active:
|
|
return
|
|
|
|
raw_text = event.raw_text or ""
|
|
media_path = None
|
|
media_type = None
|
|
media_hash = None
|
|
|
|
if event.message.media:
|
|
if isinstance(event.message.media, MessageMediaPhoto):
|
|
media_type = "photo"
|
|
elif isinstance(event.message.media, MessageMediaDocument):
|
|
media_type = "document"
|
|
else:
|
|
media_type = "other"
|
|
|
|
filename = f"{chat_id}_{event.message.id}"
|
|
download_target = os.path.join(MEDIA_DIR, filename)
|
|
downloaded_file = await event.message.download_media(file=download_target)
|
|
if downloaded_file:
|
|
media_path = downloaded_file
|
|
media_hash = compute_file_hash(downloaded_file)
|
|
|
|
content_hash = compute_content_hash(raw_text, media_hash)
|
|
|
|
post_id = await self.repo.create_raw_post(
|
|
source_channel_id=chat_id,
|
|
source_message_id=event.message.id,
|
|
raw_text=raw_text,
|
|
media_path=media_path,
|
|
media_type=media_type,
|
|
content_hash=content_hash,
|
|
)
|
|
|
|
if post_id:
|
|
COLLECTED_POSTS_TOTAL.labels(source_channel_id=str(chat_id)).inc()
|
|
logger.info(f"Collected new post ID {post_id} from channel {chat_id}")
|
|
await self.ai_processor.process_post(post_id)
|
|
except Exception as e:
|
|
logger.error(f"Error handling message from {event.chat_id}: {e}", exc_info=True)
|
|
|
|
async def stop(self):
|
|
if self.client.is_connected():
|
|
await self.client.disconnect()
|
|
logger.info("Collector Userbot disconnected.")
|