import asyncio
import os
import re
import time
import logging
import threading
import pyrogram
from pyrogram import Client, filters, handlers
from pyrogram.types import InlineKeyboardMarkup, InlineKeyboardButton, Message, CallbackQuery
from sqlalchemy.orm import sessionmaker

from config import Config
from database import engine, db_session, Account, Setting, LiveConfig, ReactionRule, SystemLog, log_system_event
from live_joiner import TGVoiceCallManager, parse_multiple_chat_targets
from auto_reaction import AutoReactionEngine
from services import TelegramServicesEngine

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("TGBotManager")

# Configured Telegram Client API Credentials
DEFAULT_API_ID = 35408643
DEFAULT_API_HASH = "9c690318646543daa057cc714e91f540"

def get_fresh_db():
    """Create fresh SQLite DB session per operation to bypass thread-local transaction caching"""
    Session = sessionmaker(bind=engine)
    return Session()

def get_api_credentials():
    api_id = Config.API_ID or Setting.get("API_ID", str(DEFAULT_API_ID))
    api_hash = Config.API_HASH or Setting.get("API_HASH", DEFAULT_API_HASH)
    return int(api_id), api_hash

def sanitize_phone_number(phone_raw: str) -> str:
    """Format phone number to international E.164 standard (e.g. +91XXXXXXXXXX)"""
    digits = re.sub(r'\D', '', phone_raw.strip())
    if len(digits) == 10 and digits[0] in ['6', '7', '8', '9']:
        return f"+91{digits}"
    elif len(digits) == 12 and digits.startswith('91'):
        return f"+{digits}"
    elif digits:
        return f"+{digits}"
    return phone_raw

# Permanent Daemon Asyncio Loop Manager with Phusion Passenger Process Fork Safety
class GlobalAsyncLoop:
    _instance = None
    _lock = threading.Lock()

    def __init__(self):
        self._pid = os.getpid()
        self._setup_loop()

    def _setup_loop(self):
        self.loop = asyncio.new_event_loop()
        self.thread = threading.Thread(target=self._start_loop, name=f"GlobalAsyncThread_{self._pid}", daemon=True)
        self.thread.start()

    def _start_loop(self):
        asyncio.set_event_loop(self.loop)
        self.loop.run_forever()

    @classmethod
    def get_loop(cls):
        with cls._lock:
            current_pid = os.getpid()
            if cls._instance is None or cls._instance._pid != current_pid or not cls._instance.thread.is_alive() or not cls._instance.loop.is_running():
                cls._instance = cls()
            return cls._instance.loop

class BotManager:
    _instance = None
    _lock = threading.Lock()

    def __init__(self):
        self.bot_client = None
        self.user_clients = {}
        self.temp_clients = {}
        self.primary_userbot = None
        self.live_manager = None
        self.services_engine = None
        self.reaction_engines = []
        self.is_running = False
        self.loop = GlobalAsyncLoop.get_loop()
        self.heartbeat_task = None
        self._is_loading_accounts = False
        self._is_initializing = False
        self._auto_join_running = False

    @classmethod
    def get_instance(cls):
        with cls._lock:
            current_loop = GlobalAsyncLoop.get_loop()
            if cls._instance is None or cls._instance.loop != current_loop:
                cls._instance = cls()
            return cls._instance

    async def initialize(self):
        """Initialize Admin Bot and Active Userbot Sessions"""
        if self.is_running or self._is_initializing:
            self.start_heartbeat_monitor()
            return
        self._is_initializing = True
        logger.info("Initializing Bot Manager...")
        try:
            bot_token = Config.BOT_TOKEN or Setting.get("BOT_TOKEN", "")
            api_id, api_hash = get_api_credentials()

            if bot_token and api_id and api_hash and not self.bot_client:
                try:
                    self.bot_client = Client(
                        name="admin_bot",
                        api_id=api_id,
                        api_hash=api_hash,
                        bot_token=bot_token,
                        workdir=Config.SESSIONS_DIR
                    )
                    self._setup_admin_bot_handlers()
                    await self.bot_client.start()
                    logger.info("Telegram Admin Bot started successfully.")
                    log_system_event("Admin Bot connected to Telegram.")
                except Exception as e:
                    logger.error(f"Error starting Admin Bot: {e}")
                    log_system_event(f"Admin Bot startup failed: {e}", level="ERROR")

            await self.load_user_accounts()
            self.is_running = True
            self.start_heartbeat_monitor()
        finally:
            self._is_initializing = False

    def start_heartbeat_monitor(self):
        """Ensure permanent background heartbeat loop is active"""
        if not self.heartbeat_task or self.heartbeat_task.done():
            try:
                running_loop = asyncio.get_running_loop()
            except RuntimeError:
                running_loop = None
            if running_loop is self.loop:
                self.heartbeat_task = self.loop.create_task(self._heartbeat_loop())
            else:
                self.heartbeat_task = asyncio.run_coroutine_threadsafe(self._heartbeat_loop(), self.loop)

    async def _heartbeat_loop(self):
        """Permanent background daemon task to maintain voice chat connection and auto-rejoin"""
        logger.info("Starting permanent TG Live Voice Chat Heartbeat & Auto-Rejoin Monitor...")
        while True:
            try:
                await asyncio.sleep(25)

                if not self.user_clients and not getattr(self, "_is_loading_accounts", False):
                    await self.load_user_accounts()

                # Re-connect disconnected userbot clients silently without network spam
                for acc_id, client in list(self.user_clients.items()):
                    try:
                        if not client.is_connected:
                            await client.connect()
                    except Exception as ex:
                        logger.warning(f"Re-connect notice for #{acc_id}: {ex}")

                # Auto-Join monitor if Live is configured with is_auto_join
                db = get_fresh_db()
                cfg = db.query(LiveConfig).first()
                auto_enabled = bool(cfg and cfg.is_auto_join and cfg.target_chat)
                target_chat = cfg.target_chat if auto_enabled else ""
                is_muted = cfg.is_muted if cfg else True
                db.close()

                if auto_enabled:
                    await self.run_auto_join_cycle(target_chat, is_muted)
                elif self.live_manager and self.live_manager.joined_accounts:
                    await self.live_manager.leave_live()

            except Exception as e:
                logger.error(f"Heartbeat loop exception: {e}")
                await asyncio.sleep(10)

    async def run_auto_join_cycle(self, target_chat: str, is_muted: bool = True):
        """Run one serialized auto-detect/join cycle without overlapping calls."""
        if self._auto_join_running:
            return False, "Auto-join cycle already running"
        self._auto_join_running = True
        try:
            success, message = await self.load_user_accounts_and_join(target_chat, is_muted)
            return success, message
        except Exception as exc:
            raise
        finally:
            self._auto_join_running = False

    async def _start_single_userbot(self, acc_id, session_str, name, phone, api_id, api_hash, stagger_delay=0.0):
        """Start a single userbot session safely with direct MTProto socket connection"""
        if acc_id in self.user_clients:
            c = self.user_clients[acc_id]
            try:
                if c.is_connected:
                    return True
                await c.connect()
                return True
            except Exception:
                pass

        if stagger_delay > 0:
            await asyncio.sleep(stagger_delay)

        # Cleanup any leftover .session files on disk to prevent SQLite database locks on cPanel
        for ext in [".session", ".session-journal"]:
            sess_file = f"acc_{acc_id}{ext}"
            if os.path.exists(sess_file):
                try:
                    os.remove(sess_file)
                except Exception:
                    pass

        try:
            client = Client(
                name=f"acc_{acc_id}",
                api_id=api_id,
                api_hash=api_hash,
                session_string=session_str,
                in_memory=True,
                workers=1,
                app_version="10.8.1",
                device_model="Samsung Galaxy S23",
                system_version="Android 14",
                lang_code="en"
            )
            # Direct MTProto connection - 0.05s response time without update listener overhead
            await client.connect()
            me = await client.get_me()

            acc_name = f"{me.first_name or ''} {me.last_name or ''}".strip() or name
            acc_phone = me.phone_number or phone or "Active"

            db = get_fresh_db()
            acc_rec = db.query(Account).get(acc_id)
            if acc_rec:
                acc_rec.name = acc_name
                acc_rec.phone_number = acc_phone
                acc_rec.is_active = True
                db.commit()
            db.close()

            self.user_clients[acc_id] = client
            self.reaction_engines.append(AutoReactionEngine(client))

            log_system_event(f"✅ Userbot Account #{acc_id} loaded & online: {acc_name} (+{acc_phone})")
            logger.info(f"Loaded Userbot: {acc_name}")
            return True

        except (pyrogram.errors.AuthKeyUnregistered, pyrogram.errors.UserDeactivated, pyrogram.errors.SessionRevoked, pyrogram.errors.UserDeactivatedBan) as auth_err:
            logger.warning(f"Account #{acc_id} session revoked/expired: {auth_err}")
            log_system_event(f"⚠️ Account #{acc_id} ({name}) session expired/invalidated on Telegram.", level="WARNING")
            return False

        except Exception as e:
            logger.error(f"Failed to load userbot account ID {acc_id}: {e}")
            log_system_event(f"❌ Userbot Account #{acc_id} connection error: {type(e).__name__} - {e}", level="ERROR")
            return False

    async def load_user_accounts(self):
        """Load and start MTProto Userbot sessions from database concurrently in parallel"""
        if getattr(self, "_is_loading_accounts", False):
            for _ in range(10):
                if not getattr(self, "_is_loading_accounts", False):
                    break
                await asyncio.sleep(0.3)

        self._is_loading_accounts = True
        try:
            db = get_fresh_db()
            api_id, api_hash = get_api_credentials()
            accounts = db.query(Account).all()
            db.close()

            if not accounts:
                log_system_event("⚠️ No userbot accounts found in database. Add accounts in Accounts tab first.", level="WARNING")
                return

            pending_accs = [acc for acc in accounts if acc.id not in self.user_clients or not self.user_clients[acc.id].is_connected]

            if pending_accs:
                log_system_event(f"Connecting {len(pending_accs)} userbot accounts to Telegram in parallel...")
                tasks = []
                stagger = 0.0
                for acc in pending_accs:
                    tasks.append(self._start_single_userbot(acc.id, acc.session_string, acc.name, acc.phone_number, api_id, api_hash, stagger_delay=stagger))
                    stagger += 0.1 # 0.1s smooth stagger

                await asyncio.gather(*tasks, return_exceptions=True)

            if self.user_clients:
                if not self.live_manager:
                    self.live_manager = TGVoiceCallManager(self.user_clients)
                    await self.live_manager.start_all()
                else:
                    self.live_manager.user_clients = self.user_clients

                if not self.services_engine:
                    self.services_engine = TelegramServicesEngine(self.user_clients)

                log_system_event(f"🎉 Successfully loaded & verified {len(self.user_clients)}/{len(accounts)} userbot accounts online in memory!")

            self.start_heartbeat_monitor()
        finally:
            self._is_loading_accounts = False

    async def load_user_accounts_and_join(self, target_chat: str, is_muted: bool = True):
        """Synchronous/Asynchronous task to load accounts and join Live Voice Chat reliably"""
        try:
            log_system_event(f"🔄 Preparing userbot accounts for Live Join...")
            await self.load_user_accounts()
            
            if not self.live_manager and self.user_clients:
                self.live_manager = TGVoiceCallManager(self.user_clients)
            elif self.live_manager:
                self.live_manager.user_clients = self.user_clients

            if self.live_manager and self.user_clients:
                log_system_event(f"🚀 Executing Live Join for {len(self.user_clients)} accounts into {target_chat}...")
                success, msg = await self.live_manager.join_live(target_chat, is_muted)
                log_system_event(f"🎙️ Live Join Result: {msg}")
                return success, msg
            else:
                err_msg = f"❌ Live Join Failed: No active accounts loaded in memory ({len(self.user_clients)} connected)."
                log_system_event(err_msg, level="ERROR")
                return False, err_msg
        except Exception as ex:
            err_msg = f"❌ Live Join Task Exception: {ex}"
            log_system_event(err_msg, level="ERROR")
            logger.error(err_msg)
            return False, err_msg

    async def send_login_otp(self, phone_number: str):
        """Connect MTProto client and request Telegram OTP code"""
        api_id, api_hash = get_api_credentials()
        phone_full = sanitize_phone_number(phone_number)
        phone_clean = re.sub(r'\D', '', phone_full)
        session_name = f"login_{phone_clean}"

        if phone_clean in self.temp_clients:
            try:
                c = self.temp_clients[phone_clean].get("client")
                if c and c.is_connected:
                    await c.disconnect()
            except Exception:
                pass
            self.temp_clients.pop(phone_clean, None)

        temp_client = Client(
            name=session_name,
            api_id=api_id,
            api_hash=api_hash,
            in_memory=True,
            app_version="10.8.1",
            device_model="Samsung Galaxy S23",
            system_version="Android 14",
            lang_code="en"
        )
        await temp_client.connect()
        sent_code = await temp_client.send_code(phone_full)
        
        Setting.set(f"otp_hash_{phone_clean}", sent_code.phone_code_hash)
        Setting.set(f"otp_phone_{phone_clean}", phone_full)

        self.temp_clients[phone_clean] = {
            "client": temp_client,
            "hash": sent_code.phone_code_hash,
            "phone": phone_full,
            "session_name": session_name
        }
        
        code_type_str = str(getattr(sent_code, "type", "APP")).replace("SentCodeType.", "")
        log_system_event(f"Requested Telegram OTP code for {phone_full} (Delivery: {code_type_str})")
        return sent_code.phone_code_hash, code_type_str, phone_full

    async def verify_login_otp(self, phone_number: str, otp_code: str, two_fa_password: str = ""):
        """Verify Telegram OTP code using active client session"""
        api_id, api_hash = get_api_credentials()
        phone_full = sanitize_phone_number(phone_number)
        phone_clean = re.sub(r'\D', '', phone_full)
        db_phone_full = Setting.get(f"otp_phone_{phone_clean}", phone_full)
        session_name = f"login_{phone_clean}"
        otp_code = re.sub(r'\D', '', otp_code.strip())

        code_hash = Setting.get(f"otp_hash_{phone_clean}", "")
        
        temp_client = None
        if phone_clean in self.temp_clients:
            temp_client = self.temp_clients[phone_clean].get("client")
            if not code_hash:
                code_hash = self.temp_clients[phone_clean].get("hash", "")

        if not code_hash:
            return False, "OTP Session expired. Please click Send OTP Code again."

        if not temp_client:
            temp_client = Client(
                name=session_name,
                api_id=api_id,
                api_hash=api_hash,
                in_memory=True,
                app_version="10.8.1",
                device_model="Samsung Galaxy S23",
                system_version="Android 14",
                lang_code="en"
            )

        try:
            if not temp_client.is_connected:
                await temp_client.connect()
            await temp_client.sign_in(db_phone_full, code_hash, otp_code)
        except Exception as e:
            err_str = str(e)
            if "SESSION_PASSWORD_NEEDED" in err_str:
                if two_fa_password:
                    await temp_client.check_password(two_fa_password)
                else:
                    return False, "2FA Password required! Please enter your 2-Step Verification Password."
            else:
                raise e

        string_session = await temp_client.export_session_string()
        await temp_client.disconnect()
        
        self.temp_clients.pop(phone_clean, None)
        Setting.set(f"otp_hash_{phone_clean}", "")
        Setting.set(f"otp_phone_{phone_clean}", "")

        return await self.add_user_session(string_session, db_phone_full)

    async def get_account_otp_messages(self, account_id: int, limit: int = 5):
        """Fetch recent login OTP codes & messages from Telegram Service Notifications (777000)"""
        if account_id not in self.user_clients:
            return False, "Account session not active or offline", []

        client = self.user_clients[account_id]
        otp_messages = []

        try:
            async for msg in client.get_chat_history(777000, limit=limit):
                if msg.text:
                    codes = re.findall(r'\b\d{5,6}\b', msg.text)
                    otp_messages.append({
                        "message_id": msg.id,
                        "date": msg.date.strftime("%Y-%m-%d %H:%M:%S") if msg.date else "",
                        "text": msg.text,
                        "extracted_codes": codes
                    })

            return True, "Successfully retrieved Telegram messages", otp_messages
        except Exception as e:
            try:
                async for dialog in client.get_dialogs(limit=10):
                    if dialog.chat and (dialog.chat.id == 777000 or "Telegram" in (dialog.chat.first_name or "")):
                        async for msg in client.get_chat_history(dialog.chat.id, limit=limit):
                            if msg.text:
                                codes = re.findall(r'\b\d{5,6}\b', msg.text)
                                otp_messages.append({
                                    "message_id": msg.id,
                                    "date": msg.date.strftime("%Y-%m-%d %H:%M:%S") if msg.date else "",
                                    "text": msg.text,
                                    "extracted_codes": codes
                                })
                        break
                return True, "Retrieved messages", otp_messages
            except Exception as ex:
                return False, f"Failed to fetch Telegram messages: {str(ex)}", []

    def _setup_admin_bot_handlers(self):
        """Setup Telegram Admin Bot Commands and Callbacks"""

        def is_admin(_, __, message: Message):
            admin_id = Config.ADMIN_USER_ID or int(Setting.get("ADMIN_USER_ID", "0"))
            return message.from_user and (message.from_user.id == admin_id or admin_id == 0)

        @self.bot_client.on_message(filters.command(["start", "help"]) & filters.create(is_admin))
        async def start_cmd(client: Client, message: Message):
            welcome_text = (
                "🤖 **Telegram Live Joiner & Multi-Services Bot**\n\n"
                "Welcome Admin! You can control all features, sell services, and manage TG Live from this menu or Web Dashboard.\n\n"
                "Select an option below:"
            )
            await message.reply_text(welcome_text, reply_markup=self._get_admin_keyboard())

        @self.bot_client.on_message(filters.command(["admin", "status", "otp", "services"]) & filters.create(is_admin))
        async def admin_cmd(client: Client, message: Message):
            await message.reply_text("🎛️ **Admin Control Panel**", reply_markup=self._get_admin_keyboard())

        @self.bot_client.on_callback_query()
        async def callback_handler(client: Client, callback: CallbackQuery):
            admin_id = Config.ADMIN_USER_ID or int(Setting.get("ADMIN_USER_ID", "0"))
            if admin_id != 0 and callback.from_user.id != admin_id:
                await callback.answer("❌ Unauthorized access!", show_alert=True)
                return

            data = callback.data

            if data == "menu_live":
                cfg = db_session.query(LiveConfig).first()
                target = cfg.target_chat if cfg else "None"
                status = cfg.last_status if cfg else "Disconnected"
                
                text = (
                    "🎙️ **Telegram Live Voice Chat Settings**\n\n"
                    f"• **Target Chat**: `{target}`\n"
                    f"• **Current Status**: `{status}`\n"
                    f"• **Active Accounts Loaded**: `{len(self.user_clients)}`\n"
                    f"• **Auto Join**: `{'ON' if cfg and cfg.is_auto_join else 'OFF'}`\n"
                    f"• **Muted**: `{'YES' if cfg and cfg.is_muted else 'NO'}`"
                )
                kb = InlineKeyboardMarkup([
                    [
                        InlineKeyboardButton(f"▶️ Join All ({len(self.user_clients)}) Live", callback_data="live_join"),
                        InlineKeyboardButton("⏹️ Leave Live", callback_data="live_leave")
                    ],
                    [
                        InlineKeyboardButton(f"🔊 Mute/Unmute", callback_data="live_toggle_mute"),
                        InlineKeyboardButton(f"🔄 AutoJoin: {'ON' if cfg and cfg.is_auto_join else 'OFF'}", callback_data="live_toggle_autojoin")
                    ],
                    [InlineKeyboardButton("🔙 Back to Main Menu", callback_data="menu_main")]
                ])
                await callback.edit_message_text(text, reply_markup=kb)

            elif data == "live_join":
                cfg = db_session.query(LiveConfig).first()
                if not cfg or not cfg.target_chat:
                    await callback.answer("❌ Target channel/group not configured! Set it in Web Panel.", show_alert=True)
                    return
                await callback.answer(f"Connecting all {len(self.user_clients)} accounts to Live...")
                if self.live_manager:
                    success, msg = await self.live_manager.join_live(cfg.target_chat, cfg.is_muted)
                    await callback.message.reply_text(f"🎙️ Live Join Result: {msg}")
                else:
                    await callback.answer("❌ No active Userbot accounts loaded!", show_alert=True)

            elif data == "live_leave":
                await callback.answer("Leaving Live...")
                if self.live_manager:
                    success, msg = await self.live_manager.leave_live()
                    await callback.message.reply_text(f"ℹ️ {msg}")

            elif data == "live_toggle_mute":
                cfg = db_session.query(LiveConfig).first()
                if cfg and self.live_manager:
                    new_mute = not cfg.is_muted
                    await self.live_manager.toggle_mute(new_mute)
                    await callback.answer(f"Audio {'Muted' if new_mute else 'Unmuted'}")

            elif data == "live_toggle_autojoin":
                cfg = db_session.query(LiveConfig).first()
                if cfg:
                    cfg.is_auto_join = not cfg.is_auto_join
                    db_session.commit()
                    await callback.answer(f"Auto-Join set to {'ON' if cfg.is_auto_join else 'OFF'}")

            elif data == "menu_services":
                text = (
                    "💼 **Telegram Commercial Services Suite**\n\n"
                    "You can deliver the following services using your loaded accounts:\n\n"
                    "1. 📊 **Poll Voting**: Mass vote on Telegram poll contests.\n"
                    "2. 👁️ **Post Views Booster**: Increase post eye view counter.\n"
                    "3. 🚀 **Mass Joiner**: Boost channel/group subscriber count.\n"
                    "4. 💬 **Live Chat Spammer**: Send custom live comments.\n\n"
                    "👉 _Use the Web Admin Panel Services tab to launch tasks!_"
                )
                kb = InlineKeyboardMarkup([[InlineKeyboardButton("🔙 Back to Main Menu", callback_data="menu_main")]])
                await callback.edit_message_text(text, reply_markup=kb)

            elif data == "menu_reactions":
                rules = db_session.query(ReactionRule).all()
                rule_text = ""
                for r in rules:
                    status = "✅ ON" if r.is_enabled else "❌ OFF"
                    rule_text += f"\n• `{r.target_chat}`: [{r.emojis}] ({status})"
                
                if not rule_text:
                    rule_text = "\n_No reaction rules set yet._"

                text = f"⚡ **Auto Reaction Rules** ({len(rules)}){rule_text}"
                kb = InlineKeyboardMarkup([
                    [InlineKeyboardButton("🔙 Back to Main Menu", callback_data="menu_main")]
                ])
                await callback.edit_message_text(text, reply_markup=kb)

            elif data == "menu_otp":
                accounts = db_session.query(Account).filter(Account.is_active == True).all()
                if not accounts:
                    await callback.answer("❌ No active accounts loaded!", show_alert=True)
                    return

                buttons = []
                for acc in accounts:
                    buttons.append([InlineKeyboardButton(f"📩 {acc.name} ({acc.phone_number})", callback_data=f"read_otp_{acc.id}")])
                buttons.append([InlineKeyboardButton("🔙 Back to Main Menu", callback_data="menu_main")])

                await callback.edit_message_text("📩 **Select Account to View Received OTP Codes:**", reply_markup=InlineKeyboardMarkup(buttons))

            elif data.startswith("read_otp_"):
                account_id = int(data.replace("read_otp_", ""))
                acc = db_session.query(Account).get(account_id)
                success, msg, messages = await self.get_account_otp_messages(account_id)
                
                if not success:
                    await callback.answer(f"❌ {msg}", show_alert=True)
                    return

                otp_text = f"📩 **Telegram OTP & Messages for {acc.name} ({acc.phone_number}):**\n\n"
                if messages:
                    for m in messages:
                        codes_str = ", ".join([f"`{c}`" for c in m["extracted_codes"]]) if m["extracted_codes"] else "None"
                        otp_text += f"⏰ `{m['date']}`\n🔑 **OTP Code**: {codes_str}\n💬 `{m['text'][:120]}`\n\n"
                else:
                    otp_text += "_No recent OTP messages found from Telegram Notification (777000)._"

                kb = InlineKeyboardMarkup([
                    [InlineKeyboardButton("🔄 Refresh OTP", callback_data=f"read_otp_{account_id}")],
                    [InlineKeyboardButton("🔙 Back to OTP Menu", callback_data="menu_otp")]
                ])
                await callback.edit_message_text(otp_text, reply_markup=kb)

            elif data == "menu_stats":
                acc_count = db_session.query(Account).count()
                rule_count = db_session.query(ReactionRule).count()
                log_count = db_session.query(SystemLog).count()
                cfg = db_session.query(LiveConfig).first()

                text = (
                    "📊 **System Status & Stats**\n\n"
                    f"• **Active Userbots**: `{acc_count}`\n"
                    f"• **Reaction Rules**: `{rule_count}`\n"
                    f"• **Total Log Entries**: `{log_count}`\n"
                    f"• **Live Status**: `{cfg.last_status if cfg else 'N/A'}`\n"
                    f"• **Engine Running**: `YES`"
                )
                kb = InlineKeyboardMarkup([[InlineKeyboardButton("🔙 Back to Main Menu", callback_data="menu_main")]])
                await callback.edit_message_text(text, reply_markup=kb)

            elif data == "menu_main":
                await callback.edit_message_text("🎛️ **Admin Control Panel**", reply_markup=self._get_admin_keyboard())

    def _get_admin_keyboard(self):
        return InlineKeyboardMarkup([
            [
                InlineKeyboardButton("🎙️ TG Live (Voice Chat)", callback_data="menu_live"),
                InlineKeyboardButton("⚡ Auto Reactions", callback_data="menu_reactions")
            ],
            [
                InlineKeyboardButton("💼 Commercial Services", callback_data="menu_services"),
                InlineKeyboardButton("📩 Read Account OTPs", callback_data="menu_otp")
            ],
            [
                InlineKeyboardButton("📊 System Stats", callback_data="menu_stats")
            ]
        ])

    async def add_user_session(self, session_string, phone_number=""):
        """Dynamically add a new userbot session string with memory storage and DB session cleanup"""
        db = get_fresh_db()
        api_id, api_hash = get_api_credentials()
        
        clean_session = "".join(session_string.split())

        if not clean_session:
            db.close()
            return False, "Session string cannot be empty"

        try:
            existing = db.query(Account).filter_by(session_string=clean_session).first()
            if existing:
                db.close()
                return True, f"Account already in database: {existing.name}"

            acc = Account(
                session_string=clean_session,
                name="Telegram Account",
                phone_number=phone_number or "Active",
                is_active=True
            )
            db.add(acc)
            db.commit()
            account_id = acc.id
            db.close()

            try:
                user_client = Client(
                    name=f"acc_{account_id}",
                    api_id=api_id,
                    api_hash=api_hash,
                    session_string=clean_session,
                    in_memory=True,
                    workers=1,
                    app_version="10.8.1",
                    device_model="Samsung Galaxy S23",
                    system_version="Android 14",
                    lang_code="en"
                )
                await user_client.connect()
                me = await user_client.get_me()

                db2 = get_fresh_db()
                acc_rec = db2.query(Account).get(account_id)
                if acc_rec:
                    acc_rec.name = f"{me.first_name or ''} {me.last_name or ''}".strip()
                    acc_rec.phone_number = me.phone_number or phone_number or "Active"
                    db2.commit()
                db2.close()

                self.user_clients[account_id] = user_client
                self.reaction_engines.append(AutoReactionEngine(user_client))

                if not self.live_manager:
                    self.live_manager = TGVoiceCallManager(self.user_clients)
                    await self.live_manager.start_all()

                if not self.services_engine:
                    self.services_engine = TelegramServicesEngine(self.user_clients)

                acc_name = f"{me.first_name or ''} {me.last_name or ''}".strip()
                log_system_event(f"Added new Telegram session: {acc_name}")
                return True, f"Successfully logged in and added account: {acc_name} ({me.phone_number or 'Active'})"
            except Exception as conn_err:
                logger.warning(f"Session saved to DB, Pyrogram start notice: {conn_err}")
                return True, f"Account saved to database successfully!"

        except Exception as e:
            db.rollback()
            db.close()
            return False, f"Failed to save account: {str(e)}"
