# gather_engine.py version: 2.10

import asyncio
import gc
import logging
import os
import random
import re
import sqlite3
import time
from datetime import datetime

from pyrogram import Client
from pyrogram.errors import (
    FloodWait,
    UserAlreadyParticipant,
    UserNotParticipant,
    PeerIdInvalid,
)

try:
    from pyrogram.errors import UserDeactivatedBan
except Exception:
    try:
        from pyrogram.errors import UserDeactivated
    except Exception:
        class UserDeactivated(Exception):
            pass
        class UserDeactivatedBan(UserDeactivated):
            pass

import config
from shared import (
    BOT_DELAY, BATCH_REST,
    clients, flood_map, stop_event, account_join_count,
    _background_tasks,
    get_setting, set_setting, get_speed, get_concurrent,
    get_work_minutes, get_rest_minutes, get_concurrent_all_bots,
    get_accounts, get_active_count, ban_account, get_bots,
    get_channels, add_channel, del_channel, get_account_lock, cleanup_client,
    add_force_sub_channel, get_all_force_sub_keys,
    del_force_sub_channels,
    rand_delay, safe_sleep, ch_link, norm_target, norm_ch_key,
    is_ban_error, get_ban_reason, create_tracked_task,
    is_channel_banned, ban_channel,
    get_client_open_lock,
    dead_sessions,
    _join_notify_buffer, get_join_notify_lock
)
import shared as _shared

log = logging.getLogger(__name__)


# مهلة زمنية قصوى لفتح أي جلسة (c.start + c.get_me) — تمنع تجمّد الحلقة
# على جلسة معلّقة إلى الأبد (كانت السبب الرئيسي لتوقف "بدء التجميع").
CLIENT_CONNECT_TIMEOUT = 25

def _extract_bot_username(bot_un: str) -> str:
    if not bot_un:
        return bot_un
    un = bot_un.strip().lstrip("@")
    for prefix in ("https://t.me/", "http://t.me/", "t.me/"):
        if un.startswith(prefix):
            un = un[len(prefix):]
            break
    return un.split("?")[0].strip()

# ============================================================
# نظام الإشعارات مع Queue
# ============================================================
_bot_instance = None
_notify_func = None
_notify_edit_func = None
_notify_queue: asyncio.Queue = None
_notify_worker_task = None
_banned_notified: set = set()



def set_bot(bot_obj, notify_func, notify_edit_func=None):
    global _bot_instance, _notify_func, _notify_edit_func
    _bot_instance = bot_obj
    _notify_func = notify_func
    _notify_edit_func = notify_edit_func


async def start_notify_worker():
    global _notify_queue, _notify_worker_task
    if _notify_worker_task and not _notify_worker_task.done():
        return
    _notify_queue = asyncio.Queue()

    async def _worker():
        # ← دعم قائمة قنوات أو قناة واحدة
        ch_ids = getattr(config, "NOTIFY_CHANNEL_IDS", None)
        if not ch_ids:
            single = getattr(config, "NOTIFY_CHANNEL_ID", None)
            ch_ids = [single] if single else []

        # ← وقت انتهاء FloodWait لكل قناة (0 = جاهزة)
        flood_until: dict = {ch_id: 0.0 for ch_id in ch_ids}

        def get_available_channel():
            """أعد أول قناة جاهزة، أو الأقرب لانتهاء FloodWait"""
            now = asyncio.get_event_loop().time()
            for ch_id in ch_ids:
                if flood_until[ch_id] <= now:
                    return ch_id
            return min(ch_ids, key=lambda x: flood_until[x])

        while True:
            try:
                dest, text = await _notify_queue.get()

                if dest == "channel" and ch_ids and _bot_instance:
                    sent = False
                    attempts = 0
                    max_attempts = len(ch_ids) * 3

                    while not sent and attempts < max_attempts:
                        attempts += 1
                        ch_id = get_available_channel()
                        now = asyncio.get_event_loop().time()

                        # ← انتظر إذا كل القنوات في FloodWait
                        wait = flood_until[ch_id] - now
                        if wait > 0:
                            log.info(
                                f"[NOTIFY] كل القنوات في FloodWait، "
                                f"انتظار {int(wait)+1}s"
                            )
                            await asyncio.sleep(wait + 1)

                        try:
                            await _bot_instance.send_message(
                                ch_id, text,
                                parse_mode="HTML",
                                disable_web_page_preview=True
                            )
                            sent = True
                            log.debug(f"[NOTIFY] ✓ أُرسل للقناة {ch_id}")

                        except FloodWait as e:
                            log.warning(
                                f"[NOTIFY] FloodWait {e.value}s "
                                f"على القناة {ch_id} ← تبديل"
                            )
                            flood_until[ch_id] = (
                                asyncio.get_event_loop().time() + e.value
                            )
                            if _notify_func:
                                try:
                                    await _notify_func(
                                        f"⚠️ FloodWait {e.value}s "
                                        f"على القناة <code>{ch_id}</code>"
                                    )
                                except Exception:
                                    pass

                        except Exception as e:
                            log.warning(
                                f"[NOTIFY CHANNEL] {ch_id} "
                                f"محاولة {attempts}: {e}"
                            )
                            # ← عاقب هذه القناة 30 ثانية عند أي خطأ
                            flood_until[ch_id] = (
                                asyncio.get_event_loop().time() + 30
                            )

                    if not sent:
                        log.warning(
                            f"[NOTIFY] فشل إرسال الرسالة لجميع القنوات "
                            f"بعد {max_attempts} محاولة"
                        )

                elif dest == "channel" and _notify_func:
                    # ← fallback للمطورين إذا لا توجد قنوات
                    try:
                        await _notify_func(text)
                    except Exception:
                        pass

                else:
                    # ← إرسال للمطورين
                    if _notify_func:
                        try:
                            await _notify_func(text)
                        except Exception:
                            pass

                # ← تأخير بين كل رسالة لتجنب FloodWait
                await asyncio.sleep(0.3)

            except asyncio.CancelledError:
                break
            except Exception as e:
                log.error(f"[NOTIFY WORKER] {e}")

    _notify_worker_task = asyncio.create_task(_worker())


async def notify(text: str, urgent: bool = False, channel: bool = False):
    """
    channel=True  → القناة عبر Queue (اشتراك، اكتمال، مغادرة)
    urgent=True   → المطورون فوراً (حظر، بدء/إيقاف، ملخص، أخطاء)
    افتراضي       → المطورون عبر Queue
    """
    if channel:
        ch_ids = getattr(config, "NOTIFY_CHANNEL_IDS", None)
        ch_id = getattr(config, "NOTIFY_CHANNEL_ID", None)
        has_channel = (ch_ids and len(ch_ids) > 0) or ch_id
        if has_channel and _notify_queue:
            await _notify_queue.put(("channel", text))
        elif _notify_func:
            try:
                await _notify_func(text)
            except Exception as e:
                log.warning(f"[NOTIFY FALLBACK] {e}")
    elif urgent:
        if _notify_func:
            try:
                return await _notify_func(text)
            except Exception as e:
                log.warning(f"[NOTIFY URGENT] {e}")
    else:
        if _notify_queue:
            await _notify_queue.put(("devs", text))

async def notify_edit(msg_ids: list, text: str):
    """تعديل رسالة موجودة بدل إرسال رسالة جديدة"""
    if _notify_edit_func:
        try:
            await _notify_edit_func(msg_ids, text)
        except Exception as e:
            log.warning(f"[NOTIFY EDIT] {e}")

# ============================================================
# التحقق من وجود البوت المستهدف
# ============================================================
async def verify_bot_exists(bot_un: str) -> bool:
    """
    التحقق من وجود البوت قبل بدء العملية.
    يستخدم أول client متاح، أو يفتح client مؤقتاً إذا لزم.
    """
    for phone, client in list(clients.items()):
        try:
            await client.get_chat(bot_un)
            log.info(f"[VERIFY BOT] {bot_un} موجود ✓")
            return True
        except Exception as e:
            err = str(e).upper()
            if "USERNAME_NOT_OCCUPIED" in err or "USERNAME_INVALID" in err:
                log.error(f"[VERIFY BOT] {bot_un} غير موجود!")
                return False
            continue

    # لا يوجد clients جاهزة → نفتح واحداً مؤقتاً
    accounts = await get_accounts(active_only=True)
    if not accounts:
        return True  # لا يمكن التحقق → افترض صحيح

    phone = accounts[0]["phone"]
    if phone in dead_sessions:
        return True  # جلسة معروفة كميتة → لا يمكن التحقق → افترض صحيح

    try:
        from pyrogram import Client as PyroClient
        c = PyroClient(
            f"sessions/{phone}",
            api_id=config.API_ID,
            api_hash=config.API_HASH,
            no_updates=True,
            sleep_threshold=60,
        )
        # نمسك قفل الفتح طوال عمر العميل المؤقت (فتح + استعمال + إيقاف)
        # حتى لا يتزامن فتح نفس ملف الجلسة مع get_client في أي مسار آخر.
        async with get_client_open_lock(phone):
            try:
                await c.start()
                try:
                    await c.get_chat(bot_un)
                    log.info(f"[VERIFY BOT] {bot_un} موجود ✓")
                    return True
                except Exception as e:
                    err = str(e).upper()
                    if "USERNAME_NOT_OCCUPIED" in err or "USERNAME_INVALID" in err:
                        log.error(f"[VERIFY BOT] {bot_un} غير موجود!")
                        return False
                    return True
            finally:
                try:
                    await c.stop()
                except Exception:
                    pass
    except Exception as e:
        log.warning(f"[VERIFY BOT] فشل التحقق: {e}")
        return True  # افترض صحيح عند الفشل


# ============================================================
# إدارة الـ Client
# ============================================================
def _patch_client_updates(c: Client):
    pass


async def _suppress_channel_sync(c: Client):
    try:
        from pyrogram.raw.functions.updates import GetChannelDifference

        class _FakeDiffEmpty:
            pts = 1
            timeout = 300
            final = True
            users = []
            chats = []
            new_messages = []
            other_updates = []

        original_invoke = c.invoke

        async def _patched_invoke(query, *args, **kwargs):
            if isinstance(query, GetChannelDifference):
                return _FakeDiffEmpty()
            return await original_invoke(query, *args, **kwargs)

        c.invoke = _patched_invoke
    except Exception as e:
        log.warning(f"[PATCH] failed: {e}")


async def get_client(phone, updates=False):
    if phone in clients:
        return clients[phone]

    c = None
    try:
        async with get_client_open_lock(phone):
            if phone in clients:
                return clients[phone]
            c = Client(
                f"sessions/{phone}",
                api_id=config.API_ID,
                api_hash=config.API_HASH,
                no_updates=True,
                sleep_threshold=120,
            )
            await asyncio.wait_for(c.start(), timeout=CLIENT_CONNECT_TIMEOUT)
            await asyncio.wait_for(c.get_me(), timeout=CLIENT_CONNECT_TIMEOUT)
            clients[phone] = c
            return c
        # ملاحظة: عند وقوع استثناء، تُحرَّر get_client_open_lock تلقائياً هنا
        # (خروج async with)، قبل الوصول لأي كود أدناه — هذا ضروري لأن
        # ban_account()/cleanup_client() تحاول إمساك نفس القفل، وasyncio.Lock
        # غير قابل لإعادة الدخول؛ استدعاؤها والقفل ممسوك يسبب تعليقاً دائماً.
    except Exception as e:
        # إغلاق العميل الفاشل — القفل محرَّر الآن فهذا آمن تماماً.
        if c is not None:
            try:
                await asyncio.wait_for(c.stop(), timeout=10)
            except Exception:
                try:
                    await asyncio.wait_for(c.disconnect(), timeout=10)
                except Exception:
                    pass

        async def _safe_ban(reason: str) -> bool:
            # مهلة زمنية خاصة لاستدعاء ban_account نفسه (قاعدة البيانات) —
            # حتى لو كانت مشغولة/بطيئة لحظياً، لن نتعلّق للأبد بانتظارها.
            try:
                return await asyncio.wait_for(ban_account(phone, reason), timeout=10)
            except Exception as ban_err:
                log.error(f"[CLIENT] {phone}: فشل حفظ الحظر في القاعدة: {ban_err}")
                return False

        if isinstance(e, asyncio.TimeoutError):
            # مهم: الـ timeout وحده دليل غامض — قد يعني جلسة تالفة فعلاً،
            # لكن قد يعني أيضاً ازدحاماً شبكياً مؤقتاً (متوقع عند فتح مئات
            # الجلسات دفعة واحدة). لذا لا نحظر حظراً دائماً (لا حذف من
            # القاعدة ولا حذف لملف الجلسة) بناءً على دليل غير مؤكد — فقط
            # نتخطّى هذه الجلسة لبقية هذا التشغيل (dead_sessions بالذاكرة)،
            # فتُعطى الجلسة فرصة أخرى تلقائياً عند إعادة تشغيل الخدمة القادمة.
            reason = "انتهت مهلة الاتصال (قد تكون مؤقتة/ازدحام شبكي — لم تُحظر نهائياً)"
            log.warning(f"[CLIENT] {phone}: {reason} — تخطٍ لهذا التشغيل فقط")
            dead_sessions.add(phone)
            raise
        err = str(e).lower()
        if is_ban_error(e):
            reason = get_ban_reason(e)
            is_new = await _safe_ban(reason)
            if is_new:
                await notify(
                    f"🚫 حساب محظور\n"
                    f"👤 الحساب: <code>{phone}</code>\n"
                    f"📋 السبب: {reason}",
                    urgent=True
                )
            raise UserDeactivatedBan(reason)
        if "eof when reading a line" in err:
            log.warning(f"[CLIENT] {phone}: session منتهية — تخطي بهدوء وحظر دائم")
            dead_sessions.add(phone)
            await _safe_ban("session تالفة (EOF)")
            raise
        if ("malformed" in err or "disk image" in err
                or "auth_key" in err or "unregistered" in err
                or "revoked" in err):
            # جلسة تالفة/ميتة — تُحظر دائماً في القاعدة فلا تُعاد محاولتها
            # حتى بعد إعادة تشغيل الخدمة.
            log.warning(f"[CLIENT] {phone}: جلسة تالفة/ميتة — تُحظر دائماً")
            dead_sessions.add(phone)
            await _safe_ban("session تالفة/ميتة")
            raise
        log.error(f"[CLIENT] {phone}: {e}")
        raise
# ============================================================
# Periodic Tasks
# ============================================================
async def periodic_cleanup():
    while True:
        await asyncio.sleep(600)  # ← 10 دقائق بدل 30
        active_phones = {
            acc["phone"] for acc in await get_accounts(active_only=True)
        }
        dead = [p for p in list(clients.keys()) if p not in active_phones]
        for phone in dead:
            await cleanup_client(phone)
        gc.collect()
        log.info(f"[CLEANUP] clients نشط: {len(clients)}")


async def health_check():
    # عدد الفشل المتتالي لكل عميل قبل إغلاقه — يمنع إغلاق العملاء بسبب
    # ازدحام عابر (fetch من مايو 2026: 20 موجة كانت تعلن عشرات العملاء
    # "ميتين" تحت الضغط ثم تقفلهم أثناء عملهم → فتح مزدوج لحرف الملف → إبطال مفاتيح جماعي).
    _fail = {}
    while True:
        await asyncio.sleep(300)  # كل 5 دقائق
        log.info(f"[HEALTH] فحص {len(clients)} client")
        for phone, client in list(clients.items()):
            # 1) لا نلمس حسابات مشغولة بمهمة حالياً (worker يحمل get_account_lock).
            #    تحت الضغط قد يتباطأ get_me ويعطي انطباع "ميت" — تخطِّها كلها.
            alock = get_account_lock(phone)
            try:
                await asyncio.wait_for(alock.acquire(), timeout=0.3)
            except asyncio.TimeoutError:
                continue  # مشغول — تركه هذا الدور
            try:
                try:
                    await asyncio.wait_for(client.get_me(), timeout=20)
                    _fail.pop(phone, None)
                except Exception as e:
                    _fail[phone] = _fail.get(phone, 0) + 1
                    log.warning(
                        f"[HEALTH] {phone} فشل {_fail[phone]} دور متتالي: {e}"
                    )
                    # 2) لا يُغلق العميل إلا بعد فشلين متتاليين — المرة الأولى
                    #    تكفي للتأكد مع احتمال ازدحام عابر؛ الثانية تؤكد الموت.
                    if _fail[phone] >= 2:
                        log.warning(f"[HEALTH] {phone} ميت — تنظيف")
                        _fail.pop(phone, None)
                        # 3) الهوية + دورة الفحص الحالية: لا يُغلق عميل حلّ غيره.
                        await cleanup_client(phone, client=client)
            finally:
                alock.release()
        log.info(f"[HEALTH] مجموع {len(_fail)} تحت المراقبة")
# ============================================================
# دوال المساعدة
# ============================================================

async def _flush_join_notify():
    while True:
        await asyncio.sleep(60)
        async with get_join_notify_lock():
            if not _join_notify_buffer:
                continue
            urls_snapshot = _join_notify_buffer.copy()
            _join_notify_buffer.clear()
        links = "\n".join(f"- {u}" for u in urls_snapshot)
        await notify(
            f"📢 <b>تقرير اشتراكات جديدة</b>\n"
            f"⚡️ قام {len(urls_snapshot)} حساب بالاشتراك بنجاح\n"
            f"🔗 <b>القنوات المشترك بها:</b>\n{links}",
            channel=True
        )


async def leave_ch(client: Client, phone: str, ch: dict, ctx="LEAVE") -> bool:
    targets = []
    if ch.get("id"):
        targets.append(ch["id"])
    if ch.get("username"):
        targets.append(ch["username"])
    link = ch_link(ch)
    if link and link not in targets:
        targets.append(link)
    for t in targets:
        try:
            await client.leave_chat(t)
            log.info(f"[{ctx}] {phone} left {t}")
            return True
        except (UserNotParticipant, PeerIdInvalid):
            return False
        except FloodWait as e:
            log.warning(f"[{ctx}] {phone} FloodWait {e.value}s على {t} — تخطي")
            continue
        except Exception as e:
            err = str(e).upper()
            if "CHANNEL_INVALID" in err or "CHANNEL_PRIVATE" in err:
                return False
            log.error(f"[{ctx}] {phone} error leaving {t}: {e}")
            continue
    return False


def _is_bot_target(target: str) -> bool:
    """التحقق إذا كان الهدف بوت وليس قناة"""
    t = target.lower().strip("/").split("/")[-1]
    return t.endswith("bot")


async def join_ch(client: Client, phone: str, target: str) -> str:
    if stop_event.is_set():
        return "failed"

    if is_channel_banned(target):
        log.info(f"[JOIN SKIP] {phone} → {target} محظورة - تخطي")
        return "skip"

    MAX_RETRIES = 3
    retry = 0
    while retry <= MAX_RETRIES:
        try:
            target_norm = norm_target(target)

            if _is_bot_target(target_norm):
                log.info(f"[SKIP_BOT] {phone} -> {target_norm}")
                return "skip_bot"

            try:
                member = await client.get_chat_member(target_norm, "me")
                if member:
                    return "ok"
            except Exception:
                pass

            try:
                await client.join_chat(target_norm)
            except Exception as e:
                if "INVITE_REQUEST_SENT" in str(e):
                    log.info(f"[WAIT_APPROVAL] {phone} -> {target}")
                    return "ok"
                raise e

            log.info(f"[JOIN] {phone} -> {target}")
            account_join_count[phone] = account_join_count.get(phone, 0) + 1
            if account_join_count[phone] >= 20:
                rest_time = random.uniform(180, 300)
                log.info(
                    f"[JOIN] {phone} | وصل للحد - استراحة {int(rest_time)} ثانية"
                )
                await safe_sleep(rest_time)
                account_join_count[phone] = 0

            url = (
                target if target.startswith("http")
                else f"https://t.me/{target_norm}"
            )
            try:
                info = await client.get_chat(target_norm)
                uname = getattr(info, "username", None)
                if uname:
                    url = f"https://t.me/{uname}"
                await add_channel({
                    "id": info.id,
                    "title": info.title,
                    "username": uname,
                    "url": url
                })
            except Exception:
                await add_channel({
                    "id": None,
                    "title": target,
                    "username": None,
                    "url": url
                })

            async with get_join_notify_lock():
                _join_notify_buffer.append(url)
                if len(_join_notify_buffer) >= 10:
                    urls_snapshot = _join_notify_buffer.copy()
                    _join_notify_buffer.clear()
                    links = "\n".join(f"- {u}" for u in urls_snapshot)
                    await notify(
                        f"📢 <b>تقرير اشتراكات جديدة</b>\n"
                        f"⚡️ قام 10 حسابات بالاشتراك بنجاح\n"
                        f"🔗 <b>القنوات المشترك بها:</b>\n{links}",
                        channel=True
                    )

            return "ok"

        except Exception as e:
            if isinstance(e, UserAlreadyParticipant):
                return "ok"
            if is_ban_error(e):
                reason = get_ban_reason(e)
                is_new = await ban_account(phone, reason)
                if is_new:
                    await notify(
                        f"🚫 حساب محظور\n"
                        f"👤 الحساب: <code>{phone}</code>\n"
                        f"📋 السبب: {reason}",
                        urgent=True
                    )
                return "banned"
            if isinstance(e, FloodWait):
                if retry < MAX_RETRIES:
                    await safe_sleep(e.value * 1.2)
                    retry += 1
                    continue
                else:
                    log.warning(f"[JOIN] {phone} تجاوز حد FloodWait - تخطي")
                    return "skip"
            err_str = str(e).upper()
            if any(k in err_str for k in [
                "USERNAME_NOT_OCCUPIED", "USERNAME_INVALID",
                "INVITE_HASH_INVALID", "USERNAME NOT FOUND",
                "USERNAME_NOT_FOUND"
            ]):
                log.warning(f"[JOIN] {phone} -> {target} | غير موجود - تخطي")
                return "skip"
            if any(k in err_str for k in [
                "CHANNEL_BANNED", "USER_BANNED_IN_CHANNEL",
                "CHAT_WRITE_FORBIDDEN", "CHAT_JOIN_FORBIDDEN",
                "YOU_BLOCKED_USER", "USER_BLOCKED"
            ]):
                log.warning(f"[JOIN] {phone} -> {target} | محظور من القناة - تخطي")
                return "skip"
            if any(k in err_str for k in [
                "CONNECTION", "CONNECTION LOST",
                "CONNECTION RESET", "BROKEN PIPE"
            ]):
                if retry < MAX_RETRIES:
                    await safe_sleep(random.uniform(5, 10))
                    retry += 1
                    continue
            log.error(f"[JOIN_ERROR] {phone} -> {target}: {e}")
            return "failed"
    return "failed"

async def get_msg(client: Client, bot_un: str, wait: float = 0, after_id: int = None):
    if stop_event.is_set():
        return None
    if wait > 0:
        await safe_sleep(wait)
    max_attempts = 5 if after_id else 1
    for attempt in range(max_attempts):
        try:
            msgs = []
            async for m in client.get_chat_history(bot_un, limit=5):
                msgs.append(m)
            msg = msgs[0] if msgs else None
            if after_id and msg and msg.id <= after_id:
                if attempt < max_attempts - 1:
                    await asyncio.sleep(2)
                    continue
            return msg
        except FloodWait as e:
            await safe_sleep(e.value * 1.2)
            continue
        except Exception as e:
            log.error(f"[MSG] {bot_un}: {e}")
            return None
    return None


# ============================================================
# دوال تحليل الرسائل والأزرار
# ============================================================
def get_chs_from_text(text: str) -> list:
    if not text:
        return []
    chs = []
    patterns = [
        r'https?://t\.me/(?:\+|joinchat/)[A-Za-z0-9_-]+',
        r'https?://t\.me/([A-Za-z0-9_]{4,})',
        r'@([A-Za-z0-9_]{4,})'
    ]
    for p in patterns:
        for m in re.finditer(p, text):
            found = m.group(0)
            if p == r'@([A-Za-z0-9_]{4,})' and found.lower().endswith('bot'):
                continue
            chs.append(found)
    return list(dict.fromkeys(chs))


def is_sub_msg(text: str) -> bool:
    keys = [
        "عليك الاشتراك", "اشترك", "اشتراك اجباري", "اشتراك إجباري",
        "لتتمكن من استخدامه", "أرسل /start", "ارسل /start",
        "عذراً", "عذرا", "اشترك ثم"
    ]
    return any(k in text for k in keys)


def is_no_channels_msg(text: str) -> bool:
    keys = [
        "لا يوجد قنوات", "لا توجد قنوات", "انتهت القنوات",
        "لا يوجد", "ليس لديك", "اكملت", "أكملت"
    ]
    return any(k in text for k in keys)


def is_not_subscribed_msg(text: str) -> bool:
    return any(k in text for k in ["لم تشترك", "تشتر", "لم يتم التحقق"])


def is_name_change(text: str) -> bool:
    return any(k in text for k in [
        "تغيير اسم", "غيّر اسم", "غير اسم", "الاسم التالي", "اسم حسابك"
    ])


def get_req_name(text: str):
    if not text:
        return None
    keywords = ["الاسم التالي", "اسم حسابك", "تغيير اسم حسابك إلى"]
    lines = [l.strip() for l in text.splitlines() if l.strip()]
    for i, line in enumerate(lines):
        if any(k in line for k in keywords):
            if ":" in line or "：" in line:
                parts = re.split(r'[:：]', line, maxsplit=1)
                if len(parts) > 1:
                    c = parts[1].strip().strip("•-").strip()
                    if c and not any(x in c for x in [
                        "قم بتغيير", "استكمال", "لتجمع", "اضغط"
                    ]):
                        return c
            if i + 1 < len(lines):
                c = lines[i + 1].strip().strip("•-").strip()
                if c and not any(x in c for x in [
                    "قم بتغيير", "استكمال", "لتجمع", "اضغط"
                ]):
                    return c
    return None


def parse_msg(msg):
    text = ""
    btns = []
    if msg:
        text = msg.text or msg.caption or ""
        if msg.reply_markup and hasattr(msg.reply_markup, "inline_keyboard"):
            btns = msg.reply_markup.inline_keyboard
    return text, btns


def btn_10x(btns):
    for row in (btns or []):
        for b in row:
            if "10x" in (getattr(b, "text", "") or ""):
                return b
    return None


def btn_points(btns):
    for row in (btns or []):
        for b in row:
            t = getattr(b, "text", "") or ""
            if (
                "قسم تجميع النقاط" in t
                or "تجميع النقاط" in t
                or "تجميع نقاط" in t
            ):
                return b
    return None



def btn_turbo(btns):
    target_texts = ["تيربو", "تجميع السريع", "التجميع السريع", "سريع تجميع"]
    for row in (btns or []):
        for b in row:
            btn_text = getattr(b, "text", "") or ""
            if any(text in btn_text for text in target_texts):
                return b
    return None


def btn_fast(btns):
    for row in (btns or []):
        for b in row:
            if "تجميع بسرعة" in (getattr(b, "text", "") or ""):
                return b
    return None


def btn_complete(btns):
    for row in (btns or []):
        for b in row:
            t = getattr(b, "text", "") or ""
            if "استكمال التجميع" in t or "استكمال" in t:
                return b
    return None


def btn_verify(btns):
    for row in (btns or []):
        for b in row:
            if "تحقق" in (getattr(b, "text", "") or ""):
                return b
    return None


def btn_next(btns):
    for row in (btns or []):
        for b in row:
            t = getattr(b, "text", "") or ""
            if "التالي" in t or "تالي" in t:
                return b
    return None


def btn_retry(btns):
    for row in (btns or []):
        for b in row:
            t = getattr(b, "text", "") or ""
            if "اعد المحاول" in t or "أعد المحاول" in t or "retry" in t.lower():
                return b
    return None


def get_ch_btns(btns) -> list:
    result = []
    exclude_terms = [
        "تجميع", "تيربو", "10x", "رجوع", "حسابي", "كود",
        "نجوم", "بنجوم", "اشتراك اجباري بنجوم",  # ← إضافة
        "شتراك اجباري بنجوم"  # ← إضافة
    ]
    for row in (btns or []):
        for b in row:
            url = getattr(b, "url", None)
            t = (getattr(b, "text", "") or "").strip()
            if url and ("t.me/" in url or "telegram.me/" in url):
                if not any(term in t for term in exclude_terms):
                    result.append(b)
    return result


def is_sub_by_btns(btns) -> bool:
    if not get_ch_btns(btns):
        return False
    return not (btn_10x(btns) or btn_points(btns) or btn_turbo(btns) or btn_fast(btns))


def btns_hash(btns) -> str:
    try:
        keys = []
        for row in (btns or []):
            for b in row:
                keys.append(getattr(b, "text", "") or getattr(b, "url", "") or "")
        return "|".join(keys)
    except Exception:
        return ""


# ============================================================
# ضغط الأزرار
# ============================================================
async def click_btn(client: Client, msg, btn):
    phone = getattr(client, "phone", "unknown")
    bot_un = getattr(msg, "chat", None)
    bot_un = bot_un.username if bot_un else "unknown"
    cb_data = btn.callback_data
    if isinstance(cb_data, bytes):
        cb_preview = cb_data[:32].decode(errors="ignore")
    else:
        cb_preview = str(cb_data)[:32]
    log.info(f"[CLICK] {phone} @{bot_un} ▶ زر: {cb_preview}")
    if stop_event.is_set():
        return None
    for attempt in range(2):
        try:
            from pyrogram.raw.functions.messages import GetBotCallbackAnswer
            peer = await client.resolve_peer(msg.chat.id)
            cb = btn.callback_data
            if isinstance(cb, str):
                cb = cb.encode()
            resp = await asyncio.wait_for(
                client.invoke(GetBotCallbackAnswer(
                    peer=peer, msg_id=msg.id, data=cb, game=False
                )),
                timeout=30
            )
            log.info(f"[CLICK] {phone} @{bot_un} ✓ نجاح (محاولة {attempt+1})")
            return resp
        except asyncio.TimeoutError:
            log.warning(f"[CLICK] {phone} @{bot_un} ⏱ Timeout (محاولة {attempt + 1})")
            if attempt == 0:
                await asyncio.sleep(5)
                continue
            return "TIMEOUT"
        except FloodWait as e:
            log.warning(f"[CLICK] {phone} @{bot_un} ⚠️ FloodWait {e.value}s")
            await safe_sleep(e.value * 1.2)
            return None
        except Exception as e:
            err_str = str(e).upper()
            if "DATA_INVALID" in err_str:
                log.warning(f"[CLICK] {phone} @{bot_un} ■ DATA_INVALID")
                return "DATA_INVALID"
            if "TIMEOUT" in err_str:
                log.warning(f"[CLICK] {phone} @{bot_un} ■ TIMEOUT")
                return "TIMEOUT"
            log.error(f"[CLICK] {phone} @{bot_un} ✗ خطأ: {e}")
            return None

async def join_from_btns(client, phone, btns, max_channels=15) -> str:
    log.info(f"[JOIN_BATCH] {phone} ▶ انضمام دفعة ({max_channels} قناة كحد أقصى)")
    if stop_event.is_set():
        return "ok"
    ch_btns = get_ch_btns(btns)[:max_channels]
    if not ch_btns:
        log.debug(f"[JOIN_BATCH] {phone} ■ لا توجد أزرار قنوات")
        return "ok"
    spd = await get_speed()
    bot_count = 0
    skip_count = 0
    total = len(ch_btns)
    log.info(f"[JOIN_BATCH] {phone} إجمالي: {total} هدف، spd={spd}")

    for i, b in enumerate(ch_btns):
        if stop_event.is_set():
            break
        if b.url:
            try:
                t = norm_target(b.url)
                log.debug(f"[JOIN_BATCH] {phone} [{i+1}/{total}] → {t}")
                result = await join_ch(client, phone, t)
                if result == "skip_bot":
                    bot_count += 1
                elif result in ("skip", "failed"):
                    skip_count += 1
                elif result == "banned":
                    log.warning(f"[JOIN_BATCH] {phone} ■ محظور أثناء الانضمام")
                    return "banned"
            except Exception as e:
                log.error(f"[JOIN_ERROR] {phone} -> {b.url}: {e}")
                skip_count += 1
        if i < total - 1:
            await safe_sleep(rand_delay(spd))

    if bot_count == total:
        log.info(f"[JOIN] {phone} | كل الأهداف بوتات - اعتبار الحساب اكتمل")
        return "all_bots"
    if skip_count == total:
        log.info(f"[JOIN] {phone} | كل الأهداف غير متاحة - اعتبار الحساب اكتمل")
        return "all_skip"
    log.info(f"[JOIN_BATCH] {phone} ✓ اكتمل (bots={bot_count}, skipped={skip_count})")
    return "ok"


# ============================================================
# دوال مساعدة للاشتراك الإجباري
# ============================================================
def _detect_loop(targets: list, sub_seen: dict) -> list:
    looped = []
    for t in targets:
        key = norm_ch_key(t)
        sub_seen[key] = sub_seen.get(key, 0) + 1
        if sub_seen[key] > 1:
            looped.append(norm_target(t))
    return looped


async def _handle_loop(client, phone, looped, sub_seen, spd):
    for t in looped:
        if stop_event.is_set():
            return
        try:
            await client.leave_chat(t)
            await safe_sleep(rand_delay(spd))
        except Exception:
            pass
        await safe_sleep(rand_delay(BOT_DELAY))
        await join_ch(client, phone, t)
        await safe_sleep(rand_delay(spd))
        sub_seen[norm_ch_key(t)] = 1

async def _start_bot_if_needed(client: Client, phone: str, url: str) -> bool:
    try:
        username = norm_target(url)
        if not _is_bot_target(username):
            return False
        try:
            history = await client.get_chat_history(username, limit=1)
            if history:
                return False
        except Exception:
            pass
        await client.send_message(username, "/start")
        log.info(f"[BOT START] {phone} → {username} تم إرسال /start")
        await safe_sleep(rand_delay(BOT_DELAY))
        return True
    except Exception as e:
        log.warning(f"[BOT START] {phone} → {url}: {e}")
        return False

# ============================================================
# معالجة تغيير الاسم
# ============================================================
async def handle_name(client: Client, phone: str, bot_un: str, msg):
    if stop_event.is_set():
        return msg, [], False
    if not msg:
        return msg, [], False
    text, btns = parse_msg(msg)
    if not is_name_change(text):
        return msg, btns, False
    req_name = get_req_name(text)
    if not req_name:
        return msg, btns, False
    updated = False
    for attempt in range(3):
        if stop_event.is_set():
            return msg, btns, False
        try:
            await client.update_profile(first_name=req_name, last_name="")
            updated = True
            break
        except FloodWait as e:
            await safe_sleep(e.value * 1.2)
        except Exception as e:
            log.error(f"[NAME] {attempt + 1} {phone}: {e}")
            await safe_sleep(rand_delay(BOT_DELAY))
    if not updated:
        return msg, btns, False
    await safe_sleep(rand_delay(BOT_DELAY))
    comp = btn_complete(btns)
    if not comp:
        ref = await get_msg(client, bot_un, after_id=msg.id)
        if ref:
            msg = ref
            text, btns = parse_msg(msg)
            comp = btn_complete(btns)
    if comp:
        result = await click_btn(client, msg, comp)
        if result in ("DATA_INVALID", "TIMEOUT"):
            ref = await get_msg(client, bot_un, after_id=msg.id)
            if ref:
                msg = ref
                text, btns = parse_msg(msg)
    else:
        await safe_sleep(rand_delay(BOT_DELAY))
        ref = await get_msg(client, bot_un, after_id=msg.id)
        if ref:
            msg = ref
            text, btns = parse_msg(msg)
    return msg, btns, True
# ============================================================
# جولة التجميع الرئيسية
# ============================================================
async def do_gather_round(
    client: Client, phone: str, bot_un: str, msg, btns
) -> str:
    log.info(f"[ROUND] {phone} @{bot_un} ▶ بدء جولة التجميع")
    if phone in flood_map:
        if time.time() < flood_map[phone]:
            log.warning(f"[ROUND] {phone} @{bot_un} ■ Skip - في flood_map")
            return "ok"
        else:
            del flood_map[phone]

    MAX_INNER = 10
    MAX_VERIFY_FAIL = 3
    inner = 0
    verify_fail = 0
    current_msg, current_btns = msg, btns
    last_hash = ""
    spd = await get_speed()
    log.debug(f"[ROUND] {phone} @{bot_un} spd={spd}")

    while inner < MAX_INNER:
        if stop_event.is_set():
            return "stopped"
        inner += 1
        log.debug(f"[ROUND] {phone} @{bot_un} inner={inner}/{MAX_INNER}")

        text, _ = parse_msg(current_msg)

        if is_no_channels_msg(text):
            log.info(f"[ROUND] {phone} @{bot_un} ✓ لا يوجد قنوات → done")
            return "done"

        ch_btns = get_ch_btns(current_btns)
        current_hash = btns_hash(current_btns)

        if current_hash == last_hash and last_hash != "" and not ch_btns:
            log.info(f"[GATHER] {phone} | لا تغيير في الأزرار - خروج")
            return "ok"
        last_hash = current_hash

        if ch_btns:
            log.info(f"[GATHER] {phone} | اشتراك في {len(ch_btns)} قناة")
            join_result = await join_from_btns(client, phone, current_btns)
            if join_result in ("all_bots", "all_skip"):
                return "done"
            if join_result == "banned":
                return "banned"
            await safe_sleep(rand_delay(BOT_DELAY))

        vb = btn_verify(current_btns)
        if vb:
            log.info(f"[GATHER] {phone} | ضغط تحقق")
            await click_btn(client, current_msg, vb)
            await safe_sleep(rand_delay(BOT_DELAY) + 1)
            new_msg = await get_msg(client, bot_un, after_id=current_msg.id)
            if new_msg:
                current_msg = new_msg
                text, current_btns = parse_msg(current_msg)

                if is_no_channels_msg(text):
                    log.info(f"[GATHER] {phone} | تحقق ناجح - لا يوجد قنوات")
                    return "done"

                if is_not_subscribed_msg(text):
                    verify_fail += 1
                    log.warning(
                        f"[GATHER] {phone} | لم تشترك ({verify_fail}/{MAX_VERIFY_FAIL})"
                    )
                    if verify_fail >= MAX_VERIFY_FAIL:
                        log.warning(
                            f"[GATHER] {phone} | تجاوز حد فشل التحقق - خروج"
                        )
                        return "ok"

                    rb = btn_retry(current_btns)
                    if rb:
                        log.info(f"[GATHER] {phone} | ضغط أعد المحاولة")
                        await click_btn(client, current_msg, rb)
                        await safe_sleep(rand_delay(BOT_DELAY))
                        retry_msg = await get_msg(
                            client, bot_un, after_id=current_msg.id
                        )
                        if retry_msg:
                            current_msg = retry_msg
                            _, current_btns = parse_msg(current_msg)
                            retry_chs = get_ch_btns(current_btns)
                            if retry_chs:
                                log.info(
                                    f"[GATHER] {phone} | "
                                    f"مغادرة وإعادة انضمام {len(retry_chs)} قناة"
                                )
                                for b in retry_chs:
                                    if b.url and not stop_event.is_set():
                                        t = norm_target(b.url)
                                        try:
                                            await client.leave_chat(t)
                                            await safe_sleep(1)
                                            await join_ch(client, phone, t)
                                            await safe_sleep(rand_delay(spd))
                                        except Exception:
                                            pass

                                await safe_sleep(rand_delay(BOT_DELAY))
                                vb2 = btn_verify(current_btns)
                                if vb2:
                                    log.info(
                                        f"[GATHER] {phone} | "
                                        f"ضغط تحقق بعد إعادة الانضمام"
                                    )
                                    await click_btn(client, current_msg, vb2)
                                    await safe_sleep(rand_delay(BOT_DELAY) + 1)
                                    # ← البوت يعدّل نفس الرسالة بدل إرسال جديدة
                                    new_msg2 = await get_msg(client, bot_un)
                                    if new_msg2:
                                        _t2, _ = parse_msg(new_msg2)
                                    if new_msg2:
                                        current_msg = new_msg2
                                        text2, current_btns = parse_msg(current_msg)
                                        if is_no_channels_msg(text2):
                                            log.info(
                                                f"[GATHER] {phone} | "
                                                f"تحقق ناجح بعد إعادة الانضمام"
                                            )
                                            return "done"

                                        # ← إذا لم تشترك بعد المرة الثانية
                                        # → حظر القناة واعتبار الحساب مكتملاً
                                        if is_not_subscribed_msg(text2):
                                            stuck_chs = get_ch_btns(current_btns)
                                            for b in stuck_chs:
                                                if b.url:
                                                    await ban_channel(b.url)
                                                    log.warning(
                                                        f"[GATHER] {phone} | "
                                                        f"قناة معطلة محظورة: {b.url}"
                                                    )
                                            return "done"
                    continue

                verify_fail = 0
                continue

        nb = btn_next(current_btns)
        if nb:
            await click_btn(client, current_msg, nb)
            await safe_sleep(rand_delay(BOT_DELAY))
            new_msg = await get_msg(client, bot_un, after_id=current_msg.id)
            if new_msg:
                current_msg = new_msg
                _, current_btns = parse_msg(current_msg)
                continue

        if (not get_ch_btns(current_btns)
                and not btn_verify(current_btns)
                and not btn_next(current_btns)):
            return "ok"

    return "ok"


# ============================================================
# تشغيل حساب واحد
# ============================================================
async def run_account_once(client: Client, phone: str, bot_un: str) -> str:
    log.info(f"[RUN_ACCOUNT] {phone} @{bot_un} ▶ بدء المعالجة")
    if stop_event.is_set():
        log.info(f"[RUN_ACCOUNT] {phone} @{bot_un} ■ stop_event")
        return "stopped"

    try:
        from pyrogram.raw.functions.updates import GetState
        await client.invoke(GetState())
    except Exception as e:
        if is_ban_error(e):
            reason = get_ban_reason(e)
            is_new = await ban_account(phone, reason)
            if is_new:
                await notify(
                    f"🚫 حساب محظور\n"
                    f"👤 الحساب: <code>{phone}</code>\n"
                    f"📋 السبب: {reason}",
                    urgent=True
                )
            log.warning(f"[RUN_ACCOUNT] {phone} @{bot_un} ■ محظور: {reason}")
            return "banned"
        log.error(f"[RUN_ACCOUNT] {phone} @{bot_un} ■ خطأ GetState: {e}")
        return "error"

    async def do_start():
        attempt = 0
        sub_seen = {}
        MAX_ATTEMPTS = 20
        MAX_SAME_CH = 3

        while attempt < MAX_ATTEMPTS:
            if stop_event.is_set():
                return None
            attempt += 1
            try:
                log.info(f"[GATHER] {phone} | إرسال /start للمحاولة {attempt}")
                await client.send_message(bot_un, "/start")
                await safe_sleep(rand_delay(BOT_DELAY) + random.uniform(3, 5))
                msg = await get_msg(client, bot_un)
                if not msg:
                    await safe_sleep(rand_delay(BOT_DELAY))
                    continue
                text, btns = parse_msg(msg)

                if is_sub_msg(text):
                    log.info(f"[GATHER] {phone} | اشتراك إجباري مطلوب")
                    chs = get_chs_from_text(text)
                    for b in get_ch_btns(btns):
                        if b.url and b.url not in chs:
                            chs.append(b.url)
                    chs = chs[:15]
                    log.info(f"[SUB LIST] {phone} | القنوات المكتشفة: {chs}")
                    looped = _detect_loop(chs, sub_seen)
                    if looped:
                        over_limit = [
                            t for t in looped
                            if sub_seen.get(norm_ch_key(t), 0) >= MAX_SAME_CH
                        ]
                        if over_limit:
                            log.warning(
                                f"[GATHER] {phone} | "
                                f"قناة عالقة {over_limit} - تخطي واكتمال"
                            )
                            return msg, btns, "done"
                        _vb = btn_verify(btns)
                        if _vb:
                            log.info(
                                f"[GATHER] {phone} | اشتراك إجباري مكرر "
                                f"(قنوات طلب انضمام) - ضغط تحقق ثم إعادة /start"
                            )
                            await click_btn(client, msg, _vb)
                            await safe_sleep(rand_delay(BOT_DELAY) + 1)
                            continue
                        await _handle_loop(
                            client, phone, looped, sub_seen,
                            spd=await get_speed()
                        )
                        await safe_sleep(rand_delay(BOT_DELAY))
                        continue
                    spd = await get_speed()
                    for i, t in enumerate(chs):
                        if stop_event.is_set():
                            return None
                        if _is_bot_target(norm_target(t)):
                            await _start_bot_if_needed(client, phone, t)
                        else:
                            r = await join_ch(client, phone, norm_target(t))
                            if r == "banned":
                                return None
                            _chs = await get_channels()
                            _ch_id = next((c["id"] for c in _chs
                                          if norm_ch_key(c.get("url","")) == norm_ch_key(t)), None)
                            await add_force_sub_channel(_extract_bot_username(bot_un), t, _ch_id)
                        if i < len(chs) - 1:
                            await safe_sleep(rand_delay(spd))
                    await safe_sleep(rand_delay(BOT_DELAY))
                    continue

                if is_sub_by_btns(btns):
                    log.info(f"[GATHER] {phone} | اشتراك إجباري عبر أزرار")
                    ch_btns = get_ch_btns(btns)[:15]
                    urls = [b.url for b in ch_btns if b.url]
                    looped = _detect_loop(urls, sub_seen)
                    if looped:
                        over_limit = [
                            t for t in looped
                            if sub_seen.get(norm_ch_key(t), 0) >= MAX_SAME_CH
                        ]
                        if over_limit:
                            log.warning(
                                f"[GATHER] {phone} | "
                                f"قناة عالقة {over_limit} - تخطي واكتمال"
                            )
                            return msg, btns, "done"
                        _vb = btn_verify(btns)
                        if _vb:
                            log.info(
                                f"[GATHER] {phone} | اشتراك إجباري مكرر "
                                f"(قنوات طلب انضمام) - ضغط تحقق ثم إعادة /start"
                            )
                            await click_btn(client, msg, _vb)
                            await safe_sleep(rand_delay(BOT_DELAY) + 1)
                            continue
                        await _handle_loop(
                            client, phone, looped, sub_seen,
                            spd=await get_speed()
                        )
                        await safe_sleep(rand_delay(BOT_DELAY))
                        continue
                    spd = await get_speed()
                    for i, b in enumerate(ch_btns):
                        if stop_event.is_set():
                            return None
                        if b.url:
                            if _is_bot_target(norm_target(b.url)):
                                await _start_bot_if_needed(client, phone, b.url)
                            else:
                                r = await join_ch(client, phone, norm_target(b.url))
                                if r == "banned":
                                    return None
                                _chs = await get_channels()
                                _ch_id = next((c["id"] for c in _chs
                                              if norm_ch_key(c.get("url","")) == norm_ch_key(b.url)), None)
                                await add_force_sub_channel(_extract_bot_username(bot_un), b.url, _ch_id)
                            if i < len(ch_btns) - 1:
                                await safe_sleep(rand_delay(spd))
                    await safe_sleep(rand_delay(BOT_DELAY))
                    continue

                msg, btns, _ = await handle_name(client, phone, bot_un, msg)

                if btn_fast(btns):
                    log.info(f"[GATHER] {phone} | نوع البوت: fast")
                    return msg, btns, "fast"
                if btn_points(btns):
                    log.info(f"[GATHER] {phone} | نوع البوت: points")
                    return msg, btns, "points"

                text, _ = parse_msg(msg)
                if is_no_channels_msg(text):
                    log.info(f"[GATHER] {phone} | لا يوجد قنوات من البداية")
                    return msg, btns, "done"

                log.info(f"[GATHER] {phone} | حالة غير معروفة - إعادة")
                await safe_sleep(rand_delay(BOT_DELAY) + 2)

            except Exception as e:
                if is_ban_error(e):
                    reason = get_ban_reason(e)
                    is_new = await ban_account(phone, reason)
                    if is_new:
                        await notify(
                            f"🚫 حساب محظور\n"
                            f"👤 الحساب: <code>{phone}</code>\n"
                            f"📋 السبب: {reason}",
                            urgent=True
                        )
                    return None
                if isinstance(e, FloodWait):
                    wait_sec = e.value
                    resume_sec = wait_sec + 60
                    flood_map[phone] = time.time() + resume_sec
                    await notify(
                        f"⚠️ حساب تعرض لـ FloodWait\n"
                        f"👤 الحساب: <code>{phone}</code>\n"
                        f"⏰ وقت الحظر: {wait_sec} ثانية\n"
                        f"⌛️ سيبدأ عمله بعد: {resume_sec} ثانية",
                        urgent=True
                    )
                    await safe_sleep(resume_sec)
                    flood_map.pop(phone, None)

                else:
                    err_str = str(e).upper()
                    if "USERNAME_NOT_OCCUPIED" in err_str:
                        log.warning(
                            f"[GATHER] {phone} | USERNAME_NOT_OCCUPIED - "
                            f"الحساب محظور/مجمد/محذوف (البوت تم التحقق منه مسبقاً)"
                        )
                        reason = "محظور/مجمد/محذوف"
                        is_new = await ban_account(phone, reason)
                        if is_new:
                            await notify(
                                f"🚫 حساب محظور/مجمد/محذوف\n"
                                f"👤 الحساب: <code>{phone}</code>\n"
                                f"📋 السبب: {reason}",
                                urgent=True
                            )
                        return None
                    else:
                        log.error(f"[GATHER] do_start {phone}: {e}")
                        await safe_sleep(5)

        log.warning(f"[GATHER] {phone} | فشل بعد {MAX_ATTEMPTS} محاولة")
        return None

    async def do_cycle(msg, btns, bot_type) -> str:
        MAX_CYCLES = 10

        if bot_type == "fast":
            fb = btn_fast(btns)
            if not fb:
                return "ok"
            log.info(f"[GATHER] {phone} | ضغط تجميع بسرعة")
            result = await click_btn(client, msg, fb)
            if stop_event.is_set():
                return "stopped"
            await safe_sleep(rand_delay(BOT_DELAY) + (3 if result == "TIMEOUT" else 0))
            new_msg = await get_msg(client, bot_un, after_id=msg.id)
            if not new_msg:
                await safe_sleep(3)
                new_msg = await get_msg(client, bot_un)
            if not new_msg:
                return "ok"
            new_msg, new_btns, _ = await handle_name(client, phone, bot_un, new_msg)
            _, new_btns = parse_msg(new_msg)
            if stop_event.is_set():
                return "stopped"
            if get_ch_btns(new_btns):
                log.info(f"[GATHER] {phone} | قنوات ظهرت بعد fast")
                return await do_gather_round(client, phone, bot_un, new_msg, new_btns)
            text, _ = parse_msg(new_msg)
            if is_no_channels_msg(text):
                return "done"
            return "ok"

        if bot_type == "points":
            cycle = 0
            pb = btn_points(btns)
            if not pb:
                return "ok"
            log.info(f"[GATHER] {phone} | ضغط تجميع النقاط")
            await click_btn(client, msg, pb)
            if stop_event.is_set():
                return "stopped"
            await safe_sleep(rand_delay(BOT_DELAY))
            msg = await get_msg(client, bot_un, after_id=msg.id)
            if not msg:
                return "ok"
            msg, btns, _ = await handle_name(client, phone, bot_un, msg)
            _, btns = parse_msg(msg)

            xb = btn_10x(btns) or btn_turbo(btns)
            if not xb:
                log.warning(f"[GATHER] {phone} | لا يوجد زر تيربو")
                return "ok"
            log.info(f"[GATHER] {phone} | ضغط تيربو")
            await click_btn(client, msg, xb)
            if stop_event.is_set():
                return "stopped"
            await safe_sleep(rand_delay(BOT_DELAY))
            msg = await get_msg(client, bot_un, after_id=msg.id)
            if not msg:
                return "ok"
            _, btns = parse_msg(msg)

            text, _ = parse_msg(msg)
            if is_no_channels_msg(text):
                log.info(f"[GATHER] {phone} | لا يوجد قنوات بعد تيربو - done")
                return "done"

            if get_ch_btns(btns):
                round_result = await do_gather_round(client, phone, bot_un, msg, btns)
                if round_result in ("done", "stopped"):
                    return round_result
                await safe_sleep(rand_delay(BOT_DELAY))
                msg = await get_msg(client, bot_un)
                if not msg:
                    return "ok"
                msg, btns, _ = await handle_name(client, phone, bot_un, msg)
                _, btns = parse_msg(msg)

            while cycle < MAX_CYCLES:
                if stop_event.is_set():
                    return "stopped"
                cycle += 1
                pb = btn_points(btns)
                if not pb:
                    return "ok"
                log.info(f"[GATHER] {phone} | دورة {cycle} - ضغط تجميع النقاط")
                await click_btn(client, msg, pb)
                if stop_event.is_set():
                    return "stopped"
                await safe_sleep(rand_delay(BOT_DELAY))
                msg = await get_msg(client, bot_un, after_id=msg.id)
                if not msg:
                    continue
                msg, btns, _ = await handle_name(client, phone, bot_un, msg)
                _, btns = parse_msg(msg)

                xb = btn_10x(btns) or btn_turbo(btns)
                if not xb:
                    return "ok"
                log.info(f"[GATHER] {phone} | دورة {cycle} - ضغط تيربو")
                await click_btn(client, msg, xb)
                if stop_event.is_set():
                    return "stopped"
                await safe_sleep(rand_delay(BOT_DELAY))
                msg = await get_msg(client, bot_un, after_id=msg.id)
                if not msg:
                    continue
                _, btns = parse_msg(msg)

                if get_ch_btns(btns):
                    round_result = await do_gather_round(
                        client, phone, bot_un, msg, btns
                    )
                    if round_result in ("done", "stopped"):
                        return round_result
                    await safe_sleep(rand_delay(BOT_DELAY))
                    msg = await get_msg(client, bot_un)
                    if not msg:
                        continue
                    msg, btns, _ = await handle_name(client, phone, bot_un, msg)
                    _, btns = parse_msg(msg)

                text, _ = parse_msg(msg)
                if is_no_channels_msg(text):
                    return "done"

            return "ok"

        return "done"

    result = await do_start()
    if result is None:
        log.warning(f"[RUN_ACCOUNT] {phone} @{bot_un} ■ do_start أعاد None → failed")
        return "failed"
    msg, btns, bot_type = result
    cycle_result = await do_cycle(msg, btns, bot_type)
    log.info(f"[RUN_ACCOUNT] {phone} @{bot_un} ■ انتهى بالنتيجة: {cycle_result}")
    return cycle_result
# ============================================================
# نظام 1: بوت واحد — Queue Worker Pool
# ============================================================
async def run_gather(bot_un: str):
    await start_notify_worker()

    if not await verify_bot_exists(bot_un):
        await notify(
            f"⚠️ البوت @{bot_un} غير موجود أو username خاطئ!\n"
            f"يرجى مراجعة البوت المضاف.",
            urgent=True
        )
        await set_setting("is_active", "0")
        await set_setting("active_bot", "")
        return

    accounts = await get_accounts(active_only=True)
    concurrent = await get_concurrent()

    if not accounts:
        await notify("⚠️ لا توجد حسابات نشطة!", urgent=True)
        await set_setting("is_active", "0")
        await set_setting("active_bot", "")
        return

    await notify(
        f"🚀 بدأت عملية التجميع\n"
        f"🤖 البوت: @{bot_un}\n"
        f"👥 الحسابات: {len(accounts)}\n"
        f"⚡ التزامن: {concurrent}",
        urgent=True
    )
    log.info(f"[RUN_GATHER] @{bot_un} ▶ بدء | حسابات={len(accounts)} concurrent={concurrent}")

    # تفعيل مهام التنظيف الدورية: تكتشف وتغلق أي عميل عالق تلقائياً
    # طوال فترة التجميع، فلا تتراكم حلقات إعادة محاولة محمومة.
    _bg_cleanup = asyncio.create_task(periodic_cleanup())
    _bg_health = asyncio.create_task(health_check())

    account_queue = asyncio.Queue()
    for acc in accounts:
        await account_queue.put(acc)
    opened, failed = 0, 0
    total = len(accounts)
    # إرسال الرسالة الأولى وحفظ الـ IDs
    msg_ids = await notify(f"⏳ عزيزي المطور جارٍ فتح {total} جلسة...\n✅ تم فتح حتى الان: 0 جلسة.", urgent=True)
    for acc in accounts:
        if stop_event.is_set():
            break
        phone = acc["phone"]
        try:
            await get_client(phone)
            opened += 1
            log.info(f"[WARMUP] @{bot_un} ✓ {phone} ({opened}/{total})")
        except Exception as e:
            failed += 1
            log.warning(f"[WARMUP] ✗ {phone}: {e}")
        await notify_edit(msg_ids or [], f"⏳ عزيزي المطور جارٍ فتح {total} جلسة...\n✅ تم فتح حتى الان: {opened} جلسة.")
        await asyncio.sleep(1)
    await notify(
        f"✅ اكتمل فتح الجلسات\n"
        f"✓ مفتوحة: {opened} | ✗ فاشلة: {failed}\n"
        f"🚀 بدء التجميع الآن...",
        urgent=True
    )

    # ── تتبّع backoff لكل حساب + كشف الانقطاع الشبكي العام ──
    fail_count: dict = {}          # عدد مرات الفشل المتتالي لكل حساب
    BACKOFF_BASE = 5               # ثانية: نقطة البداية (نفس القيمة الثابتة السابقة)
    BACKOFF_CAP = 300              # ثانية: السقف الأقصى للتأخير لكل حساب (5 دقائق)
    net_fail_window: list = []     # طوابع زمنية لفشل شبكي حديث عبر الحسابات
    NET_FAIL_THRESHOLD = 4         # عدد الإخفاقات الشبكية المتتالية لاعتبارها انقطاعاً عاماً
    NET_FAIL_WINDOW = 30           # ثانية: نافذة رصد الانقطاع
    SYSTEM_PAUSE = 60              # ثانية: هدنة على مستوى النظام عند رصد انقطاع عام
    _net_pause_until = [0.0]       # قائمة لتمرير القيمة بالمرجع بين الـ workers

    def _is_network_error(e) -> bool:
        s = str(e).lower()
        return any(k in s for k in [
            "too many open files", "errno 24",
            "connection", "cannot connect", "network",
            "timed out", "timeout", "dns", "ssl",
        ])

    def _account_backoff(phone: str) -> float:
        n = fail_count.get(phone, 0)
        if n <= 0:
            return BACKOFF_BASE
        return min(BACKOFF_BASE * (2 ** (n - 1)), BACKOFF_CAP)

    async def worker(worker_id: int):
        while not stop_event.is_set():
            try:
                # هدنة عامة إن رُصد انقطاع شبكي يطال كل الحسابات
                now = time.time()
                if now < _net_pause_until[0]:
                    await asyncio.sleep(min(5, _net_pause_until[0] - now))
                    continue
                try:
                    acc = account_queue.get_nowait()
                except asyncio.QueueEmpty:
                    await asyncio.sleep(5)
                    continue

                phone = acc["phone"]
                if phone in flood_map:
                    if time.time() < flood_map[phone]:
                        await account_queue.put(acc)
                        await asyncio.sleep(5)
                        continue
                    else:
                        del flood_map[phone]
                await asyncio.sleep(random.uniform(0.5, 2.0))
                try:
                    c = await get_client(phone)

                except Exception as e:
                    log.error(f"[WORKER-{worker_id}] {phone} فشل الاتصال: {e}")
                    if is_ban_error(e) or phone in dead_sessions:
                        # الحساب محظور رسمياً أو مُحظر تلقائياً بسبب جلسة تالفة/معلّقة
                        # (get_client يتكفّل بالحفظ في القاعدة) — لا داعي لإعادة المحاولة.
                        fail_count.pop(phone, None)
                        continue
                    # تصعيد عدّاد فشل هذا الحساب وحساب التأخير
                    fail_count[phone] = fail_count.get(phone, 0) + 1
                    delay = _account_backoff(phone)
                    # كشف الانقطاع الشبكي العام: عند تكرار أخطاء شبكية عبر الحسابات
                    if _is_network_error(e):
                        tnow = time.time()
                        net_fail_window.append(tnow)
                        while net_fail_window and tnow - net_fail_window[0] > NET_FAIL_WINDOW:
                            net_fail_window.pop(0)
                        if len(net_fail_window) >= NET_FAIL_THRESHOLD:
                            _net_pause_until[0] = tnow + SYSTEM_PAUSE
                            net_fail_window.clear()
                            log.warning(
                                f"[NET] رُصد انقطاع شبكي عام — هدنة {SYSTEM_PAUSE}s"
                            )
                    await account_queue.put(acc)
                    await asyncio.sleep(delay)
                    continue

                async with get_account_lock(phone):
                    result = await run_account_once(c, phone, bot_un)
                # نجح الاتصال → تصفير عدّاد الفشل لهذا الحساب
                fail_count.pop(phone, None)
                log.info(f"[WORKER-{worker_id}] {phone} → {result}")

                if result == "banned":
                    continue
                if result == "stopped":
                    await account_queue.put(acc)
                    break
                if result == "done":
                    await notify(
                        f"✅ اكتملت جميع القنوات بنجاح\n"
                        f"📱 الحساب: <code>{phone}</code>\n"
                        f"🤖 البوت: @{bot_un}\n"
                        f"⏳ وقت الاستراحة: 5 دقائق",
                        channel=True
                    )
                    # المحور 1: إغلاق العميل قبل فترة الراحة الطويلة (5 دقائق).
                    # إبقاؤه مفتوحاً أثناء الخمول يجعل جلسته تتعفّن وتتصادم
                    # مع restart() الداخلي عند الاضطراب الشبكي. نفتح اتصالاً
                    # نظيفاً بعد الراحة بدل إعادة استخدام المتعفّن.
                    await cleanup_client(phone)

                    async def requeue_after_rest(a=acc):
                        await safe_sleep(300)
                        if not stop_event.is_set():
                            await account_queue.put(a)
                    create_tracked_task(requeue_after_rest())
                elif result == "ok":
                    log.info(f"[GATHER] {phone} | نتيجة ok - راحة 5 دقائق")
                    async def requeue_ok(a=acc):
                        await safe_sleep(300)
                        if not stop_event.is_set():
                            await account_queue.put(a)
                    create_tracked_task(requeue_ok())
                else:
                    await account_queue.put(acc)

            except asyncio.CancelledError:
                raise
            except Exception as e:
                # المحور 2: أي استثناء غير متوقّع لا يقتل الـ worker.
                # نسجّله، نُغلق العميل المتضرّر إن أمكن، ونكمل الحلقة.
                log.error(f"[WORKER-{worker_id}] استثناء غير متوقّع: "
                          f"{type(e).__name__}: {e}")
                try:
                    if 'phone' in dir() and phone:
                        await cleanup_client(phone)
                except Exception:
                    pass
                await asyncio.sleep(2)
                continue
    workers = [worker(i) for i in range(concurrent)]
    try:
        await asyncio.gather(*workers, return_exceptions=True)
    finally:
        # المحور 3: مهما كان سبب الخروج (نهاية طبيعية، استثناء، إلغاء)،
        # نضمن إيقاف المهام الدورية، إغلاق كل العملاء، وإعادة ضبط الحالة
        # حتى لا يبقى is_active=1 عالقاً ويمنع إعادة التشغيل.
        for _t in (_bg_cleanup, _bg_health):
            _t.cancel()
            try:
                await _t
            except (asyncio.CancelledError, Exception):
                pass
        for _phone in list(clients.keys()):
            try:
                await cleanup_client(_phone)
            except Exception:
                pass
        try:
            await notify(
                f"🛑 توقف التجميع\n🤖 البوت: @{bot_un}",
                urgent=True
            )
        except Exception:
            pass
        await set_setting("is_active", "0")
        await set_setting("active_bot", "")
# ============================================================
# نظام 2: كل البوتات — Round-Robin + وقت ثابت
# ============================================================
async def run_gather_all_bots():
    await start_notify_worker()
    try:
        bots_list = await get_bots()
        accounts = await get_accounts(active_only=True)

        if not bots_list or not accounts:
            await notify("⚠️ لا توجد بوتات أو حسابات كافية!", urgent=True)
            await set_setting("is_active", "0")
            await set_setting("active_bot", "")
            return

        invalid_bots = []
        for bot_item in bots_list:
            bot_target = (
                bot_item.get("url")
                or f"https://t.me/{bot_item['username']}"
            )
            if not await verify_bot_exists(bot_target):
                invalid_bots.append(f"@{bot_item['username']}")
        if invalid_bots:
            await notify(
                f"⚠️ البوتات التالية غير موجودة:\n"
                f"{', '.join(invalid_bots)}\n"
                f"يرجى مراجعة البوتات المضافة.",
                urgent=True
            )
            await set_setting("is_active", "0")
            await set_setting("active_bot", "")
            return

        n_bots = len(bots_list)
        n_accounts = len(accounts)
        concurrent = await get_concurrent_all_bots()

        base = n_accounts // n_bots
        remainder = n_accounts % n_bots
        batches = []
        idx = 0
        for i in range(n_bots):
            size = base + (1 if i < remainder else 0)
            batches.append(accounts[idx:idx + size])
            idx += size

        bots_str = ", ".join([f"@{b['username']}" for b in bots_list])
        await notify(
            f"🚀 بدأ التجميع على جميع البوتات\n"
            f"🤖 البوتات: {bots_str}\n"
            f"👥 الحسابات: {n_accounts}\n"
            f"📦 التوزيع: {[len(b) for b in batches]}",
            urgent=True
        )

        batch_queues = []
        for batch in batches:
            q = asyncio.Queue()
            for acc in batch:
                await q.put(acc)
            batch_queues.append(q)
        opened, failed = 0, 0
        total = len(accounts)
        msg_ids = await notify(f"⏳ عزيزي المطور جار فتح {total} جلسة...\n✅ تم فتح حتى الان: 0 جلسة.", urgent=True)
        for acc in accounts:
            if stop_event.is_set():
                break
            phone = acc["phone"]
            try:
                await get_client(phone)
                opened += 1
                log.info(f"[WARMUP] ✓ {phone} ({opened}/{total})")
            except Exception as e:
                failed += 1
                log.warning(f"[WARMUP] ✗ {phone}: {e}")
            await notify_edit(msg_ids or [], f"⏳ عزيزي المطور جارٍ فتح {total} جلسة...\n✅ تم فتح حتى الان: {opened} جلسة.")
            await asyncio.sleep(1)
        await notify(
            f"✅ اكتمل فتح الجلسات\n"
            f"✓ مفتوحة: {opened} | ✗ فاشلة: {failed}\n"
            f"🚀 بدء التجميع الآن...",
            urgent=True
        )

        stop_round_event = asyncio.Event()
        round_num = 0
        rounds_since_rest = 0  # ← عداد الجولات منذ آخر استراحة
        REST_EVERY = 25        # ← استراحة كل 25 جولة
        REST_DURATION = 600    # ← 10 دقائق

        while not stop_event.is_set():
            work_minutes = await get_work_minutes()
            work_secs = work_minutes * 60
            ACCOUNT_TIMEOUT = 300
            GRACE_PERIOD = ACCOUNT_TIMEOUT + 60

            stop_round_event.clear()
            log.info(
                f"[ALL BOTS] جولة {round_num + 1} | "
                f"وقت العمل: {work_minutes} دقيقة"
            )

            async def run_single_bot(bot_item, bot_queue):
                bot_target = (
                    bot_item.get("url")
                    or f"https://t.me/{bot_item['username']}"
                )
                bot_display = bot_item["username"]

                log.info(
                    f"[ALL BOTS] @{bot_display} | "
                    f"جولة {round_num + 1} | "
                    f"{bot_queue.qsize()} حساب في القائمة"
                )

                async def bot_worker(wid: int):
                    log.info(f"[ALL BOTS W-{wid}@{bot_display}] ▶ عامل نشط")
                    while not stop_event.is_set() and not stop_round_event.is_set():
                        try:
                            acc = bot_queue.get_nowait()
                        except asyncio.QueueEmpty:
                            await asyncio.sleep(3)
                            continue

                        phone = acc["phone"]
                        log.debug(f"[ALL BOTS W-{wid}@{bot_display}] {phone} ▶ استخراج من القائمة")
                        if phone in flood_map:
                            if time.time() < flood_map[phone]:
                                await bot_queue.put(acc)
                                await asyncio.sleep(5)
                                continue
                            else:
                                del flood_map[phone]
                        try:
                            c = await get_client(phone)
                        except Exception as e:
                            log.error(f"[ALL BOTS W-{wid}] {phone}: {e}")
                            if not is_ban_error(e):
                                await bot_queue.put(acc)
                            await asyncio.sleep(5)
                            continue

                        try:
                            async with get_account_lock(phone):
                                result = await asyncio.wait_for(
                                    run_account_once(c, phone, bot_target),
                                    timeout=ACCOUNT_TIMEOUT
                                )
                        except asyncio.TimeoutError:
                            log.warning(
                                f"[ALL BOTS W-{wid}] {phone} | "
                                f"timeout {ACCOUNT_TIMEOUT//60} دقيقة - تخطي"
                            )
                            await bot_queue.put(acc)
                            continue
                        except Exception as e:
                            log.info(f"[ALL BOTS W-{wid}@{bot_display}] {phone}: {e}")
                            await bot_queue.put(acc)
                            continue

                        log.info(f"[ALL BOTS W-{wid}@{bot_display}] {phone} → {result}")

                        if result == "banned":
                            continue
                        if result == "stopped":
                            await bot_queue.put(acc)
                            break
                        if result == "done":
                            await notify(
                                f"✅ اكتملت جميع القنوات بنجاح\n"
                                f"📱 الحساب: <code>{phone}</code>\n"
                                f"🤖 البوت: @{bot_display}\n"
                                f"⏳ وقت الاستراحة: 5 دقائق",
                                channel=True
                            )
                            async def requeue(a=acc, q=bot_queue):
                                await safe_sleep(300)
                                if not stop_event.is_set():
                                    await q.put(a)
                            create_tracked_task(requeue())
                        elif result == "ok":
                            # "ok" تعني عدم وجود قنوات بعد - راحة 5 دقائق
                            # بدل إعادة فورية (وقف ضغط /start المتكرر)
                            log.info(
                                f"[GATHER] {phone} | نتيجة ok - راحة 5 دقائق"
                            )
                            async def requeue_ok(a=acc, q=bot_queue):
                                await safe_sleep(300)
                                if not stop_event.is_set():
                                    await q.put(a)
                            create_tracked_task(requeue_ok())
                        else:
                            await bot_queue.put(acc)

                bot_workers_count = max(1, concurrent // n_bots)
                bw = [bot_worker(i) for i in range(bot_workers_count)]
                await asyncio.gather(*bw, return_exceptions=True)

            bot_tasks = []
            for i, bot_item in enumerate(bots_list):
                if stop_event.is_set():
                    break
                batch_idx = (round_num + i) % n_bots
                q = batch_queues[batch_idx]
                bot_tasks.append(
                    asyncio.create_task(run_single_bot(bot_item, q))
                )

            async def round_timer():
                await safe_sleep(work_secs)
                if not stop_event.is_set():
                    log.info(
                        f"[ALL BOTS] انتهى وقت الجولة {round_num + 1} "
                        f"- إشارة للـ Workers بالتوقف"
                    )
                    stop_round_event.set()

            timer_task = asyncio.create_task(round_timer())

            try:
                await asyncio.wait_for(
                    asyncio.gather(*bot_tasks, return_exceptions=True),
                    timeout=work_secs + GRACE_PERIOD
                )
            except asyncio.TimeoutError:
                log.warning(
                    f"[ALL BOTS] انتهت grace period الجولة {round_num + 1} "
                    f"- إلغاء Workers المتبقية"
                )
                stop_round_event.set()
                for t in bot_tasks:
                    t.cancel()
                await asyncio.gather(*bot_tasks, return_exceptions=True)

            timer_task.cancel()
            try:
                await timer_task
            except asyncio.CancelledError:
                pass

            for phone in list(account_join_count.keys()):
                account_join_count[phone] = 0

            if stop_event.is_set():
                break

            rounds_since_rest += 1

            if rounds_since_rest >= REST_EVERY:
                rounds_since_rest = 0
                rest_mins = REST_DURATION // 60
                await notify(
                    f"✅ انتهت الجولة {round_num + 1}\n"
                    f"😴 استراحة عامة: {rest_mins} دقيقة بعد {REST_EVERY} جولة",
                    urgent=True
                )
                log.info(f"[ALL BOTS] استراحة عامة {REST_DURATION} ثانية")
                await safe_sleep(REST_DURATION)
            else:
                await notify(
                    f"✅ انتهت الجولة {round_num + 1}\n"
                    f"🤖 الحسابات تعمل على الجولة {round_num + 2} الآن مباشرةً",
                    urgent=True
                )
                log.info(
                    f"[ALL BOTS] بدء الجولة {round_num + 2} مباشرةً"
                )

            if stop_event.is_set():
                break
            round_num += 1

        await notify("🛑 توقف التجميع على جميع البوتات", urgent=True)

    except asyncio.CancelledError:
        log.info("[ALL BOTS] cancelled.")
    except Exception as e:
        log.error(f"[ALL BOTS] {e}")
    finally:
        await set_setting("is_active", "0")
        await set_setting("active_bot", "")

# ============================================================
# دوال المغادرة
# ============================================================
async def do_leave_ch(ch: dict, uid: int):
    accounts = await get_accounts(active_only=True)
    link = ch_link(ch)
    title = ch.get("title") or link or "القناة"
    for acc in accounts:
        try:
            c = await get_client(acc["phone"])
            has_left = await leave_ch(c, acc["phone"], ch)
            if has_left:
                # ← مغادرة → القناة
                await notify(
                    f"📤 مغادرة قناة\n"
                    f"👤 الحساب: <code>{acc['phone']}</code>\n"
                    f"🔗 القناة: {link or title}",
                    channel=True
                )
                await asyncio.sleep(rand_delay(LEAVE_SPEED))
        except Exception as e:
            log.error(f"[LEAVE] {acc['phone']}: {e}")
    await del_channel(ch_id=ch.get("id"), url=link)
    ch_html = f'<a href="{link}">{title}</a>' if link else title
    try:
        await _bot_instance.send_message(
            uid,
            f"✅ تم مغادرة القناة {ch_html} من جميع الحسابات بنجاح.",
            parse_mode="HTML",
            disable_web_page_preview=True
        )
    except Exception as e:
        log.error(f"[LEAVE] notify: {e}")


async def do_leave_all(uid: int):
    global leave_all_task
    BATCH_SIZE = 100
    await start_notify_worker()
    try:
        log.info("[LEAVE ALL] تسخين الحسابات (متدرّج)...")
        accounts = await get_accounts(active_only=True)
        to_warm = [
            acc["phone"] for acc in accounts
            if acc["phone"] not in clients and acc["phone"] not in dead_sessions
        ]
        # تسخين متدرّج: فتح الحسابات على دفعات محكومة بدل دفعة واحدة.
        # الفتح الجماعي (gather على 296 معاً) يثير قطع تيليجرام للاتصالات
        # ويبدأ سلسلة Connection lost → read() called. التدرّج يوزّع الضغط.
        WARM_CONC = 15          # عدد الحسابات المفتوحة معاً في آنٍ واحد
        WARM_GAP = 0.3          # فاصل بسيط بين بدء كل عميل
        _warm_sem = asyncio.Semaphore(WARM_CONC)
        warm_opened = [0]

        async def _warm_one(phone):
            async with _warm_sem:
                try:
                    await asyncio.wait_for(
                        get_client(phone, updates=False), timeout=30)
                    warm_opened[0] += 1
                except Exception:
                    pass
                await asyncio.sleep(WARM_GAP)

        if to_warm:
            await asyncio.gather(
                *(_warm_one(p) for p in to_warm),
                return_exceptions=True)
            log.info(f"[LEAVE ALL] تم فتح {warm_opened[0]}/{len(to_warm)} حساب")
        # السبب الجذري للمشكلة السابقة: هذه المقارنة كانت نصية خامة
        # (ch['url'] المخزَّن بصيغة طبيعية مقابل رابط الاشتراك الإجباري
        # الخام كما ورد من الرسالة) فكانت تفشل فعلياً في كل مرة تقريباً
        # بسبب اختلاف الصيغة (حالة الأحرف / بادئة الرابط)، فلا تُستثنى
        # قنوات الاشتراك الإجباري أبداً. نستخدم الآن norm_ch_key الموحّدة
        # (نفس الدالة المستخدمة أصلاً في بقية الكود لمقارنة الروابط).
        _force_sub_keys, _force_sub_ids = await get_all_force_sub_keys()
        channels = sorted(
            [ch for ch in await get_channels()
             if norm_ch_key(ch.get('url') or '') not in _force_sub_keys
             and (ch.get('id') is None or int(ch['id']) not in _force_sub_ids)],
            key=lambda x: x.get("joined_at") or "",
        )
        accounts = await get_accounts(active_only=True)
        # ── إعدادات إعادة المحاولة ──
        RETRY_GAP    = 5           # ثوانٍ بين كل محاولة في مراجعة الفاشلين
        REVIEW_AT    = 200         # نراجع الفاشلين عندما يتبقى هذا العدد من القنوات
        LEAVE_TIMEOUT = 25         # مهلة مغادرة الحساب الواحد (تمنع تجمّد الدفعة)
        _review_done = False       # نراجع مرة واحدة فقط عند 200 قناة متبقية

        # ════════════════════════════════════════════════════════
        # بدء المغادرة من القنوات المسجّلة في القاعدة
        # ════════════════════════════════════════════════════════
        await notify(
            "🔄 بدء مغادرة القنوات المسجّلة...\n"
            f"عدد القنوات: {len(channels)}",
            urgent=True)

        # ════════════════════════════════════════════════════════
        # المغادرة من القنوات المسجّلة (الأقدم أولاً)
        # ════════════════════════════════════════════════════════
        # جدول الفشل المتراكم وعدّاد القنوات
        pending_fail = []
        channels_done = 0

        async def _leave_one(p, ch_obj):
            """يحاول إخراج حساب واحد من قناة واحدة.
            يعيد True إن غادر (أو لم يكن مشتركاً)، False إن فشل فشلاً حقيقياً."""
            try:
                c = await asyncio.wait_for(
                    get_client(p, updates=False), timeout=20)
                await asyncio.wait_for(
                    leave_ch(c, p, ch_obj, ctx="BATCH LEAVE"),
                    timeout=LEAVE_TIMEOUT)
                return True
            except Exception as e:
                err = str(e).lower()
                if "flood_wait" in err or "420" in err:
                    sec = re.findall(r'\d+', err)
                    wait_sec = int(sec[0]) if sec else 3600
                    flood_map[p] = time.time() + wait_sec + 60
                    log.warning(f"[BATCH LEAVE] {p}: FloodWait {wait_sec}s")
                elif ("malformed" in err or "disk image" in err
                        or "auth_key" in err or "unregistered" in err
                        or "revoked" in err or "eof when reading" in err):
                    dead_sessions.add(p)
                    log.warning(f"[BATCH LEAVE] {p}: جلسة تالفة/ميتة — تُتخطّى")
                else:
                    log.error(f"[BATCH ERROR] {p}: {e}")
                return False

        async def _leave_channel_pass(ch_obj, phones, report=False):
            """مرور واحد لإخراج قائمة حسابات من قناة، على دفعات متوازية.
            يعيد قائمة الأرقام التي فشلت فشلاً حقيقياً.
            report=True يرسل تقرير تيليجرام بعد كل دفعة (للمرور الأول فقط)."""
            failed = []
            link_ = ch_link(ch_obj)
            for i in range(0, len(phones), BATCH_SIZE):
                if not _shared.leave_all_on:
                    break
                batch = phones[i:i + BATCH_SIZE]
                runnable = []
                for phone in batch:
                    if phone in dead_sessions:
                        continue
                    if phone in flood_map:
                        if time.time() < flood_map[phone]:
                            failed.append(phone)
                            continue
                        else:
                            del flood_map[phone]
                    runnable.append(phone)

                if not runnable:
                    continue

                results = await asyncio.gather(
                    *(_leave_one(p, ch_obj) for p in runnable),
                    return_exceptions=True)

                ok = 0
                for phone, res in zip(runnable, results):
                    if res is True:
                        ok += 1
                    else:
                        failed.append(phone)

                if report and ok > 0:
                    await notify(
                        f"📤 <b>تقرير مغادرة دفعة</b>\n"
                        f"🔗 <b>القناة:</b> {link_}\n"
                        f"✅ <b>عدد حسابات المغادرة:</b> {ok} حساب",
                        channel=True)

                await asyncio.sleep(0.1)
            return failed

        for ch in list(channels):
            if not _shared.leave_all_on:
                break
            ch_id, link = ch.get("id"), ch_link(ch)
            all_phones = [a["phone"] for a in accounts]

            log.info(f"[LEAVE ALL] ━━ بدء مغادرة القناة: {link} "
                     f"({len(all_phones)} حساب) ━━")

            # المرور الأول (مع تقرير لكل دفعة 50)
            failed = await _leave_channel_pass(ch, all_phones, report=True)
            log.info(f"[LEAVE ALL] قناة {link}: المرور الأول انتهى, "
                     f"فشل {len(failed)} حساب")

            # الفاشلون يدخلون جدول الانتظار للمراجعة عند 200 قناة متبقية
            for p in failed:
                pending_fail.append((p, ch))
            if failed:
                log.warning(f"[LEAVE ALL] قناة {link}: {len(failed)} حساب "
                            f"أُضيف لجدول الفشل")

            await del_channel(ch_id=ch_id, url=link)
            channels_done += 1
            await asyncio.sleep(0.2)

            # ── مراجعة الفاشلين عندما يتبقى 200 قناة ──
            remaining = len(channels) - channels_done
            if (not _review_done and remaining <= REVIEW_AT
                    and pending_fail and _shared.leave_all_on):
                _review_done = True
                log.info(f"[REVIEW] تبقّى {remaining} قناة — بدء مراجعة "
                         f"{len(pending_fail)} حساب فاشل")
                await notify(
                    f"🔄 <b>مراجعة الفاشلين</b>\n"
                    f"تبقّى {remaining} قناة — جاري إعادة المحاولة على "
                    f"{len(pending_fail)} حساب فاشل...",
                    urgent=True)
                still = []
                from collections import defaultdict
                by_ch = defaultdict(list)
                for p, chx in pending_fail:
                    by_ch[id(chx)].append((p, chx))
                for _gid, items in by_ch.items():
                    if not _shared.leave_all_on:
                        break
                    chx = items[0][1]
                    phones_x = [p for p, _ in items]
                    log.info(f"[REVIEW] إعادة على قناة {ch_link(chx)}: "
                             f"{len(phones_x)} حساب")
                    await asyncio.sleep(RETRY_GAP)
                    fx = await _leave_channel_pass(chx, phones_x)
                    for p in fx:
                        still.append((p, chx))
                pending_fail = still
                log.info(f"[REVIEW] انتهت المراجعة — بقي {len(pending_fail)} فاشل")

        # ── مراجعة نهائية لما تبقّى في جدول الفشل ──
        if pending_fail and _shared.leave_all_on:
            log.info(f"[FINAL REVIEW] {len(pending_fail)} إدخال فاشل متبقٍ — "
                     f"محاولة أخيرة")
            from collections import defaultdict
            by_ch = defaultdict(list)
            for p, chx in pending_fail:
                by_ch[id(chx)].append((p, chx))
            still = []
            for _gid, items in by_ch.items():
                if not _shared.leave_all_on:
                    break
                chx = items[0][1]
                phones_x = [p for p, _ in items]
                await asyncio.sleep(RETRY_GAP)
                fx = await _leave_channel_pass(chx, phones_x)
                for p in fx:
                    still.append((p, chx))
            pending_fail = still

        # تقرير نهائي بما لم يغادر إطلاقاً
        if pending_fail:
            lines = "\n".join(
                f"  • <code>{p}</code> من {ch_link(chx)}"
                for p, chx in pending_fail[:50])
            extra = "" if len(pending_fail) <= 50 else f"\n... و{len(pending_fail)-50} غيرها"
            await notify(
                f"⚠️ <b>حسابات لم تغادر رغم كل المحاولات</b>\n"
                f"العدد: {len(pending_fail)}\n{lines}{extra}",
                urgent=True)
            log.warning(f"[LEAVE ALL] انتهى — {len(pending_fail)} لم تغادر")
        else:
            log.info("[LEAVE ALL] انتهى بنجاح — كل الحسابات غادرت كل القنوات")

    except asyncio.CancelledError:
        log.info("[LEAVE ALL] أُلغيت العملية")
    except Exception as e:
        log.exception(f"[LEAVE ALL] خطأ غير متوقع: {e}")
    finally:
        _shared.leave_all_on = False
        _shared.leave_all_task = None


async def do_leave_force_sub(uid: int, bot_un: str, urls: list):
    # ── 1. التحقق من عدم وجود عملية جارية ──
    is_active = await get_setting("is_active", "0")
    if is_active == "1":
        await notify("⚠️ لا يمكن بدء المغادرة — التجميع يعمل الآن.", urgent=True)
        return
    if _shared.leave_all_on:
        await notify("⚠️ لا يمكن بدء المغادرة — مغادرة كل القنوات تعمل الآن.", urgent=True)
        return

    accounts = await get_accounts(active_only=True)
    total_urls = len(urls)
    total_acc = len(accounts)

    # ── 2. تسخين الجلسات غير المفتوحة ──
    to_warm = [
        acc["phone"] for acc in accounts
        if acc["phone"] not in clients and acc["phone"] not in dead_sessions
    ]
    if to_warm:
        await notify(
            f"⏳ جاري فتح الجلسات قبل المغادرة...\n"
            f"📊 يحتاج فتح: {len(to_warm)} جلسة",
            urgent=True
        )
        WARM_CONC = 15
        WARM_GAP  = 0.3
        _warm_sem = asyncio.Semaphore(WARM_CONC)
        warm_opened = [0]

        async def _warm_one_fs(phone):
            async with _warm_sem:
                try:
                    await asyncio.wait_for(
                        get_client(phone, updates=False), timeout=30)
                    warm_opened[0] += 1
                except Exception:
                    pass
                await asyncio.sleep(WARM_GAP)

        await asyncio.gather(
            *(_warm_one_fs(p) for p in to_warm),
            return_exceptions=True)
        log.info(f"[LEAVE FORCE SUB] تم فتح {warm_opened[0]}/{len(to_warm)} جلسة")
        await asyncio.sleep(3)

    # ── 3. جلب rows من force_sub_channels للحصول على channel_id ──
    from shared import get_force_sub_channels
    fs_rows = await get_force_sub_channels(bot_un)
    # بناء خريطة url → ch_obj
    ch_map = {}
    for row in fs_rows:
        url = row["url"]
        ch_id = row.get("channel_id")
        un = url.rstrip("/").split("/")[-1]
        ch_map[url] = {
            "id": ch_id,
            "username": None if un.startswith("+") else un,
            "url": url,
        }
    # أي url غير موجود في القاعدة نضيفه بدون id
    for url in urls:
        if url not in ch_map:
            un = url.rstrip("/").split("/")[-1]
            ch_map[url] = {"id": None, "username": None if un.startswith("+") else un, "url": url}

    # ── 4. إشعار البدء ──
    opened_count = sum(1 for acc in accounts if acc["phone"] in clients)
    action = "مغادرة قناة" if total_urls == 1 else f"مغادرة {total_urls} قناة"
    await notify(
        f"🚀 <b>بدء مغادرة الاشتراك الإجباري</b>\n"
        f"🤖 البوت: @{bot_un}\n"
        f"📋 العملية: {action}\n"
        f"👥 الحسابات المفتوحة: {opened_count}/{total_acc}",
        urgent=True
    )

    # ── 5. تنفيذ المغادرة بنفس منطق do_leave_all ──
    done = 0
    failed = 0
    log.info(f"[LEAVE FORCE SUB] @{bot_un}: {total_urls} قناة × {total_acc} حساب")
    LEAVE_CONC = 10
    _leave_sem = asyncio.Semaphore(LEAVE_CONC)

    async def _leave_one_fs(phone, ch_obj, url):
        nonlocal done, failed
        async with _leave_sem:
            try:
                async with get_account_lock(phone):
                    c = await get_client(phone)
                    await asyncio.wait_for(
                        leave_ch(c, phone, ch_obj, ctx="FORCE SUB LEAVE"),
                        timeout=25)
                done += 1
            except Exception as e:
                err = str(e).lower()
                if any(k in err for k in ("not found", "invalid", "user not participant",
                                          "channel_private", "peer_id_invalid",
                                          "usernotparticipant", "chat_not_found")):
                    done += 1  # غير مشترك = اعتبره ناجحاً
                else:
                    failed += 1
                    log.warning(f"[LEAVE FORCE SUB] {phone} | {url}: {e}")

    for url in urls:
        ch_obj = ch_map[url]
        log.info(f"[LEAVE FORCE SUB] مغادرة {url} | id={ch_obj['id']}")
        await asyncio.gather(
            *(_leave_one_fs(acc["phone"], ch_obj, url) for acc in accounts),
            return_exceptions=True
        )
        await asyncio.sleep(0.5)

    # ── 6. حذف القنوات من قاعدة البيانات بعد المغادرة ──
    await del_force_sub_channels(bot_un)
    log.info(f"[LEAVE FORCE SUB] تم حذف قنوات @{bot_un} من القاعدة")

    # ── 7. إشعار الاكتمال ──
    await notify(
        f"✅ <b>اكتملت مغادرة الاشتراك الإجباري</b>\n"
        f"🤖 البوت: @{bot_un}\n"
        f"📋 القنوات: {total_urls}\n"
        f"👥 الحسابات: {total_acc}\n"
        f"✅ نجح: {done} | ✗ فشل: {failed}",
        urgent=True
    )
    log.info(f"[LEAVE FORCE SUB] انتهى — نجح: {done} فشل: {failed}")



