diff --git a/services/admin_bot.py b/services/admin_bot.py
index 5338653..5bf448d 100644
--- a/services/admin_bot.py
+++ b/services/admin_bot.py
@@ -12,6 +12,14 @@ logger = logging.getLogger(__name__)
SESSION_DIR = os.getenv("SESSION_DIR", "/app/sessions" if os.path.exists("/app") else "/projects/telegram-bots/copykar/sessions")
+def get_main_menu_keyboard():
+ return [
+ [Button.text("🔑 Request Login Code", resize=True), Button.text("📊 Fleet Statistics", resize=True)],
+ [Button.text("📡 Monitored Sources", resize=True), Button.text("🎯 Target Channels", resize=True)],
+ [Button.text("➕ Add Source Guide", resize=True), Button.text("➕ Add Target Guide", resize=True)],
+ [Button.text("❓ Help & Documentation", resize=True)]
+ ]
+
class AdminBotService:
def __init__(
self,
@@ -51,7 +59,7 @@ class AdminBotService:
for admin_id in self.admin_user_ids:
try:
- await self.client.send_message(admin_id, text, parse_mode="html")
+ await self.client.send_message(admin_id, text, parse_mode="html", buttons=get_main_menu_keyboard())
except Exception as e:
logger.debug(f"Could not send DM to admin {admin_id}: {e}")
@@ -62,40 +70,200 @@ class AdminBotService:
self._register_handlers()
def _register_handlers(self):
+ # --- /start and Main Menu ---
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/start|/menu|menu)$"))
+ async def cmd_start(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ await event.reply(f"⛔ Unauthorized user ID: {event.sender_id}. Please add this ID to ADMIN_USER_IDS in .env.", parse_mode="html")
+ return
+
+ userbot_status = "🔴 Not Authorized"
+ if self.collector and self.collector.client.is_connected() and await self.collector.client.is_user_authorized():
+ me = await self.collector.client.get_me()
+ userbot_status = f"🟢 Online ({me.first_name})"
+
+ welcome_text = (
+ "👋 Welcome to Copykar Admin Console!\n\n"
+ f"• 🤖 Userbot Status: {userbot_status}\n"
+ f"• 📋 Review Channel: {self.review_channel_id}\n\n"
+ "Use the interactive menu buttons below to manage the fleet:"
+ )
+ await event.reply(welcome_text, parse_mode="html", buttons=get_main_menu_keyboard())
+
# --- 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"))
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/request_code|🔑 Request Login 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)
+ msg = await event.reply("⏳ Contacting Telegram to request login code...")
+ try:
+ sent = await self.collector.client.send_code_request(self.collector.phone)
+ self.collector.phone_code_hash = sent.phone_code_hash
+ await msg.edit(
+ f"📩 Code Sent!\n\n"
+ f"Telegram sent a verification code to {self.collector.phone}.\n\n"
+ f"Please reply with:\n"
+ f"/code <your_code>\n\n"
+ f"Example: /code 12345",
+ parse_mode="html"
+ )
+ except Exception as e:
+ await msg.edit(f"❌ Could not request code: {e}")
+
+ @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).strip()
+ status_msg = await event.reply("⏳ Submitting code to Telegram...")
+ result = await self.collector.submit_code(code)
+ await status_msg.edit(result, parse_mode="html", buttons=get_main_menu_keyboard())
+
+ @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).strip()
+ status_msg = await event.reply("⏳ Verifying 2FA password...")
+ result = await self.collector.submit_password(pwd)
+ await status_msg.edit(result, parse_mode="html", buttons=get_main_menu_keyboard())
+
+ # --- Statistics ---
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/stats|📊 Fleet Statistics)$"))
+ async def cmd_stats(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ pending_ai = len(await self.repo.get_posts_by_status("pending_ai", limit=1000))
+ pending_review = len(await self.repo.get_posts_by_status("pending_review", limit=1000))
+ approved = len(await self.repo.get_posts_by_status("approved", limit=1000))
+ published = len(await self.repo.get_posts_by_status("published", limit=1000))
+ rejected = len(await self.repo.get_posts_by_status("rejected", limit=1000))
+
+ text = (
+ "📊 Copykar Fleet Metrics\n\n"
+ f"• ⏳ Pending AI: {pending_ai}\n"
+ f"• 📋 Pending Review: {pending_review}\n"
+ f"• 🚀 Approved (In Queue): {approved}\n"
+ f"• ✅ Published: {published}\n"
+ f"• ❌ Rejected: {rejected}\n\n"
+ "📈 Grafana Dashboard: http://localhost:3000"
+ )
+ await event.reply(text, parse_mode="html", buttons=get_main_menu_keyboard())
+
+ # --- Sources Management ---
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/sources|📡 Monitored Sources)$"))
+ async def cmd_sources(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ sources = await self.repo.get_active_sources()
+ if not sources:
+ await event.reply("No active source channels configured.\nUse ➕ Add Source Guide to add one.", parse_mode="html")
+ return
+ lines = ["📡 Active Monitored Sources:\n"]
+ for s in sources:
+ lines.append(f"• ID: {s.channel_id} | {s.title or 'N/A'} (@{s.username or 'none'})")
+ await event.reply("\n".join(lines), parse_mode="html", buttons=get_main_menu_keyboard())
+
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(➕ Add Source Guide)$"))
+ async def cmd_add_source_guide(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ guide = (
+ "➕ How to Add a Source Channel:\n\n"
+ "Send the command in this format:\n"
+ "/add_source <channel_id> <title> [username]\n\n"
+ "Example:\n"
+ "/add_source -1001234567890 TechNews technews_chan"
+ )
+ await event.reply(guide, parse_mode="html")
+
+ @self.client.on(events.NewMessage(pattern=r"^/add_source\s+(-?\d+)\s+([^\s]+)(?:\s+([^\s]+))?"))
+ async def cmd_add_source(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ ch_id = int(event.pattern_match.group(1))
+ title = event.pattern_match.group(2)
+ username = event.pattern_match.group(3)
+ await self.repo.add_source(channel_id=ch_id, title=title, username=username)
+ await event.reply(f"✅ Added source channel {title} ({ch_id}) to monitoring.", parse_mode="html", buttons=get_main_menu_keyboard())
+
+ # --- Targets Management ---
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/targets|🎯 Target Channels)$"))
+ async def cmd_targets(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ targets = await self.repo.get_active_targets()
+ if not targets:
+ await event.reply("No target channels configured.\nUse ➕ Add Target Guide to add one.", parse_mode="html")
+ return
+ lines = ["🎯 Target Publishing Channels:\n"]
+ for t in targets:
+ lines.append(f"• ID: {t.id} (Channel: {t.channel_id})\n Title: {t.title} | Interval: {t.post_interval_min}m")
+ await event.reply("\n".join(lines), parse_mode="html", buttons=get_main_menu_keyboard())
+
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(➕ Add Target Guide)$"))
+ async def cmd_add_target_guide(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ guide = (
+ "🎯 How to Add a Target Channel:\n\n"
+ "Send the command in this format:\n"
+ "/add_target <channel_id> <title> <interval_minutes> [username]\n\n"
+ "Example (posts every 30 minutes):\n"
+ "/add_target -1009876543210 MyMainChannel 30 my_main_chan"
+ )
+ await event.reply(guide, parse_mode="html")
+
+ @self.client.on(events.NewMessage(pattern=r"^/add_target\s+(-?\d+)\s+([^\s]+)\s+(\d+)(?:\s+([^\s]+))?"))
+ async def cmd_add_target(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ ch_id = int(event.pattern_match.group(1))
+ title = event.pattern_match.group(2)
+ interval_min = int(event.pattern_match.group(3))
+ username = event.pattern_match.group(4)
+ await self.repo.add_target(channel_id=ch_id, title=title, username=username, post_interval_min=interval_min)
+ await event.reply(f"✅ Added target channel {title} with interval {interval_min}m.", parse_mode="html", buttons=get_main_menu_keyboard())
+
+ @self.client.on(events.NewMessage(pattern=r"^/set_interval\s+(\d+)\s+(\d+)"))
+ async def cmd_set_interval(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ target_id = int(event.pattern_match.group(1))
+ new_interval = int(event.pattern_match.group(2))
+ target = await self.repo.get_target_by_id(target_id)
+ if not target:
+ await event.reply("Target channel not found.")
+ return
+ await self.repo.add_target(
+ channel_id=target.channel_id,
+ title=target.title,
+ username=target.username,
+ post_interval_min=new_interval
+ )
+ await event.reply(f"✅ Updated interval for {target.title} to {new_interval} minutes.", parse_mode="html", buttons=get_main_menu_keyboard())
+
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/help|❓ Help & Documentation)$"))
+ async def cmd_help(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ help_text = (
+ "📖 Copykar Bot Quick Help\n\n"
+ "1. Authentication: Click 🔑 Request Login Code and submit via /code <12345>.\n"
+ "2. Monitored Sources: Add channels the userbot should listen to via /add_source.\n"
+ "3. Review Flow: AI scans posts, checks duplicates, and sends drafts to the review channel with inline approval buttons.\n"
+ "4. Publishing: Approved posts are published to your target channels strictly according to their interval minutes."
+ )
+ await event.reply(help_text, parse_mode="html", buttons=get_main_menu_keyboard())
# --- Review Keyboard Callbacks ---
@self.client.on(events.CallbackQuery)
@@ -137,90 +305,6 @@ class AdminBotService:
)
await event.answer("Post rejected.")
- # --- 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):
- return
- sources = await self.repo.get_active_sources()
- if not sources:
- await event.reply("No active source channels configured. Use /add_source [username]")
- return
- lines = ["Active Monitored Sources:"]
- for s in sources:
- lines.append(f"• ID: {s.channel_id} | Title: {s.title or 'N/A'} (@{s.username or 'none'})")
- await event.reply("\n".join(lines), parse_mode="html")
-
- @self.client.on(events.NewMessage(pattern=r"/add_source\s+(-?\d+)\s+([^\s]+)(?:\s+([^\s]+))?"))
- async def cmd_add_source(event: events.NewMessage.Event):
- if not self.is_admin(event.sender_id):
- return
- ch_id = int(event.pattern_match.group(1))
- title = event.pattern_match.group(2)
- username = event.pattern_match.group(3)
- await self.repo.add_source(channel_id=ch_id, title=title, username=username)
- await event.reply(f"✅ Added source channel {title} ({ch_id}).", parse_mode="html")
-
- @self.client.on(events.NewMessage(pattern="/targets"))
- async def cmd_targets(event: events.NewMessage.Event):
- if not self.is_admin(event.sender_id):
- return
- targets = await self.repo.get_active_targets()
- if not targets:
- await event.reply("No target channels configured. Use /add_target [username]")
- return
- lines = ["Target Publishing Channels:"]
- for t in targets:
- lines.append(f"• ID: {t.id} (Channel: {t.channel_id}) | {t.title} | Interval: {t.post_interval_min}m")
- await event.reply("\n".join(lines), parse_mode="html")
-
- @self.client.on(events.NewMessage(pattern=r"/add_target\s+(-?\d+)\s+([^\s]+)\s+(\d+)(?:\s+([^\s]+))?"))
- async def cmd_add_target(event: events.NewMessage.Event):
- if not self.is_admin(event.sender_id):
- return
- ch_id = int(event.pattern_match.group(1))
- title = event.pattern_match.group(2)
- interval_min = int(event.pattern_match.group(3))
- username = event.pattern_match.group(4)
- await self.repo.add_target(channel_id=ch_id, title=title, username=username, post_interval_min=interval_min)
- await event.reply(f"✅ Added target channel {title} with interval {interval_min}m.", parse_mode="html")
-
- @self.client.on(events.NewMessage(pattern=r"/set_interval\s+(\d+)\s+(\d+)"))
- async def cmd_set_interval(event: events.NewMessage.Event):
- if not self.is_admin(event.sender_id):
- return
- target_id = int(event.pattern_match.group(1))
- new_interval = int(event.pattern_match.group(2))
- target = await self.repo.get_target_by_id(target_id)
- if not target:
- await event.reply("Target channel not found.")
- return
- await self.repo.add_target(
- channel_id=target.channel_id,
- title=target.title,
- username=target.username,
- post_interval_min=new_interval
- )
- await event.reply(f"✅ Updated interval for {target.title} to {new_interval} minutes.", parse_mode="html")
-
- @self.client.on(events.NewMessage(pattern="/stats"))
- async def cmd_stats(event: events.NewMessage.Event):
- if not self.is_admin(event.sender_id):
- return
- pending_ai = len(await self.repo.get_posts_by_status("pending_ai", limit=1000))
- pending_review = len(await self.repo.get_posts_by_status("pending_review", limit=1000))
- approved = len(await self.repo.get_posts_by_status("approved", limit=1000))
- published = len(await self.repo.get_posts_by_status("published", limit=1000))
-
- text = (
- "📊 Copykar Bot Statistics\n\n"
- f"• ⏳ Pending AI: {pending_ai}\n"
- f"• 📋 Pending Admin Review: {pending_review}\n"
- f"• 🚀 Approved (Queued): {approved}\n"
- f"• ✅ Published: {published}\n"
- )
- await event.reply(text, parse_mode="html")
-
async def send_review_post(self, post_id: int):
post = await self.repo.get_post_by_id(post_id)
if not post or not self.review_channel_id: