diff --git a/services/collector.py b/services/collector.py new file mode 100644 index 0000000..ec5c4d8 --- /dev/null +++ b/services/collector.py @@ -0,0 +1,95 @@ +import os +import logging +from typing import Optional +from telethon import TelegramClient, events +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 + +logger = logging.getLogger(__name__) + +MEDIA_DIR = os.getenv("MEDIA_DIR", "/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, + session_name: str = "/projects/telegram-bots/copykar/sessions/collector.session", + ): + 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.session_name = session_name + self.client = TelegramClient(self.session_name, self.api_id, self.api_hash) + + async def start(self): + os.makedirs(os.path.dirname(self.session_name), exist_ok=True) + os.makedirs(MEDIA_DIR, exist_ok=True) + + logger.info("Starting Collector Userbot...") + await self.client.start() + logger.info("Collector Userbot connected successfully.") + + @self.client.on(events.NewMessage) + async def on_new_message(event: events.NewMessage.Event): + await self._handle_message(event) + + async def _handle_message(self, event: events.NewMessage.Event): + try: + # Check if source channel is in our monitored sources + 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 + + # Download media if present + 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) + + # Store raw post in DB + 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}") + # Trigger AI Processor pipeline + 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.")