# admin_bot.py version: 1.3
import asyncio
import logging
import os
import signal
import sqlite3
import zipfile
import tempfile
import random
import shutil
import re

from aiogram import Bot, Dispatcher, F
from aiogram.filters import Command
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
from aiogram.fsm.storage.memory import MemoryStorage
from aiogram.types import (
    InlineKeyboardButton, CallbackQuery, Message, FSInputFile
)
from aiogram.utils.keyboard import InlineKeyboardBuilder
from aiohttp import web
from aiogram.webhook.aiohttp_server import SimpleRequestHandler, setup_application
from pyrogram import Client
from pyrogram.errors import FloodWait, SessionPasswordNeeded

import config
import importlib.util

from shared import (
    _admins_cache, get_setting, set_setting,
    get_speed, get_concurrent, get_work_minutes, get_rest_minutes,
    get_accounts, get_active_count, add_account, del_account, ban_account,
    get_banned, get_bots, add_bot, del_bot,
    get_admins, add_admin, del_admin,
    get_channels, del_channel,
    get_force_sub_channels, get_all_force_sub_urls,
    cleanup_client, clients, stop_event,
    ch_link, norm_target, parse_bot_input,
    rand_delay, safe_sleep, create_tracked_task,
    dead_sessions,
)

import gather_engine as engine

PROFILE_DIR = os.path.expanduser("/var/www/kodo/profile")
PROFILE_GIRL_DIR = os.path.expanduser("/var/www/kodo/profilegirl")
SESSIONS_DIR = os.path.expanduser("/var/www/kodo/sessions")

log = logging.getLogger(__name__)
logging.basicConfig(level=logging.INFO)

def load_module_from_path(path, module_name):
    spec = importlib.util.spec_from_file_location(module_name, path)
    module = importlib.util.module_from_spec(spec)
    spec.loader.exec_module(module)
    return module

try:
    profile_names = load_module_from_path(os.path.join(PROFILE_DIR, "names.py"), "names_boy").names
except Exception:
    profile_names = ["User"]
try:
    profile_bios = load_module_from_path(os.path.join(PROFILE_DIR, "bios.py"), "bios_boy").bios
except Exception:
    profile_bios = [""]
try:
    profile_girl_names = load_module_from_path(os.path.join(PROFILE_GIRL_DIR, "names.py"), "names_girl").names
except Exception:
    profile_girl_names = ["User"]
try:
    profile_girl_bios = load_module_from_path(os.path.join(PROFILE_GIRL_DIR, "bios.py"), "bios_girl").bios
except Exception:
    profile_girl_bios = [""]

import shared
bot = Bot(token=config.BOT_TOKEN)
dp = Dispatcher(storage=MemoryStorage())

def is_auth(uid):
    return uid in config.DEV_IDS or uid in shared._admins_cache

def is_dev(uid):
    return uid in config.DEV_IDS

async def notify(text):
    ids = list(config.DEV_IDS) + shared._admins_cache
    msg_ids = []
    for uid in set(ids):
        try:
            msg = await bot.send_message(uid, text, parse_mode="HTML", disable_web_page_preview=True)
            msg_ids.append((uid, msg.message_id))
        except Exception as e:
            log.warning(f"notify {uid}: {e}")
    return msg_ids


async def notify_edit(msg_ids: list, text: str):
    for chat_id, msg_id in msg_ids:
        try:
            await bot.edit_message_text(chat_id=chat_id, message_id=msg_id, text=text, parse_mode="HTML")
        except Exception as e:
            log.warning(f"notify_edit {chat_id}: {e}")

class States(StatesGroup):
    add_bot = State()
    add_phone = State()
    enter_code = State()
    enter_2fa = State()
    set_speed = State()
    set_concurrent = State()
    add_admin = State()
    set_work_minutes = State()
    set_rest_minutes = State()
    choose_gender = State()
    add_session = State()
    set_profile_count = State()

async def main_kb(is_active):
    b = InlineKeyboardBuilder()
    b.row(
        InlineKeyboardButton(text="اضافة بوت تمويل", callback_data="add_bot"),
        InlineKeyboardButton(text="البوتات المضافة", callback_data="list_bots")
    )
    b.row(InlineKeyboardButton(text="مغادرة قناة", callback_data="leave_ch_pg_0"))
    b.row(InlineKeyboardButton(text="مغادرة كل القنوات", callback_data="leave_all"))
    b.row(InlineKeyboardButton(text="مغادرة الاشتراك الإجباري", callback_data="leave_force_sub"))
    b.row(InlineKeyboardButton(text="اضافة سيشن", callback_data="add_session"))
    b.row(InlineKeyboardButton(text="تعديل الملف الشخصي", callback_data="edit_profiles"))
    b.row(
        InlineKeyboardButton(text="اضافة رقم", callback_data="add_phone"),
        InlineKeyboardButton(text="مسح رقم", callback_data="del_phone_pg_0")
    )
    b.row(InlineKeyboardButton(text="الارقام المضافة", callback_data="phones_pg_0"))

    b.row(InlineKeyboardButton(text="الارقام المحضورة", callback_data="banned_phones_pg_0"))
    b.row(
        InlineKeyboardButton(text="سرعة التجميع", callback_data="set_speed"),
        InlineKeyboardButton(text="الحسابات المتوازية", callback_data="set_concurrent")
    )
    b.row(
        InlineKeyboardButton(text="وقت العمل", callback_data="set_work_minutes"),
        InlineKeyboardButton(text="وقت الاستراحة", callback_data="set_rest_minutes")
    )
    b.row(
        InlineKeyboardButton(text="بدء التجميع", callback_data="start_gather"),
        InlineKeyboardButton(text="ايقاف التجميع", callback_data="stop_gather")
    )
    status = "يعمل ✅" if is_active else "لايعمل ✖️"
    b.row(InlineKeyboardButton(text=f"التجميع: {status}", callback_data="gather_status"))
    b.row(
        InlineKeyboardButton(text="اضافة ادمن", callback_data="add_admin"),
        InlineKeyboardButton(text="مسح ادمن", callback_data="del_admin")
    )
    return b.as_markup()

def back_kb(target="main"):
    b = InlineKeyboardBuilder()
    b.row(InlineKeyboardButton(text="رجوع", callback_data=f"back_{target}"))
    return b.as_markup()

async def get_main_text():
    count = await get_active_count()
    return (
        f"⊱ عزيزي المطور هذه هي لوحة التحكم الخاصة بك.\n"
        f"⊱ تم تحديد السرعة وعدد الحسابات تلقائياً.\n"
        f"⊱ عدد الأرقام الكلي: {count}\n"
        f"⊱ قم بأختيار الاجراء المطلوب بالضغط على الازرار:"
    )


# ============================================================
# Session Import Functions
# ============================================================

def is_telethon_session(file_path: str) -> bool:
    """Check if session file is Telethon format"""
    try:
        import sqlite3
        conn = sqlite3.connect(file_path)
        cursor = conn.cursor()
        cursor.execute("PRAGMA table_info(sessions)")
        columns = [row[1] for row in cursor.fetchall()]
        conn.close()
        # Telethon has 'server_address' column, Pyrogram doesn't
        return 'server_address' in columns
    except:
        return False


def extract_phone_from_session_filename(filename: str) -> str | None:
    """
    Extract phone number from session filename.
    Supports:
    - +96477123456789.session
    - 12075483797_db55882d.session
    - 12082477681_2e288eb9.session
    """
    if filename.endswith('.session'):
        base = filename[:-8]
    elif filename.endswith('.session-journal'):
        return None
    else:
        return None
    
    # Pattern 1: Starts with +
    match = re.search(r'\+([\d]+)', base)
    if match:
        phone = '+' + match.group(1)
        if len(phone) >= 8:
            return phone
    
    # Pattern 2: Digits at start
    match = re.search(r'^(\d{7,15})', base)
    if match:
        digits = match.group(1)
        if len(digits) >= 10:
            return '+' + digits
        return digits
    
    # Pattern 3: Digits anywhere
    match = re.search(r'(\d{7,15})', base)
    if match:
        digits = match.group(1)
        if len(digits) >= 10:
            return '+' + digits
    
    return None


async def process_session_file(file_path: str, filename: str) -> dict:
    """Process a single session file and register it"""
    result = {"phone": None, "status": "unknown", "error": None}
    
    phone = extract_phone_from_session_filename(filename)
    if not phone:
        result["status"] = "invalid_name"
        result["error"] = f"Invalid filename: {filename}"
        return result
    
    result["phone"] = phone
    
    # Check if Telethon session
    if is_telethon_session(file_path):
        result["status"] = "telethon"
        result["error"] = "Telethon session detected. Please convert to Pyrogram first."
        return result
    
    dest_path = os.path.join(SESSIONS_DIR, f"{phone}.session")
    
    try:
        shutil.copy2(file_path, dest_path)
        await add_account(phone)
        result["status"] = "success"
        log.info(f"[SESSION IMPORT] Registered: {phone}")
    except Exception as e:
        result["status"] = "error"
        result["error"] = str(e)
        log.error(f"[SESSION IMPORT] Failed {phone}: {e}")
    
    return result

async def process_zip_sessions(zip_path: str) -> list[dict]:
    """Extract and process session files from zip"""
    results = []
    temp_dir = tempfile.mkdtemp(prefix="sessions_import_")
    
    try:
        with zipfile.ZipFile(zip_path, 'r') as zf:
            session_files = [f for f in zf.namelist() 
                           if f.endswith('.session') 
                           and not f.endswith('.session-journal')]
            
            if not session_files:
                return [{"status": "no_sessions", "error": "No .session files found"}]
            
            zf.extractall(temp_dir)
        
        for session_file in session_files:
            full_path = os.path.join(temp_dir, session_file)
            if os.path.exists(full_path):
                basename = os.path.basename(session_file)
                result = await process_session_file(full_path, basename)
                results.append(result)
        
    except zipfile.BadZipFile:
        results.append({"status": "bad_zip", "error": "Invalid zip file"})
    except Exception as e:
        results.append({"status": "error", "error": str(e)})
    finally:
        shutil.rmtree(temp_dir, ignore_errors=True)
        try:
            if os.path.exists(zip_path):
                os.remove(zip_path)
        except:
            pass
    
    return results


def build_session_report(results: list[dict]) -> str:
    """Build formatted report"""
    total = len(results)
    success = [r for r in results if r["status"] == "success"]
    failed = [r for r in results if r["status"] not in ("success", "telethon")]
    telethon = [r for r in results if r["status"] == "telethon"]
    
    lines = [
        "📦 <b>تقرير استيراد السيشنات</b>",
        f"📊 الإجمالي: {total} ملف",
        f"✅ نجح: {len(success)}",
        f"⚠️ Telethon: {len(telethon)}",
        f"❌ فشل: {len(failed)}",
        ""
    ]
    
    if success:
        lines.append("<b>الحسابات المضافة:</b>")
        for r in success:
            lines.append(f"  ✓ <code>{r['phone']}</code>")
        lines.append("")
    
    if telethon:
        lines.append("<b>ملفات Telethon (غير مدعومة):</b>")
        for r in telethon:
            lines.append(f"  ⚠️ <code>{r['phone']}</code> — تحتاج تحويل لـ Pyrogram")
        lines.append("")
    
    if failed:
        lines.append("<b>الفاشلة:</b>")
        for r in failed:
            phone = r.get('phone') or 'غير معروف'
            error = r.get('error', 'Unknown error')
            lines.append(f"  ✗ <code>{phone}</code> — {error}")
    
    return "\n".join(lines)

# ============================================================
# تعديل الملف الشخصي — نظام مرحلتين: فحص ثم تعديل
# ============================================================

# حالة الفحص المشتركة بين المرحلتين (مستخدم واحد في كل مرة: المطور)
_PROFILE_SCAN = {
    "running": False,      # هل الفحص جارٍ الآن؟
    "cancel": False,       # طلب إيقاف الفحص
    "fake": [],            # قائمة dict لكل حساب وهمي: {phone, need_name, need_bio, need_photo}
    "opened": 0,           # عدد الحسابات التي فُتحت/اتُّصل بها
    "total": 0,            # إجمالي الحسابات
}
# انشغال عمليات الملفات الشخصية (فحص أو تعديل) — يمنع التجميع من فتح
# نفس ملفات الجلسات أثناءها (الفتح المزدوج = خطر AUTH_KEY_DUPLICATED).
_PROFILE_BUSY = False

SCAN_BATCH = 5            # فحص/تحديث الرسالة كل 5 حسابات
SCAN_CONCURRENCY = 5      # توازٍ داخل الدفعة
EDIT_CONCURRENCY = 10     # توازٍ أثناء التعديل
OP_TIMEOUT = 30           # مهلة كل عملية شبكية (ثانية)


def is_weak_name(name: str) -> bool:
    """يعيد True إذا كان الاسم ضعيفاً/وهمياً:
    أقصر من 3 محارف فعلية، أو لا يحتوي أي حرف أبجدي (أرقام/رموز فقط)."""
    if not name:
        return True
    s = name.strip()
    if len(s) < 3:
        return True
    # يجب أن يحتوي على حرف أبجدي واحد على الأقل (لاتيني أو عربي أو أي لغة)
    has_alpha = any(ch.isalpha() for ch in s)
    return not has_alpha


async def _scan_one(phone: str) -> dict | None:
    """يفتح حساباً، يفحص (اسم/بايو/صورة)، ثم يغلقه دائماً.
    يعيد dict للحساب الوهمي، أو None إن كان سليماً/فشل.
    حل جذري: لا يُبقي أي عميل مفتوحاً ولا يضع شيئاً في clients المشترك،
    فلا تراكم ولا حلقات إعادة محاولة عالقة مهما حدث."""
    if phone in dead_sessions:
        return None
    c = None
    try:
        c = Client(
            f"{SESSIONS_DIR}/{phone}",
            api_id=config.API_ID,
            api_hash=config.API_HASH,
            no_updates=True,
        )
        await asyncio.wait_for(c.connect(), timeout=OP_TIMEOUT)
        me = await asyncio.wait_for(c.get_me(), timeout=OP_TIMEOUT)

        # كل حساب يُعدَّل دائماً: اسم + بايو + صورة بغض النظر عن قيمها الحالية
        return {
            "phone": phone,
            "need_name": True,
            "need_bio": True,
            "need_photo": True,
        }
    except asyncio.TimeoutError:
        log.warning(f"[SCAN] {phone}: timeout")
        return None
    except Exception as e:
        err = str(e).lower()
        if "auth_key" in err or "unregistered" in err or "revoked" in err \
                or "eof when reading a line" in err:
            dead_sessions.add(phone)
        log.warning(f"[SCAN] {phone}: {e}")
        return None
    finally:
        # إغلاق مضمون دائماً — لا استثناء، لا إبقاء.
        if c is not None:
            try:
                await asyncio.wait_for(c.disconnect(), timeout=10)
            except Exception:
                pass


def _scan_progress_text() -> str:
    return (
        "⏳ جارٍ الاتصال بالحسابات، انتظر لحظة...\n"
        f"✅ تم فتح حتى الآن: <b>{_PROFILE_SCAN['opened']}</b> حساب\n"
        f"🔍 تم اكتشاف: <b>{len(_PROFILE_SCAN['fake'])}</b> حساب "
        f"لا يحتوي اسم أو صورة أو بايو"
    )


def _scan_cancel_kb():
    b = InlineKeyboardBuilder()
    b.row(InlineKeyboardButton(text="إنهاء", callback_data="prof_scan_cancel"))
    return b.as_markup()


def _scan_result_kb():
    b = InlineKeyboardBuilder()
    b.row(InlineKeyboardButton(text="نعم، عدّل", callback_data="prof_do_edit"))
    b.row(InlineKeyboardButton(text="إلغاء", callback_data="prof_edit_cancel"))
    return b.as_markup()


async def setup_profile(client, phone: str, gender: str = "male",
                        need_name: bool = True, need_bio: bool = True,
                        need_photo: bool = True) -> dict:
    """يملأ الحقول الناقصة فقط (حسب أعلام need_*). يعيد ما عُدّل."""
    result = {"name": None, "bio": None, "photo": None}
    if gender == "female":
        names_list = profile_girl_names
        bios_list = profile_girl_bios
        photos_dir = PROFILE_GIRL_DIR
    else:
        names_list = profile_names
        bios_list = profile_bios
        photos_dir = PROFILE_DIR
    try:
        if need_name:
            name = random.choice(names_list) if names_list else "User"
            for _ in range(3):
                try:
                    await asyncio.wait_for(
                        client.update_profile(first_name=name, last_name=""),
                        timeout=OP_TIMEOUT)
                    result["name"] = name
                    break
                except FloodWait as e:
                    await asyncio.sleep(e.value * 1.2)
                except Exception:
                    await asyncio.sleep(3)

        if need_bio:
            bio = random.choice(bios_list) if bios_list else ""
            for _ in range(3):
                try:
                    await asyncio.wait_for(
                        client.update_profile(bio=bio), timeout=OP_TIMEOUT)
                    result["bio"] = bio
                    break
                except FloodWait as e:
                    await asyncio.sleep(e.value * 1.2)
                except Exception:
                    await asyncio.sleep(3)

        if need_photo:
            photos = [f for f in os.listdir(photos_dir)
                      if f.lower().endswith(".png")] if os.path.isdir(photos_dir) else []
            if photos:
                path = os.path.join(photos_dir, random.choice(photos))
                for _ in range(3):
                    try:
                        # Raw MTProto لتجاوز بق pyrofork 2.3.69 في set_profile_photo
                        import hashlib
                        from pyrogram import raw as _raw
                        with open(path, "rb") as f:
                            data = f.read()
                        chunk_size = 512 * 1024
                        chunks = [data[i:i+chunk_size] for i in range(0, len(data), chunk_size)]
                        file_id = client.rnd_id()
                        for idx, chunk in enumerate(chunks):
                            await client.invoke(
                                _raw.functions.upload.SaveFilePart(
                                    file_id=file_id,
                                    file_part=idx,
                                    bytes=chunk,
                                )
                            )
                        input_file = _raw.types.InputFile(
                            id=file_id,
                            parts=len(chunks),
                            name=os.path.basename(path),
                            md5_checksum=hashlib.md5(data).hexdigest(),
                        )
                        await client.invoke(
                            _raw.functions.photos.UploadProfilePhoto(
                                file=input_file
                            )
                        )
                        result["photo"] = path
                        break
                    except FloodWait as e:
                        await asyncio.sleep(e.value * 1.2)
                    except Exception:
                        await asyncio.sleep(3)
    except Exception:
        pass
    return result


async def _close_scan_clients():
    """تفريغ قائمة الحسابات الوهمية. لم يعد هناك عملاء مفتوحون يُبقَون
    (الحل الجذري يغلق كل عميل فور استخدامه)، لكن نُبقي الدالة احتياطاً
    لإغلاق أي عميل قد يكون بقي في clients لأي سبب."""
    for item in _PROFILE_SCAN["fake"]:
        phone = item["phone"]
        if phone in clients:
            await cleanup_client(phone)
    _PROFILE_SCAN["fake"] = []


# ============================================================
# إضافة رقم: تسجيل دخول حساب جديد (كود + 2FA اختياري + اختيار الصنف)
# ============================================================
async def strip_profile(client, phone: str):
    """يمسح الاسم/النبذة/الصور الحالية للحساب الجديد قبل تعيين ملف جديد."""
    try:
        await client.update_profile(first_name=".", last_name="", bio="")
    except Exception as e:
        log.error(f"[STRIP] name/bio {phone}: {e}")
    try:
        file_ids = []
        async for photo in client.get_chat_photos("me"):
            file_ids.append(photo.file_id)
        if file_ids:
            await client.delete_profile_photos(file_ids)
    except Exception as e:
        log.error(f"[STRIP] photos {phone}: {e}")


async def ask_gender(message: Message, phone: str):
    b = InlineKeyboardBuilder()
    b.row(
        InlineKeyboardButton(text="فتاة", callback_data=f"gender_female_{phone}"),
        InlineKeyboardButton(text="شاب", callback_data=f"gender_male_{phone}")
    )
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    await message.answer(
        f"✓: الرقم: <code>{phone}</code>\n"
        "⊱ عزيزي المطور من فضلك اختر نوع الحساب:",
        reply_markup=b.as_markup(),
        parse_mode="HTML"
    )


async def send_profile_msg(message: Message, phone: str, profile: dict,
                           gender: str = "male"):
    gender_label = "فتاة" if gender == "female" else "شاب"
    caption = (
        f"✓: تم اضافة الحساب: <code>{phone}</code> بنجاح\n"
        f"✓: الاسم: <b>{profile.get('name') or 'لم يُعيَّن'}</b>\n"
        f"✓: النبذة: <b>{profile.get('bio') or 'لم تُعيَّن'}</b>\n"
        f"✓: الصنف: <b>{gender_label}</b>"
    )
    photo = profile.get("photo")
    try:
        if photo and os.path.exists(photo):
            await message.answer_photo(FSInputFile(photo), caption=caption,
                                       parse_mode="HTML")
        else:
            await message.answer(caption, parse_mode="HTML")
    except Exception:
        await message.answer(caption, parse_mode="HTML")


@dp.callback_query(F.data == "add_phone")
async def add_phone_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.add_phone)
    await cb.message.edit_text(
        "⊱ عزيزي المطور ارسل رقم الهاتف مع مفتاح الدولة الان بهذا الشكل:\n"
        ":- +96477xxxxxxxx",
        reply_markup=back_kb())


@dp.message(States.add_phone)
async def proc_phone(msg: Message, state: FSMContext):
    phone = msg.text.replace(" ", "").strip()
    if not phone.startswith("+"):
        await msg.answer("⚠️: يجب ارسال الرقم مع مفتاح الدولة مثال:\n"
                         ":- +96477xxxxxxxx")
        return
    await state.update_data(phone=phone)
    c = Client(f"sessions/{phone}", api_id=config.API_ID,
               api_hash=config.API_HASH, device_model="Desktop")
    await c.connect()
    try:
        code = await c.send_code(phone)
        await state.update_data(code_hash=code.phone_code_hash)
        clients[phone] = c
        await state.set_state(States.enter_code)
        await msg.answer(
            f"✓: الرقم: {phone}\n"
            "←: ارسل الكود المكون من 5 ارقام الذي وصلك الان:",
            reply_markup=back_kb())
    except Exception as e:
        log.error(f"send_code {phone}: {e}")
        await msg.answer(f"خطأ: {e}")
        try:
            await c.disconnect()
        except Exception:
            pass


@dp.message(States.enter_code)
async def proc_code(msg: Message, state: FSMContext):
    code = msg.text.strip()
    data = await state.get_data()
    phone = data["phone"]
    c = clients.get(phone)
    if not c:
        await msg.answer("⚠️: انتهت الجلسة، أعد إضافة الرقم.")
        await state.clear()
        return
    try:
        await c.sign_in(phone, data["code_hash"], code)
        c.me = await c.get_me()
        await strip_profile(c, phone)
        await state.set_state(States.choose_gender)
        await ask_gender(msg, phone)
    except SessionPasswordNeeded:
        await state.set_state(States.enter_2fa)
        await msg.answer(
            "⚠️: الحساب محمي بخاصية التحقق بخطوتين، ارسل الرمز الان:",
            reply_markup=back_kb())
    except Exception as e:
        log.error(f"sign_in {phone}: {e}")
        await msg.answer(f"خطأ: {e}")


@dp.message(States.enter_2fa)
async def proc_2fa(msg: Message, state: FSMContext):
    data = await state.get_data()
    phone = data["phone"]
    c = clients.get(phone)
    if not c:
        await msg.answer("⚠️: انتهت الجلسة، أعد إضافة الرقم.")
        await state.clear()
        return
    try:
        await c.check_password(msg.text)
        c.me = await c.get_me()
        await strip_profile(c, phone)
        await state.set_state(States.choose_gender)
        await ask_gender(msg, phone)
    except Exception as e:
        log.error(f"2fa {phone}: {e}")
        await msg.answer(f"خطأ: {e}")


@dp.callback_query(F.data.startswith("gender_"))
async def gender_chosen(cb: CallbackQuery, state: FSMContext):
    parts = cb.data.split("_", 2)
    gender = parts[1]
    phone = parts[2]
    c = clients.get(phone)
    if not c:
        await cb.answer("⚠️: انتهت الجلسة، أعد إضافة الرقم.", show_alert=True)
        await state.clear()
        return
    try:
        await cb.message.delete()
    except Exception:
        pass
    await add_account(phone)
    profile = await setup_profile(c, phone, gender=gender)
    await send_profile_msg(cb.message, phone, profile, gender=gender)
    # إغلاق العميل بعد الانتهاء (لا نُبقيه مفتوحاً — درس الاستقرار السابق)
    try:
        await cleanup_client(phone)
    except Exception:
        pass
    await state.clear()


@dp.callback_query(F.data == "edit_profiles")
async def edit_profiles_cb(cb: CallbackQuery, state: FSMContext):
    global _PROFILE_BUSY
    is_active = await shared.get_setting("is_active")
    if is_active == "1":
        await cb.answer(
            "⚠️ التجميع يعمل الآن، يرجى إيقافه أولاً قبل تعديل الملفات الشخصية.",
            show_alert=True)
        return
    if _PROFILE_SCAN["running"] or _PROFILE_BUSY:
        await cb.answer("⏳ هناك فحص جارٍ بالفعل.", show_alert=True)
        return

    await state.clear()
    stop_event.clear()  # إعادة ضبط stop_event قبل بدء الفحص
    accounts = await shared.get_accounts()
    total = len(accounts)

    # تهيئة حالة الفحص
    _PROFILE_SCAN.update({
        "running": True, "cancel": False, "fake": [],
        "opened": 0, "total": total,
    })
    _PROFILE_BUSY = True

    status_msg = await cb.message.edit_text(
        _scan_progress_text(), reply_markup=_scan_cancel_kb(), parse_mode="HTML")

    # فحص على دفعات من SCAN_BATCH، بتوازٍ داخل كل دفعة
    for i in range(0, total, SCAN_BATCH):
        if _PROFILE_SCAN["cancel"] or stop_event.is_set():
            break
        batch = accounts[i:i + SCAN_BATCH]
        sem = asyncio.Semaphore(SCAN_CONCURRENCY)

        async def _guarded(phone):
            async with sem:
                return await _scan_one(phone)

        results = await asyncio.gather(
            *(_guarded(a["phone"]) for a in batch),
            return_exceptions=True)

        for r in results:
            _PROFILE_SCAN["opened"] += 1
            if isinstance(r, dict):
                _PROFILE_SCAN["fake"].append(r)

        try:
            await bot.edit_message_text(
                chat_id=status_msg.chat.id,
                message_id=status_msg.message_id,
                text=_scan_progress_text(),
                reply_markup=_scan_cancel_kb(),
                parse_mode="HTML")
        except Exception:
            pass

    _PROFILE_SCAN["running"] = False
    _PROFILE_BUSY = False
    fake_count = len(_PROFILE_SCAN["fake"])

    # حذف رسالة التقدّم وعرض الملخّص
    try:
        await bot.delete_message(status_msg.chat.id, status_msg.message_id)
    except Exception:
        pass

    if fake_count == 0:
        await _close_scan_clients()
        await cb.message.answer(
            "✅ تم فتح والاتصال بجميع الحسابات.\n"
            f"عدد الحسابات الحالي: <b>{_PROFILE_SCAN['opened']}</b> حساب.\n"
            "لا توجد حسابات وهمية تحتاج تعديلاً. 🎉",
            reply_markup=back_kb(), parse_mode="HTML")
        return

    await cb.message.answer(
        "✅ تم فتح والاتصال بجميع الحسابات الحالية.\n"
        f"عدد الحسابات الحالي: <b>{_PROFILE_SCAN['opened']}</b> حساب.\n"
        f"حسابات وهمية غير حقيقية: <b>{fake_count}</b> حساب.\n\n"
        "هل تريد تعديل هذه الحسابات الوهمية الآن؟",
        reply_markup=_scan_result_kb(), parse_mode="HTML")


@dp.callback_query(F.data == "prof_scan_cancel")
async def prof_scan_cancel_cb(cb: CallbackQuery):
    if _PROFILE_SCAN["running"]:
        _PROFILE_SCAN["cancel"] = True
        await cb.answer("⏹️ يتم إنهاء الفحص...", show_alert=False)
    else:
        await cb.answer()


@dp.callback_query(F.data == "prof_edit_cancel")
async def prof_edit_cancel_cb(cb: CallbackQuery):
    await _close_scan_clients()
    try:
        await cb.message.edit_text(
            "تم الإلغاء. لم يُعدّل أي حساب.",
            reply_markup=back_kb())
    except Exception:
        await cb.message.answer("تم الإلغاء.", reply_markup=back_kb())
    await cb.answer()


@dp.callback_query(F.data == "prof_do_edit")
async def prof_do_edit_cb(cb: CallbackQuery):
    global _PROFILE_BUSY
    fake = list(_PROFILE_SCAN["fake"])
    if not fake:
        await cb.answer("لا توجد حسابات للتعديل.", show_alert=True)
        return

    total = len(fake)
    status_msg = await cb.message.edit_text(
        "⏳ تتم عملية تعديل الملف الشخصي للحسابات الوهمية...\n"
        f"✅ تم تعديل: 0 / {total} حساب.",
        parse_mode="HTML")

    done = 0
    sem = asyncio.Semaphore(EDIT_CONCURRENCY)
    lock = asyncio.Lock()
    _PROFILE_BUSY = True

    async def _edit_one(item):
        nonlocal done
        phone = item["phone"]
        async with sem:
            c = None
            try:
                # عميل محلي بالكامل — لا يُسجَّل في clients المشترك إطلاقاً.
                c = Client(
                    f"{SESSIONS_DIR}/{phone}",
                    api_id=config.API_ID,
                    api_hash=config.API_HASH,
                    no_updates=True)
                await asyncio.wait_for(c.connect(), timeout=OP_TIMEOUT)
                await asyncio.wait_for(c.get_me(), timeout=OP_TIMEOUT)
                gender = random.choice(["male", "female"])
                await setup_profile(
                    c, phone, gender=gender,
                    need_name=item["need_name"],
                    need_bio=item["need_bio"],
                    need_photo=item["need_photo"])
            except Exception as e:
                log.warning(f"[EDIT PROFILE] {phone}: {e}")
            finally:
                if c is not None:
                    try:
                        await asyncio.wait_for(c.disconnect(), timeout=10)
                    except Exception:
                        pass
        async with lock:
            done += 1
            if done % SCAN_BATCH == 0 or done == total:
                try:
                    await bot.edit_message_text(
                        chat_id=status_msg.chat.id,
                        message_id=status_msg.message_id,
                        text="⏳ تتم عملية تعديل الملف الشخصي للحسابات الوهمية...\n"
                             f"✅ تم تعديل: {done} / {total} حساب.",
                        parse_mode="HTML")
                except Exception:
                    pass

    try:
        await asyncio.gather(*(_edit_one(it) for it in fake),
                             return_exceptions=True)
    finally:
        _PROFILE_BUSY = False

    _PROFILE_SCAN["fake"] = []
    await cb.message.answer(
        "✅ اكتملت عملية تعديل الملف الشخصي\n"
        f"✓ تم تعديل: {done} من أصل {total} حساب وهمي.",
        reply_markup=back_kb(), parse_mode="HTML")
    await cb.answer()



def build_leave_kb(channels, page):
    start = page * 10
    current = channels[start:start + 10]
    b = InlineKeyboardBuilder()
    for i, ch in enumerate(current):
        idx = start + i
        link = ch_link(ch)
        title = (ch.get("title") or link or "قناة")[:30]
        btn = InlineKeyboardButton(text=title, url=link) if link else InlineKeyboardButton(text=title, callback_data="none")
        b.row(btn, InlineKeyboardButton(text="مغادرة", callback_data=f"leave_exec_{idx}"))
    nav = []
    if page > 0:
        nav.append(InlineKeyboardButton(text="← السابق", callback_data=f"leave_ch_pg_{page - 1}"))
    if len(channels) > start + 10:
        nav.append(InlineKeyboardButton(text="التالي →", callback_data=f"leave_ch_pg_{page + 1}"))
    if nav:
        b.row(*nav)
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    return b.as_markup(), current

def build_del_phone_kb(accounts, page):
    start = page * 10
    current = accounts[start:start + 10]
    b = InlineKeyboardBuilder()
    for acc in current:
        b.row(
            InlineKeyboardButton(text=acc["phone"], callback_data="none"),
            InlineKeyboardButton(text="🗑️", callback_data=f"del_ph_{acc['phone']}_{page}")
        )
    nav = []
    if page > 0:
        nav.append(InlineKeyboardButton(text="← السابق", callback_data=f"del_phone_pg_{page - 1}"))
    if len(accounts) > start + 10:
        nav.append(InlineKeyboardButton(text="التالي →", callback_data=f"del_phone_pg_{page + 1}"))
    if nav:
        b.row(*nav)
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    return b.as_markup(), current

async def render_del_phones(message: Message, page: int):
    accounts = await get_accounts()
    markup, current = build_del_phone_kb(accounts, page)
    if not current and page == 0:
        await message.edit_text("عزيزي المطور لاتوجد ارقام مضافة حالياً.", reply_markup=back_kb())
        return
    if not current:
        raise ValueError("عزيزي المطور لاتوجد ارقام مضافة حالياً.")
    await message.edit_text(
        f"قائمة {page + 1} | {len(current)} رقم:",
        reply_markup=markup
    )

@dp.callback_query(F.data == "back_main")
async def back_main(cb: CallbackQuery, state: FSMContext):
    await state.clear()
    active = await get_setting("is_active", "0") == "1"
    await cb.message.edit_text(await get_main_text(), reply_markup=await main_kb(active))

@dp.callback_query(F.data == "back_list_bots")
async def back_list_bots(cb: CallbackQuery):
    await list_bots_cb(cb)

@dp.callback_query(F.data == "back_del_admin")
async def back_del_admin(cb: CallbackQuery):
    await del_admin_list(cb)

@dp.callback_query(F.data == "none")
async def none_cb(cb: CallbackQuery):
    await cb.answer()

@dp.message(Command("start"))
async def cmd_start(msg: Message):
    if is_auth(msg.from_user.id):
        active = await get_setting("is_active", "0") == "1"
        await msg.answer(await get_main_text(), reply_markup=await main_kb(active))
    else:
        await msg.answer(
            "هذا البوت مخصص للمطورين والأدمنية فقط.\n"
            "للاستفسار تواصل مع المطور:\n• @n_u_7"
        )

@dp.callback_query(F.data == "add_bot")
async def add_bot_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.add_bot)
    await cb.message.edit_text("⊱ عزيزي المطور ارسل يوزر او رابط البوت المطلوب لاضافته الى القائمة:", reply_markup=back_kb())

@dp.message(States.add_bot)
async def proc_add_bot(msg: Message, state: FSMContext):
    parsed = parse_bot_input(msg.text.strip())
    username = parsed.get("username")
    if not username:
        await msg.answer("⊱ اليوزر أو الرابط غير صحيح، ارسل اليوزر بهذا الشكل:\n- @usernamebot", reply_markup=back_kb())
        return
    bots_list = await get_bots()
    if any(b["username"].lower() == username.lower() for b in bots_list):
        await msg.answer(f"⚠️ عزيزي هذا البوت @{username} مضاف مسبقاً.", reply_markup=back_kb())
        await state.clear()
        return
    await add_bot(username, parsed.get("url"))
    await msg.answer(f"✓: تم إضافة البوت @{username} الى القائمة بنجاح.", reply_markup=back_kb())
    await state.clear()

@dp.callback_query(F.data == "list_bots")
async def list_bots_cb(cb: CallbackQuery):
    bots_list = await get_bots()
    if not bots_list:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد بوتات مضافة حالياً.", reply_markup=back_kb())
        return
    b = InlineKeyboardBuilder()
    for bot_item in bots_list:
        u = bot_item["username"]
        url = bot_item.get("url") or f"https://t.me/{u}"
        b.row(
            InlineKeyboardButton(text=f"@{u}", url=url),
            InlineKeyboardButton(text="🗑️", callback_data=f"del_bot_{u}")
        )
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    await cb.message.edit_text("⊱ عزيزي المطور هذه هي قائمة بوتات التمويل المضافة:", reply_markup=b.as_markup())

@dp.callback_query(F.data.startswith("del_bot_"))
async def del_bot_cb(cb: CallbackQuery):
    username = cb.data[len("del_bot_"):]
    await del_bot(username)
    await cb.answer("✓: تم حذف البوت المطلوب بنجاح.")
    await list_bots_cb(cb)

@dp.callback_query(F.data.startswith("leave_ch_pg_"))
async def leave_ch_pg(cb: CallbackQuery):
    page = int(cb.data.split("_")[-1])
    channels = await get_channels()
    if not channels:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد قنوات في قائمة المغادرة حالياً.", reply_markup=back_kb())
        return
    markup, current = build_leave_kb(channels, page)
    if not current:
        await cb.answer("⊱ عزيزي المطور لاتوجد قنوات في هذه الصفحة حالياً.", show_alert=True)
        return
    await cb.message.edit_text(
        f"قائمة القنوات (صفحة {page + 1}) | المجموع: {len(channels)}\n"
        "اضغط مغادرة لإزالة القناة من جميع الحسابات:",
        reply_markup=markup
    )

@dp.callback_query(F.data.startswith("leave_exec_"))
async def leave_exec(cb: CallbackQuery):
    idx = int(cb.data.split("_")[-1])
    channels = await get_channels()
    if idx >= len(channels):
        await cb.answer("⊱ عزيزي المطور هذه القناة غير موجودة.", show_alert=True)
        return
    ch = channels[idx]
    link = ch_link(ch)
    title = ch.get("title") or link or "القناة"
    ch_html = f'<a href="{link}">{title}</a>' if link else title
    b = InlineKeyboardBuilder()
    b.row(InlineKeyboardButton(text="رجوع للقائمة", callback_data="leave_ch_pg_0"))
    await cb.message.edit_text(
        f"⏳: تتم الان عملية مغادرة القناة: {ch_html} سيتم ابلاغك عند انتهاء المغادرة.",
        reply_markup=b.as_markup(), parse_mode="HTML", disable_web_page_preview=True
    )
    create_tracked_task(engine.do_leave_ch(ch, cb.from_user.id))

@dp.callback_query(F.data == "leave_all")
async def leave_all_cb(cb: CallbackQuery):
    channels = await get_channels()
    b = InlineKeyboardBuilder()
    if shared.leave_all_on:
        b.row(InlineKeyboardButton(text="⛔ ايقاف المغادرة", callback_data="stop_leave_all"))
        b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
        await cb.message.edit_text(
            f"⚠️ عملية المغادرة جارية! ({len(channels)} قناة متبقية)",
            reply_markup=b.as_markup()
        )
    else:
        b.row(InlineKeyboardButton(text="تأكيد", callback_data="confirm_leave_all"))
        b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
        await cb.message.edit_text(
            f"🔻: سيتم مغادرة {len(channels)} قناة هل انت متأكد من عملية المغادرة؟",
            reply_markup=b.as_markup()
        )

@dp.callback_query(F.data == "confirm_leave_all")
async def confirm_leave_all(cb: CallbackQuery):
    if shared.leave_all_on:
        await cb.answer("⚠️: عزيزي المطور العملية الحالية تعمل بالفعل.", show_alert=True)
        return
    shared.leave_all_on = True
    try:
        await cb.message.edit_text(
            "⏳: تم بدء عملية مغادرة جميع القنوات بنظام الدفعات، سيتم ابلاغك عند الانتهاء.",
            reply_markup=back_kb()
        )
    except Exception as e:
        log.warning(f"Could not edit message: {e}")
    try:
        await cb.answer("🚀 بدأت العملية بنجاح")
    except Exception:
        pass
    shared.leave_all_task = create_tracked_task(engine.do_leave_all(cb.from_user.id))

@dp.callback_query(F.data == "leave_force_sub")
async def leave_force_sub_cb(cb: CallbackQuery):
    bots = await shared.get_bots()
    if not bots:
        await cb.answer("⊱ لا يوجد بوتات مضافة.", show_alert=True)
        return
    b = InlineKeyboardBuilder()
    for bot in bots:
        un = bot.get("username") or ""
        b.row(InlineKeyboardButton(
            text=f"@{un}",
            callback_data=f"lfs_bot_{un}"
        ))
    b.row(InlineKeyboardButton(text="↩️ رجوع", callback_data="back_main"))
    await cb.message.edit_text(
        "🔘 اختر البوت لعرض قنوات الاشتراك الإجباري:",
        reply_markup=b.as_markup()
    )


LFS_PAGE_SIZE = 15


@dp.callback_query(F.data.startswith("lfs_bot_"))
async def lfs_bot_selected_cb(cb: CallbackQuery):
    import urllib.parse
    raw = cb.data.removeprefix("lfs_bot_")
    if "|" in raw:
        bot_un, page_s = raw.split("|", 1)
        try:
            page = max(int(page_s), 0)
        except ValueError:
            page = 0
    else:
        bot_un, page = raw, 0
    rows = sorted(
        await get_force_sub_channels(bot_un),
        key=lambda r: r["url"],
    )
    if not rows:
        b = InlineKeyboardBuilder()
        b.row(InlineKeyboardButton(text="↩️ رجوع", callback_data="leave_force_sub"))
        await cb.message.edit_text(
            f"⊱ البوت @{bot_un} لا يملك قنوات اشتراك إجباري مسجّلة.",
            reply_markup=b.as_markup()
        )
        return
    total_pages = (len(rows) + LFS_PAGE_SIZE - 1) // LFS_PAGE_SIZE
    page = min(page, total_pages - 1)
    start = page * LFS_PAGE_SIZE
    current = rows[start:start + LFS_PAGE_SIZE]
    b = InlineKeyboardBuilder()
    for r in current:
        url = r["url"]
        enc = urllib.parse.quote(url, safe="")
        b.row(
            InlineKeyboardButton(text=url, url=url),
            InlineKeyboardButton(text="🚪 مغادرة", callback_data=f"lfs_one_{bot_un}|{enc}")
        )
    if total_pages > 1:
        nav = []
        if page > 0:
            nav.append(InlineKeyboardButton(
                text="◀️ السابق",
                callback_data=f"lfs_bot_{bot_un}|{page - 1}"
            ))
        nav.append(InlineKeyboardButton(
            text=f"صفحة {page + 1}/{total_pages}",
            callback_data="lfs_noop"
        ))
        if page < total_pages - 1:
            nav.append(InlineKeyboardButton(
                text="التالي ▶️",
                callback_data=f"lfs_bot_{bot_un}|{page + 1}"
            ))
        b.row(*nav)
    b.row(InlineKeyboardButton(text="🚪 مغادرة الكل", callback_data=f"lfs_confirm_{bot_un}"))
    b.row(InlineKeyboardButton(text="↩️ رجوع", callback_data="leave_force_sub"))
    await cb.message.edit_text(
        f"🤖 البوت: @{bot_un}\n"
        f"📋 عدد القنوات: {len(rows)} قناة\n"
        f"صفحة {page + 1}/{total_pages}\n"
        "اضغط 'مغادرة' لقناة معينة أو 'مغادرة الكل' لمغادرة كل القنوات الموجودة:",
        reply_markup=b.as_markup()
    )


@dp.callback_query(F.data == "lfs_noop")
async def lfs_noop_cb(cb: CallbackQuery):
    await cb.answer()


@dp.callback_query(F.data.startswith("lfs_one_"))
async def lfs_one_cb(cb: CallbackQuery):
    import urllib.parse
    data = cb.data.removeprefix("lfs_one_")
    bot_un, enc_url = data.split("|", 1)
    url = urllib.parse.unquote(enc_url)
    if shared.leave_all_on:
        await cb.answer("⚠️ عملية مغادرة جارية، انتظر حتى تنتهي.", show_alert=True)
        return
    await cb.message.edit_text(
        f"⏳ جاري مغادرة القناة:\n{url}",
        reply_markup=back_kb()
    )
    try:
        await cb.answer()
    except Exception:
        pass
    create_tracked_task(engine.do_leave_force_sub(cb.from_user.id, bot_un, [url]))


@dp.callback_query(F.data.startswith("lfs_confirm_"))
async def lfs_confirm_cb(cb: CallbackQuery):
    bot_un = cb.data.removeprefix("lfs_confirm_")
    rows = await get_force_sub_channels(bot_un)
    if not rows:
        await cb.answer("⊱ لا توجد قنوات لهذا البوت.", show_alert=True)
        return
    if shared.leave_all_on:
        await cb.answer("⚠️ عملية مغادرة جارية، انتظر حتى تنتهي.", show_alert=True)
        return
    urls = [r["url"] for r in rows]
    await cb.message.edit_text(
        f"⏳ جاري مغادرة {len(urls)} قناة اشتراك إجباري للبوت @{bot_un}...",
        reply_markup=back_kb()
    )
    try:
        await cb.answer()
    except Exception:
        pass
    create_tracked_task(engine.do_leave_force_sub(cb.from_user.id, bot_un, urls))


@dp.callback_query(F.data == "stop_leave_all")
async def stop_leave_all(cb: CallbackQuery):
    if not shared.leave_all_on:
        await cb.answer("⊱ عزيزي المطور حالياً لاتوجد عملية مغادرة جارية!", show_alert=True)
        return
    if shared.leave_all_task and not shared.leave_all_task.done():
        shared.leave_all_task.cancel()
    shared.leave_all_on = False
    shared.leave_all_task = None
    await cb.message.edit_text("✓: عزيزي تم ايقاف عملية المغادرة بنجاح.", reply_markup=back_kb())
    await notify("⏸️ تم إيقاف عملية مغادرة كل القنوات.")

@dp.callback_query(F.data == "add_session")
async def add_session_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.add_session)
    await cb.message.edit_text(
        "📦 <b>استيراد سيشنات</b>\n\n"
        "⊱ ارسل ملف <b>.zip</b> يحتوي على ملفات <b>.session</b>\n"
        "⊱ أو ملف <b>.session</b> واحد مباشرة\n\n"
        "⚠️ يجب أن يحتوي اسم الملف على الرقم:\n"
        "<code>+96477123456789.session</code>\n"
        "<code>12075483797_abc123.session</code>",
        reply_markup=back_kb(),
        parse_mode="HTML"
    )


@dp.message(States.add_session)
async def proc_add_session(msg: Message, state: FSMContext):
    if not msg.document:
        await msg.answer(
            "⚠️ أرسل ملف (.zip أو .session) وليس نص.",
            reply_markup=back_kb()
        )
        return
    
    document = msg.document
    file_name = document.file_name or ""
    
    temp_path = os.path.join(
        tempfile.gettempdir(), 
        f"session_import_{msg.from_user.id}_{file_name}"
    )
    
    try:
        file = await bot.get_file(document.file_id)
        await bot.download_file(file.file_path, temp_path)
        
        results = []
        
        if file_name.endswith('.zip'):
            results = await process_zip_sessions(temp_path)
            
        elif file_name.endswith('.session'):
            result = await process_session_file(temp_path, file_name)
            results.append(result)
            try:
                os.remove(temp_path)
            except:
                pass
        else:
            await msg.answer(
                "⚠️ صيغة غير مدعومة. أرسل .zip أو .session فقط.",
                reply_markup=back_kb()
            )
            try:
                os.remove(temp_path)
            except:
                pass
            return
        
        report = build_session_report(results)
        await msg.answer(report, reply_markup=back_kb(), parse_mode="HTML")
        
        active = await get_setting("is_active", "0") == "1"
        await msg.answer(await get_main_text(), reply_markup=await main_kb(active))
        
        await state.clear()
        
    except Exception as e:
        log.error(f"[SESSION IMPORT] Error: {e}")
        await msg.answer(
            f"❌ خطأ: <code>{e}</code>",
            reply_markup=back_kb(),
            parse_mode="HTML"
        )
        try:
            if os.path.exists(temp_path):
                os.remove(temp_path)
        except:
            pass
        await state.clear()



@dp.callback_query(F.data.startswith("del_phone_pg_"))
async def del_phone_pg(cb: CallbackQuery):
    page = int(cb.data.split("_")[-1])
    try:
        await render_del_phones(cb.message, page)
    except ValueError as e:
        await cb.answer(str(e), show_alert=True)

@dp.callback_query(F.data.startswith("del_ph_"))
async def del_ph_exec(cb: CallbackQuery):
    parts = cb.data.split("_")
    phone = parts[2]
    page = int(parts[3])
    await del_account(phone)
    for path in [f"sessions/{phone}.session", f"sessions/{phone}.session-journal"]:
        if os.path.exists(path):
            os.remove(path)
    await cleanup_client(phone)
    await cb.answer(f"✓: تم حذف الرقم {phone} بنجاح.")
    await render_del_phones(cb.message, page)

@dp.callback_query(F.data.startswith("phones_pg_"))
async def phones_pg(cb: CallbackQuery):
    page = int(cb.data.split("_")[-1])
    accounts = await get_accounts()
    total = len(accounts)
    start = page * 10
    current = accounts[start:start + 10]
    if not current and page == 0:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد ارقام مضافة حالياً في القائمة.", reply_markup=back_kb())
        return
    if not current:
        await cb.answer(f"القائمة {page + 1} فارغة.", show_alert=True)
        return
    b = InlineKeyboardBuilder()
    for acc in current:
        b.row(InlineKeyboardButton(text=f"✅ {acc['phone']}", callback_data="none"))
    nav = []
    if page > 0:
        nav.append(InlineKeyboardButton(text="← السابق", callback_data=f"phones_pg_{page - 1}"))
    if len(accounts) > start + 10:
        nav.append(InlineKeyboardButton(text="التالي →", callback_data=f"phones_pg_{page + 1}"))
    if nav:
        b.row(*nav)
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    await cb.message.edit_text(f"إجمالي الأرقام: {total} | قائمة {page + 1} | {len(current)} رقم:", reply_markup=b.as_markup())

@dp.callback_query(F.data.startswith("banned_phones_pg_"))
async def banned_phones_pg(cb: CallbackQuery):
    page = int(cb.data.split("_")[-1])
    banned = await get_banned()
    total = len(banned)
    start = page * 10
    current = banned[start:start + 10]
    if not current and page == 0:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد ارقام محضورة حالياً في القائمة.", reply_markup=back_kb())
        return
    if not current:
        await cb.answer(f"القائمة {page + 1} فارغة.", show_alert=True)
        return
    b = InlineKeyboardBuilder()
    for acc in current:
        b.row(
            InlineKeyboardButton(text=acc["phone"], callback_data="none"),
            InlineKeyboardButton(text="🗑️", callback_data=f"del_banned_{acc['phone']}")
        )
    nav = []
    if page > 0:
        nav.append(InlineKeyboardButton(text="← السابق", callback_data=f"banned_phones_pg_{page - 1}"))
    if total > start + 10:
        nav.append(InlineKeyboardButton(text="التالي →", callback_data=f"banned_phones_pg_{page + 1}"))
    if nav:
        b.row(*nav)
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    await cb.message.edit_text(
        f"⊱ عزيزي المطور هذه هي قائمة الارقام المحضورة:\nالإجمالي: {total} | قائمة {page + 1} | {len(current)} رقم:",
        reply_markup=b.as_markup()
    )

@dp.callback_query(F.data.startswith("del_banned_"))
async def del_banned_exec(cb: CallbackQuery):
    phone = cb.data.rsplit("_", 1)[-1]
    import shared as _shared

    # ← حذف من banned_accounts في PostgreSQL
    async with _shared._pg_pool.acquire() as con:
        await con.execute(
            "DELETE FROM banned_accounts WHERE phone=$1", phone
        )

    # ← حذف ملفات الجلسة إن وجدت
    for path in [
        f"sessions/{phone}.session",
        f"sessions/{phone}.session-journal"
    ]:
        try:
            if os.path.exists(path):
                os.remove(path)
                log.info(f"[DEL BANNED] حُذف ملف الجلسة: {path}")
        except Exception as e:
            log.error(f"[DEL BANNED] فشل حذف {path}: {e}")

    # ← تنظيف الـ client إن كان مفتوحاً
    await cleanup_client(phone)

    await cb.answer(f"✓ تم حذف الرقم {phone} نهائياً من كل مكان.")
    await banned_phones_pg(cb)


@dp.callback_query(F.data == "set_speed")
async def set_speed_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.set_speed)
    spd = await get_speed()
    await cb.message.edit_text(
        f"السرعة الحالية: {spd} ثانية بين كل قناة وأخرى.\n"
        "⚠️ الموصى به: 5 ثواني\n\nأرسل العدد بالثواني:",
        reply_markup=back_kb()
    )

@dp.message(States.set_speed)
async def proc_speed(msg: Message, state: FSMContext):
    if not msg.text.isdigit():
        await msg.answer("يرجى إرسال رقم صحيح.")
        return
    await set_setting("gathering_speed", msg.text)
    await msg.answer(f"✓: تم تحديد السرعة: {msg.text} ثانية.", reply_markup=back_kb())
    await state.clear()

@dp.callback_query(F.data == "set_concurrent")
async def set_concurrent_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.set_concurrent)
    conc = await get_concurrent()
    await cb.message.edit_text(
        f"عدد الحسابات المتوازية الحالي: {conc}\n"
        "⚠️ الموصى به: 20 | الحد الأقصى المقترح: 30\n\nأرسل العدد المطلوب:",
        reply_markup=back_kb()
    )

@dp.message(States.set_concurrent)
async def proc_concurrent(msg: Message, state: FSMContext):
    if not msg.text.isdigit() or int(msg.text) < 1:
        await msg.answer("يرجى إرسال رقم صحيح أكبر من 0.")
        return
    await set_setting("concurrent", msg.text)
    await msg.answer(f"✓: تم تحديد الحسابات المتوازية: {msg.text} حساب.", reply_markup=back_kb())
    await state.clear()

@dp.callback_query(F.data == "set_work_minutes")
async def set_work_minutes_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.set_work_minutes)
    wm = await get_work_minutes()
    await cb.message.edit_text(
        f"وقت العمل الحالي: {wm} دقيقة.\n"
        "⚠️ الموصى به: 15 دقيقة\n\nأرسل المدة بالدقائق:",
        reply_markup=back_kb()
    )

@dp.message(States.set_work_minutes)
async def proc_work_minutes(msg: Message, state: FSMContext):
    if not msg.text.isdigit() or int(msg.text) < 1:
        await msg.answer("يرجى إرسال رقم صحيح أكبر من 0.")
        return
    await set_setting("work_minutes", msg.text)
    await msg.answer(f"✓: تم تحديد وقت العمل: {msg.text} دقيقة.", reply_markup=back_kb())
    await state.clear()

@dp.callback_query(F.data == "set_rest_minutes")
async def set_rest_minutes_cb(cb: CallbackQuery, state: FSMContext):
    await state.set_state(States.set_rest_minutes)
    rm = await get_rest_minutes()
    await cb.message.edit_text(
        f"وقت الاستراحة الحالي: {rm} دقيقة.\n"
        "⚠️ الموصى به: 5 دقائق\n\nأرسل المدة بالدقائق:",
        reply_markup=back_kb()
    )

@dp.message(States.set_rest_minutes)
async def proc_rest_minutes(msg: Message, state: FSMContext):
    if not msg.text.isdigit() or int(msg.text) < 1:
        await msg.answer("يرجى إرسال رقم صحيح أكبر من 0.")
        return
    await set_setting("rest_minutes", msg.text)
    await msg.answer(f"✓: تم تحديد وقت الاستراحة: {msg.text} دقيقة.", reply_markup=back_kb())
    await state.clear()

@dp.callback_query(F.data == "start_gather")
async def start_gather(cb: CallbackQuery):
    accounts = await get_accounts(active_only=True)
    bots_list = await get_bots()
    if not accounts:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد حسابات مضافة نشطة حالياً.", reply_markup=back_kb())
        return
    if not bots_list:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد بوتات تمويل مضافة الى القائمة حالياً.", reply_markup=back_kb())
        return
    b = InlineKeyboardBuilder()
    for bot_item in bots_list:
        u = bot_item["username"]
        b.row(
            InlineKeyboardButton(text=f"@{u}", callback_data="none"),
            InlineKeyboardButton(text="بدء", callback_data=f"start_gather_bot_{u}")
        )
    b.row(InlineKeyboardButton(text="بدء لكل البوتات", callback_data="start_gather_all_bots"))
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    await cb.message.edit_text(
        "عزيزي المطور يجب اختيار بوت قبل البدء في عملية التجميع:",
        reply_markup=b.as_markup()
    )

@dp.callback_query(F.data.startswith("start_gather_bot_"))
async def start_gather_bot(cb: CallbackQuery):
    bot_un = cb.data[len("start_gather_bot_"):]
    if _PROFILE_SCAN["running"] or _PROFILE_BUSY:
        await cb.answer(
            "⚠️ فحص/تعديل الملفات الشخصية يعمل الآن — انتظر انتهاءَه أولاً.",
            show_alert=True)
        return
    if await get_setting("is_active", "0") == "1":
        active_bot = await get_setting("active_bot", "")
        await cb.answer(f"⚠️: التجميع يعمل بالفعل على: @{active_bot}", show_alert=True)
        return
    accounts = await get_accounts(active_only=True)
    if not accounts:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد حسابات مضافة نشطة حالياً.", reply_markup=back_kb())
        return
    stop_event.clear()
    await set_setting("is_active", "1")
    await set_setting("active_bot", bot_un)
    conc = await get_concurrent()
    spd = await get_speed()
    await cb.message.edit_text(
        f"✓: تم بدء التجميع على البوت @{bot_un} الان.\n"
        f"←: الحسابات المتوازية: {conc} حساب.\n"
        f"←: السرعة: {spd} ثانية بين كل قناة واخرى.",
        reply_markup=back_kb()
    )
    # نمرر username فقط — gather_engine يبني الـ url داخلياً
    create_tracked_task(engine.run_gather(bot_un))

@dp.callback_query(F.data == "start_gather_all_bots")
async def start_gather_all_bots_cb(cb: CallbackQuery):
    if _PROFILE_SCAN["running"] or _PROFILE_BUSY:
        await cb.answer(
            "⚠️ فحص/تعديل الملفات الشخصية يعمل الآن — انتظر انتهاءَه أولاً.",
            show_alert=True)
        return
    if await get_setting("is_active", "0") == "1":
        active_bot = await get_setting("active_bot", "")
        await cb.answer(f"⚠️: التجميع يعمل بالفعل على: @{active_bot}", show_alert=True)
        return
    accounts = await get_accounts(active_only=True)
    bots_list = await get_bots()
    if not accounts:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد حسابات مضافة نشطة حالياً.", reply_markup=back_kb())
        return
    if not bots_list:
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد بوتات تمويل مضافة الى القائمة حالياً.", reply_markup=back_kb())
        return
    stop_event.clear()
    # ← إزالة shared.all_bots_on = True (لم يعد مستخدماً)
    await set_setting("is_active", "1")
    await set_setting("active_bot", "ALL")
    wm = await get_work_minutes()
    rm = await get_rest_minutes()
    bots_list_str = ", ".join([f"@{b['username']}" for b in bots_list])
    await cb.message.edit_text(
        f"✓: تم بدء التجميع على جميع البوتات الان.\n"
        f"←: البوتات: {bots_list_str}\n"
        f"←: الحسابات: {len(accounts)}\n"
        f"←: وقت العمل: {wm} دقيقة | وقت الاستراحة: {rm} دقيقة.",
        reply_markup=back_kb()
    )
    # ← نحفظ الـ task في shared.all_bots_task للإيقاف
    shared.all_bots_task = create_tracked_task(engine.run_gather_all_bots())

@dp.callback_query(F.data == "stop_gather")
async def stop_gather(cb: CallbackQuery):
    if await get_setting("is_active", "0") != "1":
        await cb.message.edit_text("⊱ عزيزي المطور لاتوجد عملية تجميع تعمل حالياً.", reply_markup=back_kb())
        return
    # ← stop_event يوقف كل الـ Workers تلقائياً
    stop_event.set()
    await set_setting("is_active", "0")
    await set_setting("active_bot", "")
    # إلغاء task كل البوتات إن وجد
    if shared.all_bots_task and not shared.all_bots_task.done():
        shared.all_bots_task.cancel()
    shared.all_bots_task = None
    shared.all_bots_on = False
    await cb.message.edit_text("✓: عزيزي تم ايقاف عملية التجميع بنجاح.", reply_markup=back_kb())

@dp.callback_query(F.data == "gather_status")
async def gather_status(cb: CallbackQuery):
    active = await get_setting("is_active", "0") == "1"
    status = "يعمل ✅" if active else "لايعمل ✖️"
    active_bot = await get_setting("active_bot", "")
    bot_info = f"\nالبوت النشط: @{active_bot}" if active and active_bot else ""
    conc = await get_concurrent()
    spd = await get_speed()
    await cb.answer(
        f"التجميع: {status}{bot_info}\nمتوازي: {conc} حساب\nالسرعة: {spd}ث",
        show_alert=True
    )

@dp.callback_query(F.data == "add_admin")
async def add_admin_cb(cb: CallbackQuery, state: FSMContext):
    if not is_dev(cb.from_user.id):
        await cb.answer("⚠️: عزيزي هذه الاجراء مخصص للمطورين فقط.", show_alert=True)
        return
    await state.set_state(States.add_admin)
    await cb.message.edit_text("⊱ حسناً ارسل أيدي المستخدم المطلوب لاضافته كأدمن في البوت:", reply_markup=back_kb())

@dp.message(States.add_admin)
async def proc_add_admin(msg: Message, state: FSMContext):
    if not msg.text.isdigit():
        await msg.answer("⚠️: حدث خطأ، يجب ارسال أيدي صحيح.")
        return
    await add_admin(int(msg.text))
    await msg.answer(f"✓: تم اضافة الادمن {msg.text} بنجاح.", reply_markup=back_kb())
    await state.clear()

@dp.callback_query(F.data == "del_admin")
async def del_admin_list(cb: CallbackQuery):
    if not is_dev(cb.from_user.id):
        await cb.answer("⚠️: عزيزي هذه الاجراء مخصص للمطورين فقط.", show_alert=True)
        return
    admins = await get_admins()
    if not admins:
        await cb.message.edit_text("⊱ عزيزي المطور لايوجد ادمنية مضافة حالياً في البوت.", reply_markup=back_kb())
        return
    b = InlineKeyboardBuilder()
    for uid in admins:
        try:
            user = await bot.get_chat(uid)
            name = f"@{user.username}" if user.username else str(uid)
            b.row(
                InlineKeyboardButton(text=name, url=f"tg://user?id={uid}"),
                InlineKeyboardButton(text="🗑️", callback_data=f"del_adm_{uid}")
            )
        except Exception:
            b.row(
                InlineKeyboardButton(text=str(uid), callback_data="none"),
                InlineKeyboardButton(text="🗑️", callback_data=f"del_adm_{uid}")
            )
    b.row(InlineKeyboardButton(text="رجوع", callback_data="back_main"))
    await cb.message.edit_text("⊱ عزيزي المطور هذه هي قائمة الادمنية المضافة في البوت:", reply_markup=b.as_markup())

@dp.callback_query(F.data.startswith("del_adm_"))
async def del_adm_exec(cb: CallbackQuery):
    uid = int(cb.data.split("_")[2])
    await del_admin(uid)
    await cb.answer("✓: تم حذف الادمن بنجاح.")
    await del_admin_list(cb)

def handle_exception(loop, context):
    msg = context.get("exception", context["message"])
    msg_str = str(msg)
    ignored = [
        "Peer id invalid", "ID not found", "Request timed out",
        "Connection lost", "Broken pipe", "TimeoutError",
    ]
    if any(x in msg_str for x in ignored):
        return
    loop.default_exception_handler(context)

def shutdown_handler():
    log.info("[SHUTDOWN] إيقاف السكربت")
    # ← استخدام psycopg2 بدل sqlite3
    import psycopg2
    try:
        con = psycopg2.connect(shared.PG_DSN)
        cur = con.cursor()
        cur.execute("UPDATE settings SET value='0' WHERE key='is_active'")
        cur.execute("UPDATE settings SET value='' WHERE key='active_bot'")
        con.commit()
        con.close()
        log.info("[SHUTDOWN] تم إيقاف التجميع في قاعدة البيانات")
    except Exception as e:
        log.error(f"[SHUTDOWN] {e}")

async def on_startup(bot: Bot):
    await bot.set_webhook(
        url=config.WEBHOOK_URL,
        ip_address=config.WEBHOOK_IP,
        secret_token=config.WEBHOOK_SECRET,
        allowed_updates=["message", "channel_post", "callback_query"],
        drop_pending_updates=True,
    )
    log.info(f"[WEBHOOK] تم تفعيل الويب هوك: {config.WEBHOOK_URL}")

async def on_shutdown(bot: Bot):
    shutdown_handler()
    await bot.delete_webhook()
    log.info("[WEBHOOK] تم حذف الويب هوك عند الإيقاف")


async def main():
    os.makedirs("sessions", exist_ok=True)
    os.makedirs(PROFILE_DIR, exist_ok=True)
    os.makedirs(PROFILE_GIRL_DIR, exist_ok=True)
    await shared.init_db()                           # ← await بدل init_db()
    shared._admins_cache = await shared.get_admins() # ← await بدل get_admins_sync()
    engine.set_bot(bot, notify, notify_edit)
    loop = asyncio.get_event_loop()
    loop.set_exception_handler(handle_exception)

    asyncio.create_task(engine.periodic_cleanup())
    asyncio.create_task(engine.health_check())
    asyncio.create_task(engine._flush_join_notify())

    dp.startup.register(on_startup)
    dp.shutdown.register(on_shutdown)

    app = web.Application()
    SimpleRequestHandler(
        dispatcher=dp,
        bot=bot,
        secret_token=config.WEBHOOK_SECRET,
    ).register(app, path=config.WEBHOOK_PATH)
    setup_application(app, dp, bot=bot)

    runner = web.AppRunner(app)
    await runner.setup()
    site = web.TCPSite(runner, config.WEBHOOK_LISTEN_HOST, config.WEBHOOK_LISTEN_PORT)
    await site.start()
    log.info(f"✅ Bot يعمل الآن (webhook) على {config.WEBHOOK_LISTEN_HOST}:{config.WEBHOOK_LISTEN_PORT}")

    stop_forever = asyncio.Event()
    for sig in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(sig, stop_forever.set)
    await stop_forever.wait()

    await runner.cleanup()

if __name__ == "__main__":
    asyncio.run(main())
