from telethon import TelegramClient, events, Button
from telethon.errors import (
    SessionPasswordNeededError, FloodWaitError, PeerFloodError,
    UserAlreadyParticipantError, InviteHashExpiredError, InviteHashInvalidError,
    ChatAdminRequiredError, ChannelPrivateError,
    UserPrivacyRestrictedError, UserNotMutualContactError,
    UserChannelsTooMuchError, ChatWriteForbiddenError,
    UserBannedInChannelError, InputUserDeactivatedError,
    UserIdInvalidError, PeerIdInvalidError, UsernameNotOccupiedError,
    UsernameInvalidError, AuthKeyUnregisteredError, UserDeactivatedBanError
)
from telethon.tl.types import (
    InputUser,
    UserStatusOnline, UserStatusRecently, UserStatusLastWeek, UserStatusOffline,
    InputUserSelf
)
from telethon.tl.functions.channels import InviteToChannelRequest, JoinChannelRequest
from telethon.tl.functions.messages import AddChatUserRequest, ImportChatInviteRequest
from telethon.tl.functions.contacts import AddContactRequest, GetContactsRequest, DeleteContactsRequest
from telethon.tl.functions.users import GetUsersRequest
import asyncio
import os
import json
import re
from datetime import datetime, timedelta, timezone
import logging

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s | %(levelname)s | %(message)s",
    datefmt="%H:%M:%S"
)
logger = logging.getLogger("Bot")

SESSIONS_FOLDER = "sessions"
LEECH_FOLDER = "leeches"
CONFIG_FILE = "config.json"

if not os.path.exists(SESSIONS_FOLDER):
    os.makedirs(SESSIONS_FOLDER)
if not os.path.exists(LEECH_FOLDER):
    os.makedirs(LEECH_FOLDER)

bot = None
current_task = None
user_states = {}


def load_config():
    if os.path.exists(CONFIG_FILE):
        with open(CONFIG_FILE, "r", encoding="utf-8") as f:
            return json.load(f)
    return {"api_id": None, "api_hash": None, "bot_token": None}


def save_config(config):
    with open(CONFIG_FILE, "w", encoding="utf-8") as f:
        json.dump(config, f, ensure_ascii=False, indent=2)


def get_api():
    config = load_config()
    if config.get("api_id") and config.get("api_hash"):
        return int(config["api_id"]), config["api_hash"]
    return None, None


def get_leech_files():
    return sorted([f for f in os.listdir(LEECH_FOLDER) if f.endswith(".json")])


def get_session_files():
    return sorted([f for f in os.listdir(SESSIONS_FOLDER) if f.endswith(".session")])


def load_leech(selected_files=None):
    all_data = []
    files = selected_files if selected_files is not None else get_leech_files()
    for fname in files:
        path = os.path.join(LEECH_FOLDER, fname)
        try:
            with open(path, "r", encoding="utf-8") as f:
                data = json.load(f)
                if isinstance(data, list):
                    all_data.extend(data)
        except Exception as e:
            logger.error(f"خطا در خواندن {fname}: {e}")
    return all_data


def save_new_leech(new_users, group_title):
    if not new_users:
        return
    safe_title = re.sub(r"[^\w\u0600-\u06FF\-]", "_", group_title)[:40]
    timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
    fname = f"{safe_title}_{timestamp}.json"
    path = os.path.join(LEECH_FOLDER, fname)
    with open(path, "w", encoding="utf-8") as f:
        json.dump(new_users, f, ensure_ascii=False, indent=2)
    logger.info(f"فایل لیچ ذخیره شد → {fname} | {len(new_users)} نفر")


def remove_from_leech(user_id):
    for fname in get_leech_files():
        path = os.path.join(LEECH_FOLDER, fname)
        try:
            with open(path, "r", encoding="utf-8") as f:
                data = json.load(f)
            new_data = [m for m in data if m.get("id") != user_id]
            if len(new_data) != len(data):
                if new_data:
                    with open(path, "w", encoding="utf-8") as f:
                        json.dump(new_data, f, ensure_ascii=False, indent=2)
                else:
                    os.remove(path)
                return
        except Exception as e:
            logger.error(f"خطا در حذف از {fname}: {e}")


def is_active_within_week(user):
    status = getattr(user, "status", None)
    if status is None:
        return False
    if isinstance(status, (UserStatusOnline, UserStatusRecently, UserStatusLastWeek)):
        return True
    if isinstance(status, UserStatusOffline) and status.was_online:
        return status.was_online >= datetime.now(timezone.utc) - timedelta(days=7)
    return False


def format_user(user, group_title=""):
    return {
        "id": user.id,
        "access_hash": getattr(user, "access_hash", None),
        "first_name": user.first_name or "",
        "last_name": user.last_name or "",
        "username": user.username or "",
        "phone": getattr(user, "phone", "") or "",
        "group": group_title,
        "leeched_at": datetime.now(timezone.utc).isoformat()
    }


def extract_invite_hash(link: str):
    link = link.strip()
    patterns = [
        r"(?:https?://)?(?:t\.me|telegram\.me)/(?:\+|joinchat/)([a-zA-Z0-9_-]+)",
        r"(?:https?://)?(?:t\.me|telegram\.me)/joinchat/([a-zA-Z0-9_-]+)",
    ]
    for pattern in patterns:
        match = re.search(pattern, link, re.IGNORECASE)
        if match:
            return match.group(1)
    return None


async def join_group_if_needed(client, group_input: str):
    group_input = group_input.strip()
    invite_hash = extract_invite_hash(group_input)
    if invite_hash:
        try:
            result = await client(ImportChatInviteRequest(invite_hash))
            entity = result.chats[0]
            await asyncio.sleep(1)
            return entity
        except UserAlreadyParticipantError:
            pass
        except Exception as e:
            raise Exception(f"خطا در جوین لینک خصوصی: {e}")
    try:
        entity = await client.get_entity(group_input)
        try:
            await client(JoinChannelRequest(entity))
            await asyncio.sleep(1)
        except UserAlreadyParticipantError:
            pass
        return entity
    except Exception as e:
        raise Exception(f"نتوانست وارد گروه شود: {e}")


def main_menu():
    return [
        [Button.inline("⚙️ تنظیمات API", b"set_api")],
        [Button.inline("➕ اد اکانت", b"add_account"), Button.inline("📂 لیست سشن‌ها", b"list_sessions")],
        [Button.inline("📥 لیچ گروه", b"leech_group"), Button.inline("👥 اد گروه (لیچ)", b"add_group")],
        [Button.inline("👤 اد به گروه از طریق مخاطب", b"add_from_contacts")],
        [Button.inline("📞 اضافه کردن مخاطبین", b"add_contacts"), Button.inline("🗑️ حذف مخاطبین", b"delete_contacts")],
        [Button.inline("📇 تعداد مخاطبین سشن‌ها", b"count_contacts")],
        [Button.inline("🚦 وضعیت فلود", b"check_flood"), Button.inline("📊 آمار لیچ", b"leech_stats")],
        [Button.inline("🗑️ حذف لیچ", b"delete_leech")]
    ]


def stop_button():
    return [[Button.inline("⏹ پایان عملیات", b"stop_operation")]]


def build_file_select_buttons(files, selected):
    buttons = []
    for i, f in enumerate(files):
        mark = "✅ " if i in selected else "⬜ "
        buttons.append([Button.inline(f"{mark}{f}", f"toggle_file_{i}".encode())])
    buttons.append([Button.inline("✅ تأیید و ادامه", b"confirm_files")])
    buttons.append([Button.inline("🔙 بازگشت", b"back_main")])
    return buttons


def build_contact_file_buttons(files, selected):
    buttons = []
    for i, f in enumerate(files):
        mark = "✅ " if i in selected else "⬜ "
        buttons.append([Button.inline(f"{mark}{f}", f"ctoggle_{i}".encode())])
    buttons.append([Button.inline("✅ تأیید و شروع", b"confirm_contacts")])
    buttons.append([Button.inline("🔙 بازگشت", b"back_main")])
    return buttons


async def show_main_menu(event, text="🏠 منوی اصلی:"):
    await event.respond(text, buttons=main_menu())


@events.register(events.CallbackQuery)
async def callback_handler(event):
    global current_task
    data = event.data.decode("utf-8")
    user_id = event.sender_id

    if data == "stop_operation":
        if current_task and not current_task.done():
            current_task.cancel()
            current_task = None
            await event.answer("⏹ متوقف شد", alert=True)
            try:
                await event.edit("⏹ عملیات متوقف شد.", buttons=main_menu())
            except Exception:
                await event.respond("⏹ عملیات متوقف شد.", buttons=main_menu())
        else:
            await event.answer("هیچ عملیاتی در حال اجرا نیست", alert=True)
        return

    if current_task and not current_task.done():
        await event.answer("⚠️ اول عملیات فعلی را متوقف کنید", alert=True)
        return

    if data == "set_api":
        user_states[user_id] = "waiting_api_id"
        await event.edit("⚙️ API ID را ارسال کنید:")

    elif data == "add_account":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        user_states[user_id] = "waiting_phone"
        await event.edit(
            "📱 شماره را با کد کشور بفرستید:\n`+989123456789`",
            buttons=[[Button.inline("🔙 بازگشت", b"back_main")]]
        )

    elif data == "list_sessions":
        files = get_session_files()
        text = "هیچ سشنی نیست" if not files else "سشن‌ها:\n\n" + "\n".join(f"🔹 `{f}`" for f in files)
        await event.edit(text, buttons=main_menu())

    elif data == "leech_stats":
        files = get_leech_files()
        total = len(load_leech())
        if not files:
            await event.edit("📊 هیچ فایل لیچی وجود ندارد.", buttons=main_menu())
            return
        lines = []
        lines.append("📊 آمار لیچ")
        lines.append(f"تعداد کل: {total}")
        lines.append("")
        lines.append("📁 فایل‌ها:")
        for idx, f in enumerate(files, 1):
            path = os.path.join(LEECH_FOLDER, f)
            try:
                with open(path, "r", encoding="utf-8") as ff:
                    count = len(json.load(ff))
            except Exception:
                count = -1
            if count >= 0:
                lines.append(f"{idx}) {f}")
                lines.append(f"   ➜ تعداد: \u200e{count}")
            else:
                lines.append(f"{idx}) {f}")
                lines.append("   ➜ تعداد: خطا در خواندن")
            lines.append("")
        text = "\n".join(lines)
        if len(text) > 3500:
            text = text[:3400] + "\n...\n(لیست طولانی بود)"
        await event.edit(text, buttons=main_menu())

    elif data == "count_contacts":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        sessions = get_session_files()
        if not sessions:
            await event.answer("هیچ سشنی وجود ندارد", alert=True)
            return
        msg = await event.edit(
            f"📇 در حال شمارش مخاطبین {len(sessions)} سشن...",
            buttons=stop_button()
        )

        async def do_count_contacts():
            global current_task
            api_id, api_hash = get_api()
            results = []
            total_all = 0
            try:
                for idx, session_file in enumerate(sessions):
                    if current_task and current_task.cancelled():
                        raise asyncio.CancelledError()
                    session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
                    client = TelegramClient(session_path, api_id, api_hash)
                    status = "❓"
                    try:
                        await client.connect()
                        if not await client.is_user_authorized():
                            status = "❌ سشن نامعتبر"
                        else:
                            result = await client(GetContactsRequest(hash=0))
                            users = list(result.users) if result.users else []
                            me = await client.get_me()
                            users = [u for u in users if not getattr(u, "bot", False) and u.id != me.id]
                            count = len(users)
                            total_all += count
                            status = f"✅ \u200e{count} مخاطب"
                    except FloodWaitError as e:
                        status = f"🔴 فلود {e.seconds}s"
                    except Exception as e:
                        status = f"❌ {type(e).__name__}"
                    finally:
                        try:
                            await client.disconnect()
                        except Exception:
                            pass
                    results.append((session_file, status))
                    try:
                        lines = [f"🔹 `{s}`\n   {st}" for s, st in results]
                        await msg.edit(
                            f"📇 شمارش مخاطبین ({idx + 1}/{len(sessions)})...\n\n"
                            + "\n\n".join(lines)
                            + f"\n\n📊 مجموع تا الان: \u200e{total_all}",
                            buttons=stop_button()
                        )
                    except Exception:
                        pass
                    await asyncio.sleep(1)
                lines = [f"🔹 `{s}`\n   {st}" for s, st in results]
                await msg.edit(
                    "📇 **تعداد مخاطبین سشن‌ها**\n\n"
                    + "\n\n".join(lines)
                    + f"\n\n📊 مجموع کل: \u200e{total_all}",
                    buttons=main_menu()
                )
            except asyncio.CancelledError:
                await msg.edit("⏹ شمارش متوقف شد.", buttons=main_menu())
            finally:
                current_task = None

        current_task = asyncio.create_task(do_count_contacts())

    elif data == "delete_leech":
        files = get_leech_files()
        if not files:
            await event.edit("🗑️ هیچ فایل لیچی نیست.", buttons=main_menu())
            return
        buttons = [[Button.inline(f"🗑️ {f}", f"del_file_{i}".encode())] for i, f in enumerate(files)]
        buttons.append([Button.inline("🗑️ حذف همه", b"del_all_leech")])
        buttons.append([Button.inline("🔙 بازگشت", b"back_main")])
        await event.edit("کدام فایل حذف شود؟", buttons=buttons)

    elif data.startswith("del_file_"):
        idx = int(data.split("_")[-1])
        files = get_leech_files()
        if 0 <= idx < len(files):
            try:
                os.remove(os.path.join(LEECH_FOLDER, files[idx]))
                await event.edit(f"✅ `{files[idx]}` حذف شد.", buttons=main_menu())
            except Exception as e:
                await event.edit(f"❌ {e}", buttons=main_menu())
        else:
            await event.edit("فایل پیدا نشد.", buttons=main_menu())

    elif data == "del_all_leech":
        files = get_leech_files()
        for f in files:
            try:
                os.remove(os.path.join(LEECH_FOLDER, f))
            except Exception:
                pass
        await event.edit(f"🗑️ {len(files)} فایل حذف شد.", buttons=main_menu())

    elif data == "add_from_contacts":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        if not get_session_files():
            await event.answer("هیچ سشنی وجود ندارد", alert=True)
            return
        user_states[user_id] = {"state": "afc_waiting_target"}
        await event.edit(
            "👤 **اد به گروه از طریق مخاطب**\n\n"
            "گروه مقصد را بفرستید (لینک / یوزرنیم / آیدی):\n\n"
            "⚠️ هر سشن باید ادمین باشد و دسترسی Add Members داشته باشد.",
            buttons=[[Button.inline("🔙 بازگشت", b"back_main")]]
        )

    elif data == "check_flood":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        sessions = get_session_files()
        if not sessions:
            await event.answer("هیچ سشنی وجود ندارد", alert=True)
            return
        msg = await event.edit(f"🚦 بررسی فلود {len(sessions)} سشن...", buttons=stop_button())

        async def do_check_flood():
            global current_task
            api_id, api_hash = get_api()
            results = []
            try:
                for idx, session_file in enumerate(sessions):
                    if current_task and current_task.cancelled():
                        raise asyncio.CancelledError()
                    session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
                    client = TelegramClient(session_path, api_id, api_hash)
                    status_text = "❓"
                    try:
                        await client.connect()
                        if not await client.is_user_authorized():
                            status_text = "❌ سشن نامعتبر"
                        else:
                            try:
                                await client(GetUsersRequest(id=[InputUserSelf()]))
                                await client(GetContactsRequest(hash=0))
                                me = await client.get_me()
                                status_text = f"✅ سالم | {me.first_name or ''}"
                            except FloodWaitError as e:
                                until = datetime.now() + timedelta(seconds=e.seconds)
                                h = e.seconds // 3600
                                m = (e.seconds % 3600) // 60
                                s = e.seconds % 60
                                t = f"{h}س {m}د {s}ث" if h else f"{m}د {s}ث"
                                status_text = f"🔴 فلود\n⏱ {t}\n📅 تا {until.strftime('%H:%M:%S')}"
                            except PeerFloodError:
                                status_text = "🟠 PeerFlood"
                            except (AuthKeyUnregisteredError, UserDeactivatedBanError):
                                status_text = "🚫 بن"
                            except Exception as e:
                                status_text = f"⚠️ {type(e).__name__}"
                    except Exception as e:
                        status_text = f"❌ {type(e).__name__}"
                    finally:
                        try:
                            await client.disconnect()
                        except Exception:
                            pass
                    results.append((session_file, status_text))
                    try:
                        lines = [f"🔹 `{s}`\n{st}\n" for s, st in results]
                        await msg.edit(
                            f"🚦 ({idx + 1}/{len(sessions)})...\n\n" + "\n".join(lines),
                            buttons=stop_button()
                        )
                    except Exception:
                        pass
                    await asyncio.sleep(1)
                lines = [f"🔹 `{s}`\n{st}\n" for s, st in results]
                await msg.edit("🚦 **نتیجه فلود**\n\n" + "\n".join(lines), buttons=main_menu())
            except asyncio.CancelledError:
                await msg.edit("⏹ متوقف شد.", buttons=main_menu())
            finally:
                current_task = None

        current_task = asyncio.create_task(do_check_flood())

    elif data == "add_contacts":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        if not get_session_files():
            await event.answer("هیچ سشنی وجود ندارد", alert=True)
            return
        files = get_leech_files()
        if not files:
            await event.answer("هیچ فایل لیچی وجود ندارد", alert=True)
            return
        user_states[user_id] = {"state": "contact_selecting_files", "files": files, "selected": set()}
        await event.edit("📞 فایل‌های لیچ را انتخاب کنید:", buttons=build_contact_file_buttons(files, set()))

    elif data.startswith("ctoggle_"):
        state = user_states.get(user_id)
        if not state or state.get("state") != "contact_selecting_files":
            return
        idx = int(data.split("_")[-1])
        selected = state["selected"]
        if idx in selected:
            selected.remove(idx)
        else:
            selected.add(idx)
        await event.edit(
            "📞 فایل‌های لیچ را انتخاب کنید:",
            buttons=build_contact_file_buttons(state["files"], selected)
        )

    elif data == "confirm_contacts":
        state = user_states.get(user_id)
        if not state or state.get("state") != "contact_selecting_files":
            return
        selected = state["selected"]
        if not selected:
            await event.answer("حداقل یک فایل انتخاب کنید", alert=True)
            return
        selected_files = [state["files"][i] for i in sorted(selected)]
        members = list({m["id"]: m for m in load_leech(selected_files)}.values())
        sessions = get_session_files()
        total = len(members)
        per_session = total // len(sessions) if sessions else 0
        remainder = total % len(sessions) if sessions else 0
        user_states.pop(user_id, None)
        msg = await event.edit(
            f"📞 شروع اضافه کردن مخاطبین\n\n"
            f"📊 کل: **{total}** | 📂 سشن: **{len(sessions)}** | هر سشن ~**{per_session}**",
            buttons=stop_button()
        )

        async def do_add_contacts():
            global current_task
            api_id, api_hash = get_api()
            total_success = 0
            total_failed = 0
            session_stats = {s: [0, 0] for s in sessions}
            lock = asyncio.Lock()
            chunks = []
            start = 0
            for i in range(len(sessions)):
                extra = 1 if i < remainder else 0
                size = per_session + extra
                chunks.append(members[start:start + size])
                start += size

            async def contact_worker(session_file, batch, start_delay):
                nonlocal total_success, total_failed
                if not batch:
                    return
                await asyncio.sleep(start_delay)
                if current_task and current_task.cancelled():
                    return
                session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
                client = TelegramClient(session_path, api_id, api_hash)
                success = 0
                failed = 0
                try:
                    await client.start()
                    for i, member in enumerate(batch, 1):
                        if current_task and current_task.cancelled():
                            break
                        try:
                            input_user = None
                            if member.get("username"):
                                try:
                                    input_user = await client.get_input_entity(member["username"])
                                except Exception:
                                    pass
                            if input_user is None:
                                try:
                                    input_user = await client.get_input_entity(member["id"])
                                except Exception:
                                    if member.get("access_hash") is not None:
                                        input_user = InputUser(member["id"], int(member["access_hash"]))
                            if input_user is None:
                                failed += 1
                                continue
                            await client(AddContactRequest(
                                id=input_user,
                                first_name=(member.get("first_name") or str(member["id"]))[:64],
                                last_name=(member.get("last_name") or "")[:64],
                                phone="",
                                add_phone_privacy_exception=False
                            ))
                            success += 1
                            async with lock:
                                total_success += 1
                                session_stats[session_file] = [success, failed]
                        except FloodWaitError as e:
                            await asyncio.sleep(e.seconds + 2)
                            continue
                        except Exception:
                            failed += 1
                            async with lock:
                                total_failed += 1
                                session_stats[session_file] = [success, failed]
                        await asyncio.sleep(1.2)
                    session_stats[session_file] = [success, failed]
                except Exception as e:
                    logger.error(f"[{session_file}] {e}")
                    session_stats[session_file] = [success, failed]
                finally:
                    try:
                        await client.disconnect()
                    except Exception:
                        pass

            workers = [
                asyncio.create_task(
                    contact_worker(sessions[i], chunks[i] if i < len(chunks) else [], i * 2)
                )
                for i in range(len(sessions))
            ]

            async def progress_updater():
                while any(not w.done() for w in workers):
                    if current_task and current_task.cancelled():
                        break
                    try:
                        lines = [
                            f"• `{s}` → ✅{session_stats[s][0]} | ❌{session_stats[s][1]}"
                            for s in sessions
                        ]
                        await msg.edit(
                            f"📞 اضافه کردن مخاطبین...\n\n{chr(10).join(lines)}\n\n"
                            f"✅ کل: **{total_success}** | ❌ **{total_failed}**",
                            buttons=stop_button()
                        )
                    except Exception:
                        pass
                    await asyncio.sleep(3)

            updater = asyncio.create_task(progress_updater())
            try:
                await asyncio.gather(*workers)
            except asyncio.CancelledError:
                for w in workers:
                    w.cancel()
                raise
            finally:
                updater.cancel()
            lines = [f"• `{s}` → ✅{session_stats[s][0]} | ❌{session_stats[s][1]}" for s in sessions]
            await msg.edit(
                f"✅ تمام شد\n\n{chr(10).join(lines)}\n\n✅ **{total_success}** | ❌ **{total_failed}**",
                buttons=main_menu()
            )
            current_task = None

        current_task = asyncio.create_task(do_add_contacts())

    elif data == "delete_contacts":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        sessions = get_session_files()
        if not sessions:
            await event.answer("هیچ سشنی وجود ندارد", alert=True)
            return
        msg = await event.edit(f"🗑️ حذف مخاطبین از {len(sessions)} سشن...", buttons=stop_button())

        async def do_delete_contacts():
            global current_task
            api_id, api_hash = get_api()
            results = []
            total_deleted = 0
            try:
                for idx, session_file in enumerate(sessions):
                    if current_task and current_task.cancelled():
                        raise asyncio.CancelledError()
                    session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
                    client = TelegramClient(session_path, api_id, api_hash)
                    try:
                        await client.connect()
                        if not await client.is_user_authorized():
                            results.append((session_file, "❌ نامعتبر"))
                        else:
                            result = await client(GetContactsRequest(hash=0))
                            users = list(result.users) if result.users else []
                            if not users:
                                results.append((session_file, "✅ ۰ (خالی)"))
                            else:
                                await client(DeleteContactsRequest(id=users))
                                total_deleted += len(users)
                                results.append((session_file, f"✅ {len(users)} نفر"))
                    except FloodWaitError as e:
                        results.append((session_file, f"🔴 فلود {e.seconds}s"))
                    except Exception as e:
                        results.append((session_file, f"❌ {type(e).__name__}"))
                    finally:
                        try:
                            await client.disconnect()
                        except Exception:
                            pass
                    try:
                        lines = [f"• `{s}` → {st}" for s, st in results]
                        await msg.edit(
                            f"🗑️ ({idx + 1}/{len(sessions)})...\n\n" + "\n".join(lines) +
                            f"\n\n🗑️ مجموع: **{total_deleted}**",
                            buttons=stop_button()
                        )
                    except Exception:
                        pass
                    await asyncio.sleep(1.5)
                lines = [f"• `{s}` → {st}" for s, st in results]
                await msg.edit(
                    f"✅ حذف تمام شد\n\n" + "\n".join(lines) + f"\n\n🗑️ مجموع: **{total_deleted}**",
                    buttons=main_menu()
                )
            except asyncio.CancelledError:
                await msg.edit(f"⏹ متوقف شد\n🗑️ تا الان: {total_deleted}", buttons=main_menu())
            finally:
                current_task = None

        current_task = asyncio.create_task(do_delete_contacts())

    elif data == "leech_group":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        sessions = get_session_files()
        if not sessions:
            await event.answer("هیچ سشنی وجود ندارد", alert=True)
            return
        buttons = [[Button.inline(f"🔹 {s}", f"leech_sess_{i}".encode())] for i, s in enumerate(sessions)]
        buttons.append([Button.inline("🔙 بازگشت", b"back_main")])
        await event.edit("سشن را انتخاب کنید:", buttons=buttons)

    elif data.startswith("leech_sess_"):
        idx = int(data.split("_")[-1])
        sessions = get_session_files()
        session_file = sessions[idx]
        user_states[user_id] = {"state": "leech_choose_method", "session": session_file}
        buttons = [
            [Button.inline("👥 از لیست اعضا", b"leech_method_members")],
            [Button.inline("💬 از پیام‌ها", b"leech_method_messages")],
            [Button.inline("🔙 بازگشت", b"back_main")]
        ]
        await event.edit(f"سشن: `{session_file}`\nروش لیچ:", buttons=buttons)

    elif data == "leech_method_members":
        user_states[user_id]["state"] = "leech_waiting_group_members"
        await event.edit("لینک / یوزرنیم / آیدی گروه را بفرستید:")

    elif data == "leech_method_messages":
        user_states[user_id]["state"] = "leech_waiting_days"
        await event.edit("تا چند روز پیش؟ (فقط عدد)")

    elif data == "add_group":
        if not get_api()[0]:
            await event.answer("اول تنظیمات API را انجام دهید", alert=True)
            return
        files = get_leech_files()
        if not files:
            await event.answer("هیچ فایل لیچی نیست", alert=True)
            return
        user_states[user_id] = {"state": "add_selecting_files", "files": files, "selected": set()}
        await event.edit("📁 فایل‌های لیچ را انتخاب کنید:", buttons=build_file_select_buttons(files, set()))

    elif data.startswith("toggle_file_"):
        state = user_states.get(user_id)
        if not state or state.get("state") != "add_selecting_files":
            return
        idx = int(data.split("_")[-1])
        selected = state["selected"]
        if idx in selected:
            selected.remove(idx)
        else:
            selected.add(idx)
        await event.edit(
            "📁 فایل‌های لیچ را انتخاب کنید:",
            buttons=build_file_select_buttons(state["files"], selected)
        )

    elif data == "confirm_files":
        state = user_states.get(user_id)
        if not state or state.get("state") != "add_selecting_files":
            return
        selected = state["selected"]
        if not selected:
            await event.answer("حداقل یک فایل انتخاب کنید", alert=True)
            return
        selected_files = [state["files"][i] for i in sorted(selected)]
        total_count = len(load_leech(selected_files))
        user_states[user_id] = {"state": "add_waiting_target", "selected_files": selected_files}
        await event.edit(
            f"✅ {len(selected_files)} فایل | 📊 **{total_count}** نفر\n\nگروه مقصد را بفرستید:"
        )

    elif data == "back_main":
        user_states.pop(user_id, None)
        await event.edit("🏠 منوی اصلی:", buttons=main_menu())
@events.register(events.NewMessage)
async def message_handler(event):
    global current_task
    if not event.is_private:
        return

    user_id = event.sender_id
    text = event.raw_text.strip()
    state = user_states.get(user_id)

    if not state:
        if text == "/start":
            await show_main_menu(event, "👋 سلام! ربات آماده است.")
        return

    if state == "waiting_api_id":
        try:
            api_id = int(text)
            user_states[user_id] = {"state": "waiting_api_hash", "api_id": api_id}
            await event.respond("API HASH را بفرستید:")
        except Exception:
            await event.respond("API ID باید عدد باشد.")

    elif isinstance(state, dict) and state.get("state") == "waiting_api_hash":
        config = load_config()
        config["api_id"] = state["api_id"]
        config["api_hash"] = text
        save_config(config)
        user_states.pop(user_id, None)
        await event.respond("✅ تنظیمات API ذخیره شد.", buttons=main_menu())

    elif state == "waiting_phone":
        api_id, api_hash = get_api()
        session_name = text.replace("+", "").replace(" ", "")
        session_path = os.path.join(SESSIONS_FOLDER, session_name)
        client = TelegramClient(session_path, api_id, api_hash)
        await client.connect()
        try:
            await client.send_code_request(text)
            user_states[user_id] = {"state": "waiting_code", "phone": text, "client": client}
            await event.respond("کد تأیید را بفرستید:")
        except Exception as e:
            await client.disconnect()
            user_states.pop(user_id, None)
            await event.respond(f"خطا: `{e}`", buttons=main_menu())

    elif isinstance(state, dict) and state.get("state") == "waiting_code":
        client = state["client"]
        try:
            await client.sign_in(state["phone"], text)
            me = await client.get_me()
            await event.respond(f"✅ سشن ساخته شد\n{me.first_name} | `{me.id}`", buttons=main_menu())
            await client.disconnect()
            user_states.pop(user_id, None)
        except SessionPasswordNeededError:
            user_states[user_id]["state"] = "waiting_password"
            await event.respond("رمز دو مرحله‌ای را بفرستید:")
        except Exception as e:
            await client.disconnect()
            user_states.pop(user_id, None)
            await event.respond(f"خطا: `{e}`", buttons=main_menu())

    elif isinstance(state, dict) and state.get("state") == "waiting_password":
        client = state["client"]
        try:
            await client.sign_in(password=text)
            me = await client.get_me()
            await event.respond(f"✅ سشن ساخته شد\n{me.first_name}", buttons=main_menu())
        except Exception as e:
            await event.respond(f"خطا: `{e}`", buttons=main_menu())
        await client.disconnect()
        user_states.pop(user_id, None)

    elif isinstance(state, dict) and state.get("state") == "leech_waiting_days":
        try:
            days = int(text)
            if days <= 0:
                raise ValueError
            user_states[user_id]["days"] = days
            user_states[user_id]["state"] = "leech_waiting_group_messages"
            await event.respond("لینک یا یوزرنیم گروه را بفرستید:")
        except Exception:
            await event.respond("فقط عدد معتبر بفرستید.")

    elif isinstance(state, dict) and state.get("state") == "afc_waiting_target":
        user_states[user_id]["target"] = text
        user_states[user_id]["state"] = "afc_waiting_per_session"
        await event.respond(
            f"📂 تعداد سشن: **{len(get_session_files())}**\n\n"
            f"هر سشن چند نفر **موفق** از مخاطبین خودش اد کند؟\n(فقط عدد — مثال: ۱۰)"
        )

    elif isinstance(state, dict) and state.get("state") == "afc_waiting_per_session":
        try:
            per = int(text)
            if per <= 0:
                raise ValueError
            user_states[user_id]["per_session"] = per
            user_states[user_id]["state"] = "afc_waiting_delay"
            await event.respond(
                "⏱️ **تاخیر بین هر اد چند ثانیه باشد؟**\n\n"
                "• عدد کمتر (۸) → سریع‌تر ولی ریسک فلود بیشتر\n"
                "• عدد بیشتر (۱۵ تا ۲۵) → امن‌تر\n"
                "• پیشنهاد امن: **۱۲ تا ۲۰**\n\n"
                "فقط عدد بفرستید:"
            )
        except Exception:
            await event.respond("عدد معتبر بفرستید.")

    elif isinstance(state, dict) and state.get("state") == "afc_waiting_delay":
        try:
            delay = int(text)
            if delay < 3:
                delay = 8
        except Exception:
            delay = 12

        target = state["target"]
        per = state["per_session"]
        user_states.pop(user_id, None)
        sessions = get_session_files()

        msg = await event.respond(
            f"🚀 اد از طریق مخاطبین\n\n"
            f"🎯 گروه: `{target}`\n"
            f"📂 سشن‌ها: **{len(sessions)}**\n"
            f"✅ هدف هر سشن: **{per} موفق**\n"
            f"⏱️ تاخیر بین هر اد: **{delay} ثانیه**\n"
            f"⏳ شروع سشن‌ها با فاصله\n"
            f"🗑️ هر مخاطب پردازش‌شده از مخاطبین حذف می‌شود\n"
            f"ℹ️ فقط اعضای یکتا شمرده می‌شوند",
            buttons=stop_button()
        )

        async def do_add_from_contacts():
            global current_task
            api_id, api_hash = get_api()
            total_success = 0
            total_skipped = 0
            total_failed = 0
            session_stats = {s: 0 for s in sessions}
            peerflood_count = {s: 0 for s in sessions}
            added_ids = set()
            lock = asyncio.Lock()

            async def safe_delete_contact(client, contact):
                try:
                    await client(DeleteContactsRequest(id=[contact]))
                except Exception as e:
                    logger.warning(f"حذف مخاطب: {e}")

            async def worker(session_file, start_delay):
                nonlocal total_success, total_skipped, total_failed
                await asyncio.sleep(start_delay)
                if current_task and current_task.cancelled():
                    return

                session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
                client = TelegramClient(session_path, api_id, api_hash)
                success_this = 0

                try:
                    await client.start()
                    entity = await join_group_if_needed(client, target)

                    result = await client(GetContactsRequest(hash=0))
                    contacts = list(result.users) if result.users else []
                    me = await client.get_me()
                    contacts = [u for u in contacts if not getattr(u, "bot", False) and u.id != me.id]
                    logger.info(f"[{session_file}] {len(contacts)} مخاطب | هدف: {per} | تاخیر: {delay}s")

                    if not contacts:
                        return

                    for contact in contacts:
                        if success_this >= per:
                            break
                        if current_task and current_task.cancelled():
                            break

                        async with lock:
                            already_done = contact.id in added_ids
                        if already_done:
                            await safe_delete_contact(client, contact)
                            async with lock:
                                total_skipped += 1
                            continue

                        try:
                            input_user = InputUser(contact.id, contact.access_hash)
                            is_channel = (
                                getattr(entity, "megagroup", False)
                                or getattr(entity, "broadcast", False)
                                or entity.__class__.__name__ in ("Channel", "ChannelForbidden")
                            )

                            really_added = False
                            if is_channel:
                                res = await client(InviteToChannelRequest(entity, [input_user]))
                                missing = getattr(res, "missing_invitees", None) or []
                                missing_ids = set()
                                for m in missing:
                                    uid = getattr(getattr(m, "user_id", None), "user_id", None)
                                    if uid is None:
                                        uid = getattr(m, "user_id", None)
                                    if uid:
                                        missing_ids.add(uid)
                                really_added = contact.id not in missing_ids
                            else:
                                await client(AddChatUserRequest(entity.id, input_user, fwd_limit=10))
                                really_added = True

                            await safe_delete_contact(client, contact)

                            if really_added:
                                async with lock:
                                    if contact.id not in added_ids:
                                        added_ids.add(contact.id)
                                        total_success += 1
                                        success_this += 1
                                        session_stats[session_file] = success_this
                                    else:
                                        total_skipped += 1
                                logger.info(f"[{session_file}] ✅ یکتا ({success_this}/{per}) {contact.first_name}")
                            else:
                                async with lock:
                                    total_skipped += 1
                                logger.info(f"[{session_file}] ⏭ missing_invitees → حذف")

                        except UserAlreadyParticipantError:
                            await safe_delete_contact(client, contact)
                            async with lock:
                                added_ids.add(contact.id)
                                total_skipped += 1

                        except (
                            UserPrivacyRestrictedError, UserNotMutualContactError,
                            UserChannelsTooMuchError, InputUserDeactivatedError,
                            UserIdInvalidError, PeerIdInvalidError,
                            UserBannedInChannelError, ChatWriteForbiddenError
                        ):
                            await safe_delete_contact(client, contact)
                            async with lock:
                                total_skipped += 1

                        except FloodWaitError as e:
                            logger.warning(f"[{session_file}] FloodWait {e.seconds}s")
                            await asyncio.sleep(e.seconds + 5)
                            continue

                        except PeerFloodError:
                            peerflood_count[session_file] += 1
                            if peerflood_count[session_file] >= 3:
                                break
                            await asyncio.sleep(300)
                            continue

                        except Exception as e:
                            await safe_delete_contact(client, contact)
                            async with lock:
                                total_failed += 1
                            logger.error(f"[{session_file}] {type(e).__name__}: {e}")

                        try:
                            lines = [
                                f"• `{s}` → ✅{session_stats[s]}/{per}" + (" ⚠️" if peerflood_count[s] else "")
                                for s in sessions
                            ]
                            async with lock:
                                ts, tsk, tf = total_success, total_skipped, total_failed
                            await msg.edit(
                                f"🔄 اد از طریق مخاطبین...\n"
                                f"⏱️ تاخیر: {delay}s\n\n"
                                f"{chr(10).join(lines)}\n\n"
                                f"✅ موفق یکتا در گروه: **{ts}**\n"
                                f"⏭ رد/تکراری/قبلاً عضو: {tsk}\n"
                                f"❌ خطا: {tf}",
                                buttons=stop_button()
                            )
                        except Exception:
                            pass

                        await asyncio.sleep(delay)

                except Exception as e:
                    logger.error(f"[{session_file}] خطای کلی: {e}")
                finally:
                    try:
                        await client.disconnect()
                    except Exception:
                        pass

            workers = [asyncio.create_task(worker(s, i * 3)) for i, s in enumerate(sessions)]
            try:
                await asyncio.gather(*workers)
            except asyncio.CancelledError:
                for w in workers:
                    w.cancel()
                raise

            await msg.edit(
                f"✅ اد از طریق مخاطبین تمام شد\n\n"
                f"✅ تعداد یکتای اضافه‌شده به گروه: **{total_success}**\n"
                f"⏭ رد / تکراری / قبلاً عضو: {total_skipped}\n"
                f"❌ خطا: {total_failed}\n\n"
                f"ℹ️ این عدد فقط اعضای جدید و غیرتکراری را می‌شمارد.",
                buttons=main_menu()
            )
            current_task = None

        current_task = asyncio.create_task(do_add_from_contacts())

    elif isinstance(state, dict) and state.get("state") in [
        "leech_waiting_group_members", "leech_waiting_group_messages"
    ]:
        group = text
        session_file = state["session"]
        api_id, api_hash = get_api()
        session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
        method = state["state"]
        msg = await event.respond("⏳ ورود به گروه و لیچ...", buttons=stop_button())

        async def do_leech():
            global current_task
            client = TelegramClient(session_path, api_id, api_hash)
            try:
                await client.start()
                entity = await join_group_if_needed(client, group)
                title = getattr(entity, "title", str(entity.id))
                await msg.edit(f"✅ وارد شدم: **{title}**\n⏳ لیچ فعال‌ها...", buttons=stop_button())
                existing_ids = {m["id"] for m in load_leech()}
                new_users = []
                checked = 0
                filtered = 0

                if method == "leech_waiting_group_members":
                    try:
                        async for user in client.iter_participants(entity, aggressive=True):
                            if current_task.cancelled():
                                raise asyncio.CancelledError()
                            checked += 1
                            if user.bot or user.id in existing_ids:
                                continue
                            if not is_active_within_week(user):
                                filtered += 1
                                continue
                            new_users.append(format_user(user, title))
                            existing_ids.add(user.id)
                            if len(new_users) % 20 == 0:
                                try:
                                    await msg.edit(
                                        f"⏳ لیچ...\n🆕 **{len(new_users)}** | 🔍 {checked} | 🚫 {filtered}",
                                        buttons=stop_button()
                                    )
                                except Exception:
                                    pass
                        save_new_leech(new_users, title)
                        await msg.edit(
                            f"✅ تمام شد\n🆕 **{len(new_users)}** | 📊 کل **{len(load_leech())}** | 🚫 {filtered}",
                            buttons=main_menu()
                        )
                    except ChatAdminRequiredError:
                        await msg.edit(
                            "❌ لیست اعضا مخفی است. از روش پیام‌ها استفاده کنید.",
                            buttons=main_menu()
                        )
                    except Exception as e:
                        await msg.edit(f"❌ `{e}`", buttons=main_menu())
                else:
                    days = state.get("days", 7)
                    cutoff = datetime.now(timezone.utc) - timedelta(days=days)
                    async for message in client.iter_messages(entity):
                        if current_task.cancelled():
                            raise asyncio.CancelledError()
                        msg_date = message.date
                        if msg_date.tzinfo is None:
                            msg_date = msg_date.replace(tzinfo=timezone.utc)
                        if msg_date < cutoff:
                            break
                        if not message.sender_id:
                            continue
                        try:
                            sender = message.sender or await client.get_entity(message.sender_id)
                        except Exception:
                            continue
                        if getattr(sender, "bot", False) or sender.id in existing_ids:
                            continue
                        if not is_active_within_week(sender):
                            filtered += 1
                            continue
                        new_users.append(format_user(sender, title))
                        existing_ids.add(sender.id)
                        if len(new_users) % 15 == 0:
                            try:
                                await msg.edit(
                                    f"⏳ لیچ پیام‌ها...\n🆕 **{len(new_users)}**",
                                    buttons=stop_button()
                                )
                            except Exception:
                                pass
                    save_new_leech(new_users, title)
                    await msg.edit(
                        f"✅ تمام شد\n🆕 **{len(new_users)}** | 📊 کل **{len(load_leech())}**",
                        buttons=main_menu()
                    )
            except asyncio.CancelledError:
                await msg.edit("⏹ متوقف شد.", buttons=main_menu())
            except Exception as e:
                await msg.edit(f"❌ `{e}`", buttons=main_menu())
            finally:
                try:
                    await client.disconnect()
                except Exception:
                    pass
                user_states.pop(user_id, None)
                current_task = None

        current_task = asyncio.create_task(do_leech())

    elif isinstance(state, dict) and state.get("state") == "add_waiting_target":
        user_states[user_id]["target"] = text
        user_states[user_id]["state"] = "add_waiting_per_session"
        await event.respond(
            f"سشن‌ها: {len(get_session_files())}\n\nهر سشن چند نفر **موفق** اد کند؟"
        )

    elif isinstance(state, dict) and state.get("state") == "add_waiting_per_session":
        try:
            per = int(text)
            if per <= 0:
                raise ValueError
            user_states[user_id]["per_session"] = per
            user_states[user_id]["state"] = "add_waiting_delay"
            await event.respond("تاخیر بین هر اد (ثانیه)؟ پیشنهاد: ۲۰ تا ۴۰")
        except Exception:
            await event.respond("عدد معتبر بفرستید.")

    elif isinstance(state, dict) and state.get("state") == "add_waiting_delay":
        try:
            delay = max(10, int(text))
        except Exception:
            delay = 25
        target = state["target"]
        per_session = state["per_session"]
        selected_files = state.get("selected_files", [])
        user_states.pop(user_id, None)
        total_members = len(load_leech(selected_files))
        msg = await event.respond(
            f"🚀 اد از لیچ | اعضا: **{total_members}** | هر سشن: **{per_session}** | تاخیر: {delay}s",
            buttons=stop_button()
        )

        async def do_add():
            global current_task
            api_id, api_hash = get_api()
            sessions = get_session_files()
            members = load_leech(selected_files)
            lock = asyncio.Lock()
            total_success = 0
            total_failed = 0
            total_skipped = 0
            session_stats = {s: 0 for s in sessions}
            peerflood_count = {s: 0 for s in sessions}

            async def session_worker(session_file, start_delay):
                nonlocal total_success, total_failed, total_skipped
                await asyncio.sleep(start_delay)
                if current_task and current_task.cancelled():
                    return
                session_path = os.path.join(SESSIONS_FOLDER, session_file.replace(".session", ""))
                client = TelegramClient(session_path, api_id, api_hash)
                success_this = 0
                try:
                    await client.start()
                    entity = await join_group_if_needed(client, target)
                except Exception as e:
                    logger.error(f"[{session_file}] {e}")
                    return

                while success_this < per_session:
                    if current_task and current_task.cancelled():
                        break
                    member = None
                    async with lock:
                        if not members:
                            break
                        member = members.pop(0)
                    if not member:
                        break
                    try:
                        input_user = None
                        if member.get("username"):
                            try:
                                input_user = await client.get_input_entity(member["username"])
                            except Exception:
                                pass
                        if input_user is None:
                            try:
                                input_user = await client.get_input_entity(member["id"])
                            except Exception:
                                if member.get("access_hash") is not None:
                                    input_user = InputUser(member["id"], int(member["access_hash"]))
                        if input_user is None:
                            async with lock:
                                total_failed += 1
                            continue
                        is_channel = getattr(entity, "megagroup", False) or getattr(entity, "broadcast", False)
                        if is_channel:
                            await client(InviteToChannelRequest(entity, [input_user]))
                        else:
                            await client(AddChatUserRequest(entity.id, input_user, fwd_limit=10))
                        async with lock:
                            total_success += 1
                            success_this += 1
                            session_stats[session_file] = success_this
                            remove_from_leech(member["id"])
                    except UserAlreadyParticipantError:
                        async with lock:
                            total_skipped += 1
                            remove_from_leech(member["id"])
                    except (
                        UserPrivacyRestrictedError, UserNotMutualContactError,
                        UserChannelsTooMuchError, InputUserDeactivatedError,
                        UserIdInvalidError, PeerIdInvalidError, UserBannedInChannelError
                    ):
                        async with lock:
                            total_skipped += 1
                            remove_from_leech(member["id"])
                    except FloodWaitError as e:
                        async with lock:
                            members.insert(0, member)
                        await asyncio.sleep(e.seconds + 5)
                        continue
                    except PeerFloodError:
                        peerflood_count[session_file] += 1
                        async with lock:
                            members.insert(0, member)
                        if peerflood_count[session_file] >= 3:
                            break
                        await asyncio.sleep(300)
                        continue
                    except ChatWriteForbiddenError:
                        break
                    except Exception:
                        async with lock:
                            total_failed += 1
                            members.insert(0, member)
                    try:
                        lines = [f"• `{s}` → {session_stats[s]}/{per_session}" for s in sessions]
                        await msg.edit(
                            f"🔄 اد از لیچ...\n\n{chr(10).join(lines)}\n\n"
                            f"✅ {total_success} | ⏭ {total_skipped} | ❌ {total_failed} | باقی: {len(members)}",
                            buttons=stop_button()
                        )
                    except Exception:
                        pass
                    await asyncio.sleep(delay)
                try:
                    await client.disconnect()
                except Exception:
                    pass

            workers = [asyncio.create_task(session_worker(s, i * 3)) for i, s in enumerate(sessions)]
            try:
                await asyncio.gather(*workers)
            except asyncio.CancelledError:
                for w in workers:
                    w.cancel()
                raise
            await msg.edit(
                f"✅ تمام شد\n✅ موفق: **{total_success}**\n⏭ رد: {total_skipped}\n❌ ناموفق: {total_failed}",
                buttons=main_menu()
            )
            current_task = None

        current_task = asyncio.create_task(do_add())


async def main():
    global bot
    config = load_config()
    if not config.get("bot_token"):
        token = input("Bot Token: ").strip()
        config["bot_token"] = token
        save_config(config)

    api_id = config.get("api_id") or 6
    api_hash = config.get("api_hash") or "eb06d4abfb49dc3eeb1aeb98ae0f581e"

    bot = TelegramClient("bot_session", int(api_id), api_hash)
    await bot.start(bot_token=config["bot_token"])
    bot.add_event_handler(callback_handler)
    bot.add_event_handler(message_handler)
    logger.info("✅ ربات روشن شد")
    print("ربات آماده است.")
    await bot.run_until_disconnected()


if __name__ == "__main__":
    asyncio.run(main())        