diff --git a/main.py b/main.py
index c699afc..18c99e3 100644
--- a/main.py
+++ b/main.py
@@ -51,11 +51,12 @@ async def main():
ai_processor.process_post = process_and_notify
collector = CollectorService(repo=repo, ai_processor=ai_processor)
+ admin_bot.set_collector(collector)
publisher = PublisherService(repo=repo)
# 4. Start all services
await admin_bot.start()
- await collector.start()
+ await collector.start(notify_fn=admin_bot.notify_admins)
await publisher.start()
logger.info("All Copykar services are active and running.")
diff --git a/services/admin_bot.py b/services/admin_bot.py
index 4895143..5338653 100644
--- a/services/admin_bot.py
+++ b/services/admin_bot.py
@@ -33,10 +33,28 @@ class AdminBotService:
self.session_name = session_name or os.path.join(SESSION_DIR, "admin_bot.session")
os.makedirs(os.path.dirname(self.session_name), exist_ok=True)
self.client = TelegramClient(self.session_name, self.api_id, self.api_hash, proxy=get_telegram_proxy())
+ self.collector = None
+
+ def set_collector(self, collector):
+ self.collector = collector
def is_admin(self, user_id: int) -> bool:
return not self.admin_user_ids or user_id in self.admin_user_ids
+ async def notify_admins(self, text: str):
+ """Broadcast message to review channel and all admin DMs."""
+ if self.review_channel_id:
+ try:
+ await self.client.send_message(self.review_channel_id, text, parse_mode="html")
+ except Exception as e:
+ logger.error(f"Failed to notify review channel: {e}")
+
+ for admin_id in self.admin_user_ids:
+ try:
+ await self.client.send_message(admin_id, text, parse_mode="html")
+ except Exception as e:
+ logger.debug(f"Could not send DM to admin {admin_id}: {e}")
+
async def start(self):
logger.info("Starting Admin Review Bot...")
await self.client.start(bot_token=self.bot_token)
@@ -44,6 +62,42 @@ class AdminBotService:
self._register_handlers()
def _register_handlers(self):
+ # --- Interactive Userbot Authentication Commands ---
+ @self.client.on(events.NewMessage(pattern=r"/code\s+(\S+)"))
+ async def cmd_code(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ if not self.collector:
+ await event.reply("Collector service not linked.")
+ return
+ code = event.pattern_match.group(1)
+ status_msg = await event.reply("⏳ Verifying code with Telegram...")
+ result = await self.collector.submit_code(code)
+ await status_msg.edit(result, parse_mode="html")
+
+ @self.client.on(events.NewMessage(pattern=r"/password\s+(.+)"))
+ async def cmd_password(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ if not self.collector:
+ await event.reply("Collector service not linked.")
+ return
+ pwd = event.pattern_match.group(1)
+ status_msg = await event.reply("⏳ Verifying 2FA password...")
+ result = await self.collector.submit_password(pwd)
+ await status_msg.edit(result, parse_mode="html")
+
+ @self.client.on(events.NewMessage(pattern="/request_code"))
+ async def cmd_request_code(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ if not self.collector:
+ await event.reply("Collector service not linked.")
+ return
+ await event.reply("Requesting new login code...")
+ await self.collector.start(notify_fn=self.notify_admins)
+
+ # --- Review Keyboard Callbacks ---
@self.client.on(events.CallbackQuery)
async def on_callback(event: events.CallbackQuery.Event):
if not self.is_admin(event.sender_id):
@@ -83,7 +137,7 @@ class AdminBotService:
)
await event.answer("Post rejected.")
- # --- Admin Commands ---
+ # --- Admin Configuration Commands ---
@self.client.on(events.NewMessage(pattern="/sources"))
async def cmd_sources(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
diff --git a/services/collector.py b/services/collector.py
index 35db427..6d7322b 100644
--- a/services/collector.py
+++ b/services/collector.py
@@ -1,7 +1,8 @@
import os
import logging
-from typing import Optional
+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
@@ -33,19 +34,71 @@ class CollectorService:
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):
- logger.info("Starting Collector Userbot...")
- if self.phone:
- await self.client.start(phone=self.phone)
- else:
- await self.client.start()
- logger.info("Collector Userbot connected successfully.")
+ 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"🔐 Collector Userbot Login Required\n\n"
+ f"A login code was sent to phone {self.phone}.\n\n"
+ f"Please reply with: /code <your_code>\n"
+ f"(Or /password <2fa_password> 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 {me.first_name} (@{me.username or 'none'}). Collector is now active!"
+ except SessionPasswordNeededError:
+ return "🔐 Two-Factor Authentication (2FA) is enabled. Please send: /password <your_2fa_password>"
+ 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 {me.first_name} (@{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