import asyncio
import base64
import functools
import hashlib
import json
import logging
import os
import secrets
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional

import psutil
from telethon import TelegramClient, events, errors, functions, types
from telethon.errors import (
    AuthKeyUnregisteredError,
    AuthTokenExpiredError,
    AuthTokenInvalidError,
    PasswordHashInvalidError,
    PhoneCodeExpiredError,
    PhoneCodeInvalidError,
    SessionPasswordNeededError,
)
from telethon.errors.rpcerrorlist import FloodWaitError, SendCodeUnavailableError
from telethon.sessions import StringSession
from telethon.network import ConnectionTcpAbridged
from telethon.tl.custom.qrlogin import QRLogin

import config
from ctrader_spot_client import CTraderSpotClient
from forwarding import MessageForwarder
from llm_classifier import LlmEngine
from signal_processor import SignalProcessor

logger = logging.getLogger(__name__)


class ForwarderCore:
    """Multi-account Telegram forwarder + async client manager."""

    def __init__(self, db):
        self.db = db
        self.clients: Dict[int, TelegramClient] = {}
        self.loop = asyncio.new_event_loop()
        self.thread = threading.Thread(target=self._run_loop, daemon=True)
        self._executor = ThreadPoolExecutor(max_workers=8, thread_name_prefix="db-")
        # Per-account routing: _pairs[account_id][source_chat_id] = [pair, ...].
        # Pairs with account_id=NULL ("auto") live under the None key and are
        # handled by whichever connected client sees the source first.
        self._pairs: Dict[Optional[int], Dict[int, List[dict]]] = {}
        # Auth/QR state is scoped per account so one account's login flow can
        # never clobber or observe another's.
        self._auth_state: Dict[int, dict] = {}
        self._qr_state: Dict[int, dict] = {}
        self._forwarding = False
        self._handlers_accounts: set = set()
        self._memory_task: Optional[asyncio.Task] = None
        self._pairs_sync_task: Optional[asyncio.Task] = None
        self._watchdog_task: Optional[asyncio.Task] = None
        self._ctrader_client: Optional[CTraderSpotClient] = None
        self._ctrader_started = False
        self._signal_processor: Optional[SignalProcessor] = None
        self._llm_engine: LlmEngine = LlmEngine(db)
        self._forwarder: Optional[MessageForwarder] = None
        self._retrying: bool = False
        self._init_started: bool = False
        # Cross-process core lock (MySQL GET_LOCK). Only the worker holding it
        # runs Telegram clients; the others serve HTTP and read the shared
        # dialog cache. Held connection lives for the app's lifetime.
        self._lock_conn = None
        self._primary_account_id: Optional[int] = None

    def _run_loop(self):
        try:
            asyncio.set_event_loop(self.loop)
            logger.info("Event loop thread started, running loop forever")
            self.loop.run_forever()
        except Exception as exc:
            logger.error("Event loop thread CRASHED: %s", exc, exc_info=True)
        finally:
            logger.warning("Event loop thread exiting")

    def start(self):
        if not self.thread.is_alive():
            self.thread.start()
            time.sleep(0.3)
        logger.info("Scheduling _init_accounts on event loop (alive=%s)", self.thread.is_alive())
        asyncio.run_coroutine_threadsafe(self._init_accounts(), self.loop)

    async def _run_db(self, func, *args, **kwargs):
        """Run a blocking DB call in the thread pool from the event loop."""
        if kwargs:
            func = functools.partial(func, **kwargs)
        return await self.loop.run_in_executor(self._executor, func, *args)

    def _run_db_sync(self, func, *args, timeout=10, **kwargs):
        """Run a blocking DB call in the thread pool from a synchronous caller."""
        if kwargs:
            func = functools.partial(func, **kwargs)
        future = self._executor.submit(func, *args)
        return future.result(timeout=timeout)

    async def _init_accounts(self):
        # Prevent duplicate concurrent runs — Telegram rejects a second
        # connection from the same IP for the same account.
        if self._init_started:
            logger.info("_init_accounts already running, skipping duplicate")
            return
        self._init_started = True
        logger.info("_init_accounts coroutine started")
        try:
            await self._init_accounts_inner()
        except Exception as exc:
            logger.error("_init_accounts failed: %s", exc, exc_info=True)
        finally:
            self._init_started = False
            self._retrying = False

    async def _init_accounts_inner(self):
        logger.info("_init_accounts_inner: starting DB migrations")

        # Load account list first so passive workers also know the primary
        # account — this lets them serve the correct shared dialog cache.
        accounts = await self._run_db(self.db.get_accounts)
        authorized = [a for a in accounts if a.get("is_authorized")]
        self._primary_account_id = next(
            (a["id"] for a in authorized if a.get("is_primary")),
            authorized[0]["id"] if authorized else None,
        )

        # Only ONE Passenger worker may run the Telegram core: a second
        # concurrent connection with the same session auth key makes the DC
        # drop the first worker's pending requests (Telethon then swallows the
        # cancellation, retries internally, and wedges the event loop for
        # minutes). MySQL GET_LOCK is the cross-process mutex — losers serve
        # HTTP only. If the leader dies, MySQL releases the lock and a
        # follower's next retry (ensure_core) promotes it.
        if self._lock_conn is None:
            self._lock_conn = await self._run_db(self.db.acquire_core_lock)
        if not self._lock_conn:
            logger.info(
                "Another worker holds the core lock — running passive (HTTP only)"
            )
            return

        # One-time migration: move legacy session to tg_accounts table
        await self._run_db(self.db.migrate_legacy_session)
        logger.info("_init_accounts_inner: migrate_legacy_session done")

        # Idempotently add LLM + pair columns to existing databases, then load config.
        logger.info("_init_accounts_inner: running ensure_llm_schema")
        try:
            await self._run_db(self.db.ensure_llm_schema)
        except Exception as exc:
            logger.warning("ensure_llm_schema failed: %s", exc)
        logger.info("_init_accounts_inner: running ensure_pair_schema")
        try:
            await self._run_db(self.db.ensure_pair_schema)
        except Exception as exc:
            logger.warning("ensure_pair_schema failed: %s", exc)
        logger.info("_init_accounts_inner: reloading LLM config")
        try:
            await self._run_db(self._llm_engine.reload_from_db)
            logger.info("LLM engine: %s", self._llm_engine.status())
        except Exception as exc:
            logger.warning("LLM config load failed: %s", exc)
        logger.info("_init_accounts_inner: schema/LLM steps complete")

        if not authorized:
            logger.info("No authorized accounts found")
            self._memory_task = self.loop.create_task(self._memory_logger())
            self._prune_task = self.loop.create_task(self._prune_loop())
            self._pairs_sync_task = self.loop.create_task(self._pairs_sync_loop())
            self._watchdog_task = self.loop.create_task(self._watchdog_loop())
            return

        for acct in authorized:
            await self._start_account_client(acct["id"])

        await self._load_pairs()
        # Install handlers only on accounts that own at least one pair (or, for
        # auto pairs, any connected client — those are deduped in the forwarder).
        for acct_id in self._accounts_needing_handlers():
            self._install_handlers(acct_id)

        self._memory_task = self.loop.create_task(self._memory_logger())
        self._prune_task = self.loop.create_task(self._prune_loop())
        self._pairs_sync_task = self.loop.create_task(self._pairs_sync_loop())
        self._watchdog_task = self.loop.create_task(self._watchdog_loop())
        logger.info(
            "Initialized %d account(s), %d pair(s)",
            len(self.clients),
            sum(
                len(pair_list)
                for per_src in self._pairs.values()
                for pair_list in per_src.values()
            ),
        )

    async def _bounded(
        self, coro_factory, timeout: float, desc: str, account_id: Optional[int] = None
    ) -> Any:
        """Run a Telethon call with a HARD deadline the loop can never miss.

        ``asyncio.wait_for`` cancels its argument, and Telethon swallows that
        cancellation and retries internally — so the deadline never fires and
        the event loop is wedged for minutes. ``asyncio.wait`` with a timeout
        does NOT cancel the task on timeout; it returns and leaves the task
        detached. That guarantees the deadline fires on time.

        Returns the call's result, or None on timeout. On timeout the orphaned
        task is cancelled best-effort and left to die on its own (bounded: one
        per failed attempt; the retry creates a fresh client).
        """
        result = coro_factory()
        # Telethon's disconnect() returns a Future, not a coroutine.
        task = result if asyncio.isfuture(result) else self.loop.create_task(result)
        done, _pending = await asyncio.wait({task}, timeout=timeout)
        if task in done:
            task.add_done_callback(
                lambda t: t.exception() if not t.cancelled() else None
            )
            return task.result()
        logger.warning(
            "Account %s: %s timed out after %.0fs",
            account_id if account_id is not None else "-",
            desc,
            timeout,
        )
        task.cancel()  # best-effort; may be swallowed, task is detached
        return None

    async def _start_account_client(self, account_id: int):
        session_str = await self._run_db(self.db.get_account_session, account_id)
        if not session_str:
            logger.warning("Account %d has no session string", account_id)
            return

        # Build a list of server candidates: the session's original DC IP first
        # (this is usually the fastest path), then known DC5 alternatives.  The
        # first candidate that connects AND authorizes wins.
        base_session = StringSession(session_str)
        original_addr = getattr(base_session, "_server_address", None)
        candidates = []
        if original_addr:
            candidates.append(original_addr)
        candidates.extend([
            "91.108.56.178", "91.108.56.181", "91.108.56.182",
            "91.108.56.186", "91.108.56.187",
        ])
        # De-duplicate while preserving order.
        seen = set()
        candidates = [ip for ip in candidates if not (ip in seen or seen.add(ip))]

        for ip in candidates:
            session = StringSession(session_str)
            session._server_address = ip
            client = TelegramClient(
                session,
                config.TG_API_ID,
                config.TG_API_HASH,
                connection=ConnectionTcpAbridged,
                # connection_retries=0: we provide our own IP fallback loop.
                # Telethon's internal retry swallows cancellation and can wedge
                # the event loop for minutes on a blocked host.
                connection_retries=0,
                retry_delay=0,
                auto_reconnect=True,
            )
            try:
                logger.info(
                    "Account %d: connecting to Telegram via %s...", account_id, ip
                )
                # Short connect timeout: shared-host Passenger workers are
                # recycled if they do not respond to requests for too long.  We
                # must either succeed or fail fast enough to try the next DC IP
                # within the worker's lifetime.
                # NOTE: Telethon's connect() returns None on success (not True),
                # and _bounded() also returns None on its deadline.  We must NOT
                # distinguish the two by the return value — check is_connected()
                # instead, otherwise a successful connect is misread as a timeout
                # and the client is dropped before it is ever added to self.clients.
                await self._bounded(
                    lambda: client.connect(), 35, "connect", account_id
                )
                if not client.is_connected():
                    await self._bounded(
                        lambda: client.disconnect(), 10, "disconnect", account_id
                    )
                    continue

                authorized = await self._bounded(
                    lambda: client.is_user_authorized(), 15, "is_user_authorized", account_id
                )
                if not authorized:
                    await self._bounded(
                        lambda: client.disconnect(), 10, "disconnect", account_id
                    )
                    continue

                if authorized:
                    self.clients[account_id] = client
                    logger.info(
                        "Account %d connected and authorized via %s", account_id, ip
                    )
                    # Warm the on-disk dialog cache for this account in the
                    # background so passive workers (and this worker if the
                    # connection later drops) can keep serving channels even
                    # before a UI request triggers a dialog fetch.
                    # Warm the on-disk dialog cache for this account in the
                    # background so passive workers (and this worker if the
                    # connection later drops) can keep serving channels even
                    # before a UI request triggers a dialog fetch.
                    try:
                        warm_task = self.loop.create_task(
                            self._async_get_dialogs(refresh=True, account_id=account_id)
                        )
                        warm_task.add_done_callback(
                            lambda t: (
                                logger.warning(
                                    "Account %d: dialog cache warm failed: %s",
                                    account_id, t.exception(),
                                )
                                if not t.cancelled() and t.exception()
                                else None
                            )
                        )
                    except Exception as exc:
                        logger.warning(
                            "Account %d: dialog cache warm failed: %s", account_id, exc
                        )
                    return
                # authorized is False (or None on deadline) — already handled
                # above by `if not authorized: ... continue`.
            except Exception as exc:
                logger.error(
                    "Account %d init failed on %s: %s",
                    account_id, ip, exc, exc_info=True,
                )
            finally:
                # Disconnect any client that was NOT added to self.clients so it
                # doesn't leak and silently retry in the background.
                if account_id not in self.clients:
                    try:
                        await client.disconnect()
                    except Exception:
                        pass
        logger.warning(
            "Account %d: exhausted all Telegram server candidates", account_id
        )

    async def _stop_account_client(self, account_id: int):
        self._remove_handlers(account_id)
        client = self.clients.pop(account_id, None)
        if client:
            try:
                await client.disconnect()
            except Exception as e:
                logger.warning("Disconnect error for account %d: %s", account_id, e)
            logger.info("Account %d disconnected", account_id)

    def _get_client(self, account_id: Optional[int] = None) -> TelegramClient:
        """Get a client by account_id, or any available client. Does NOT block."""
        if account_id:
            if account_id in self.clients:
                return self.clients[account_id]
            raise RuntimeError(f"Account {account_id} is not connected")
        if not self.clients:
            raise RuntimeError("No authorized Telegram clients")
        return next(iter(self.clients.values()))

    def _resolve_account_id(self, account_id: Optional[int]) -> Optional[int]:
        """Resolve a requested account id to a usable account for fetching dialogs.

        - If the requested ``account_id`` is connected, use it.
        - If no clients are connected at all (e.g. a passive worker or while the
          Telegram client is reconnecting), keep the requested id (or the
          primary) so callers can still serve the correct per-account on-disk
          dialog cache.
        - ``None`` means "use the default account": prefer a *connected*
          account, favouring the primary among the connected set. This prevents
          the channel picker from silently mixing channels from arbitrary
          clients when multiple Telegram accounts are connected.
        """
        if account_id is not None:
            # Explicit request — respect it so we never silently show a
            # different account's channels. If that account isn't connected,
            # the caller serves its on-disk cache (or reports "not connected").
            return account_id
        if not self.clients:
            # No live client anywhere — keep the primary so the on-disk dialog
            # cache for the default account is served.
            return self._primary_account_id
        if len(self.clients) == 1:
            return next(iter(self.clients.keys()))
        if self._primary_account_id in self.clients:
            return self._primary_account_id
        return next(iter(self.clients.keys()))

    def _accounts_needing_handlers(self) -> set:
        """Accounts whose clients should have event handlers installed.

        Bound pairs need their own account's client. Auto pairs (account_id
        NULL) may be seen by any connected client, so every client gets a
        handler and the forwarder dedupes by message.
        """
        accts = {a for a in self._pairs if a is not None}
        if None in self._pairs and self._pairs[None]:
            accts |= set(self.clients.keys())
        return accts

    async def _load_pairs(self):
        if not self._lock_conn:
            # Passive worker: keep the routing table empty; the DB is the
            # source of truth and the leader reloads it periodically.
            self._pairs = {}
            return
        rows = await self._run_db(self.db.get_pairs, only_enabled=True)
        pairs: Dict[Optional[int], Dict[int, List[dict]]] = {}
        for row in rows:
            acct = row.get("account_id")
            pairs.setdefault(acct, {}).setdefault(row["source_chat_id"], []).append(row)
        self._pairs = pairs
        logger.info(
            "Loaded %s enabled pair(s) across %s account(s)",
            len(rows),
            len(pairs),
        )
        await self._ensure_ctrader()

    async def _ensure_ctrader(self):
        """Start/stop the cTrader spot client and signal processor as needed."""
        needs_ctrader = any(
            p.get("price_augment")
            for per_src in self._pairs.values()
            for pairs in per_src.values()
            for p in pairs
        )
        needs_bot = any(
            p.get("forward_via_bot")
            for per_src in self._pairs.values()
            for pairs in per_src.values()
            for p in pairs
        )

        if needs_ctrader and not self._ctrader_started:
            if not config.CTRADER_ACCOUNT_ID:
                logger.warning(
                    "Price augmentation requested but CTRADER_ACCOUNT_ID is not set"
                )
            else:
                symbols = self._collect_augment_symbols()
                # Acquire access_token: vault first (refreshable), env var fallback.
                access_token = ""
                if config.CTRADER_TOKEN_SECRET:
                    acct = self._run_db_sync(
                        self.db.get_ctrader_account, config.CTRADER_ACCOUNT_ID
                    )
                    if acct and acct.get("encrypted_payload"):
                        from crypto_vault import decrypt_payload

                        try:
                            payload = decrypt_payload(
                                config.CTRADER_TOKEN_SECRET, acct["encrypted_payload"]
                            )
                            access_token = payload.get("access_token", "")
                        except Exception as exc:
                            logger.warning("Failed to decrypt vault token: %s", exc)

                if not access_token:
                    access_token = config.CTRADER_ACCESS_TOKEN

                if not access_token:
                    logger.warning(
                        "No cTrader access_token available.  "
                        "Set CTRADER_ACCESS_TOKEN or authenticate via /callback."
                    )
                    return

                self._ctrader_client = CTraderSpotClient(
                    client_id=config.CTRADER_CLIENT_ID,
                    client_secret=config.CTRADER_CLIENT_SECRET,
                    access_token=access_token,
                    account_id=config.CTRADER_ACCOUNT_ID,
                    use_live=config.CTRADER_USE_LIVE,
                    symbols=symbols,
                    symbol_ids=[config.CTRADER_SYMBOL_ID]
                    if config.CTRADER_SYMBOL_ID
                    else None,
                )
                await self._ctrader_client.start()
                self._ctrader_started = True
                logger.info(
                    "cTrader spot client started for price augmentation (symbols=%s)",
                    symbols,
                )

        if (needs_ctrader or needs_bot) and not self._signal_processor:
            self._signal_processor = SignalProcessor(
                db=self.db,
                bot_token=config.SIGNAL_BOT_TOKEN,
                admin_chat_id=config.SIGNAL_ADMIN_CHAT_ID or None,
                symbol=config.CTRADER_SYMBOL,
                price_timeout=config.CTRADER_PRICE_TIMEOUT,
                webhook_url=config.SIGNAL_WEBHOOK_URL,
                webhook_secret=config.SIGNAL_WEBHOOK_SECRET,
            )
            logger.info("Signal processor started")

        if not needs_ctrader and self._ctrader_started:
            try:
                await self._ctrader_client.stop()
            except Exception as exc:
                logger.warning("cTrader stop error: %s", exc)
            self._ctrader_client = None
            self._ctrader_started = False
            logger.info("cTrader spot client stopped (no price-augment pairs)")

        if not (needs_ctrader or needs_bot) and self._signal_processor:
            self._signal_processor = None
            logger.info("Signal processor stopped")

    def _collect_augment_symbols(self) -> list[str]:
        """Collect unique symbols from enabled price-augment pairs."""
        symbols = set()
        for per_src in self._pairs.values():
            for pairs in per_src.values():
                for p in pairs:
                    if not p.get("price_augment"):
                        continue
                    sym = (p.get("augment_symbol") or config.CTRADER_SYMBOL).upper()
                    if sym:
                        symbols.add(sym)
        if not symbols:
            symbols.add(config.CTRADER_SYMBOL.upper())
        return sorted(symbols)

    def _install_handlers(self, account_id: int):
        if account_id in self._handlers_accounts:
            return
        client = self.clients.get(account_id)
        if not client:
            return
        if not self._forwarder:
            self._forwarder = MessageForwarder(self)
        self._forwarder.install(client, account_id)
        self._handlers_accounts.add(account_id)
        self._forwarding = True
        logger.info("Handlers installed for account %d", account_id)

    def _remove_handlers(self, account_id: int):
        if account_id not in self._handlers_accounts:
            return
        client = self.clients.get(account_id)
        if client and self._forwarder:
            self._forwarder.remove(client)
        self._handlers_accounts.discard(account_id)
        if not self._handlers_accounts:
            self._forwarding = False
        logger.info("Handlers removed for account %d", account_id)

    async def _memory_logger(self):
        proc = psutil.Process()
        while True:
            try:
                rss = proc.memory_info().rss
                logger.info("Memory RSS: %.2f MB", rss / 1024 / 1024)
            except Exception:
                pass
            try:
                await asyncio.sleep(60)
            except asyncio.CancelledError:
                break

    async def _pairs_sync_loop(self):
        """Reload pairs from the DB periodically.

        Passive workers can create/edit pairs (the DB is the source of truth),
        so the leader's in-memory routing table would go stale; this keeps it
        fresh without restarting.  After reloading, install handlers on any
        newly-connected accounts that now own pairs.
        """
        while True:
            try:
                await asyncio.sleep(30)
                if self._lock_conn:
                    await self._load_pairs()
                    for acct_id in self._accounts_needing_handlers():
                        self._install_handlers(acct_id)
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.error("pairs sync error: %s", e)

    async def _watchdog_loop(self):
        """Keep authorized Telegram accounts connected.

        Runs only on the leader worker (the one holding the core lock).  Every
        30 seconds it compares the DB list of authorized accounts against the
        in-memory clients dict and reconnects any missing or stale client.  This
        recovers from transient DC failures, session drops, and Passenger worker
        restarts without waiting for an HTTP request to trigger ensure_core().
        """
        while True:
            try:
                await asyncio.sleep(30)
                if not self._lock_conn:
                    continue
                accounts = await self._run_db(self.db.get_accounts)
                authorized_ids = {a["id"] for a in accounts if a.get("is_authorized")}
                # Connect authorized accounts that are missing.
                for acct_id in authorized_ids:
                    if acct_id not in self.clients:
                        logger.info("Watchdog: connecting missing account %d", acct_id)
                        await self._start_account_client(acct_id)
                        if acct_id in self._accounts_needing_handlers():
                            self._install_handlers(acct_id)
                # Drop accounts that are no longer authorized.
                for acct_id in list(self.clients.keys()):
                    if acct_id not in authorized_ids:
                        logger.info("Watchdog: disconnecting de-authorized account %d", acct_id)
                        await self._stop_account_client(acct_id)
                # Re-connect clients that lost authorization on Telegram's side.
                for acct_id, client in list(self.clients.items()):
                    try:
                        if not await self._bounded(
                            lambda: client.is_user_authorized(), 15, "watchdog auth check", acct_id
                        ):
                            logger.warning("Watchdog: account %d no longer authorized, reconnecting", acct_id)
                            await self._stop_account_client(acct_id)
                            await self._start_account_client(acct_id)
                            if acct_id in self._accounts_needing_handlers():
                                self._install_handlers(acct_id)
                    except Exception as exc:
                        logger.warning("Watchdog: account %d check failed: %s", acct_id, exc)
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.error("Watchdog error: %s", e, exc_info=True)

    async def _prune_loop(self):
        while True:
            try:
                await asyncio.sleep(3600)
                deleted = await self._run_db(
                    self.db.prune_message_map, config.MAP_RETENTION_DAYS
                )
                if deleted:
                    logger.info("Pruned %s old message_map rows", deleted)
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.error("Prune error: %s", e)

    def _run(self, coro, timeout=30):
        """Schedule coroutine on the event-loop thread and block for result."""
        return asyncio.run_coroutine_threadsafe(coro, self.loop).result(timeout=timeout)

    def status(self) -> dict:
        authorized = bool(self.clients)
        connected = False
        accounts_status = []
        try:
            pairs = self._run_db_sync(self.db.get_pairs, only_enabled=False)
        except Exception as e:
            logger.warning("status() DB error: %s", e)
            pairs = []
        for acct_id, client in list(self.clients.items()):
            acct = self._run_db_sync(self.db.get_account, acct_id)
            name = acct["name"] if acct else f"Account #{acct_id}"
            try:
                c = client.is_connected()
            except Exception:
                c = False
            if c:
                connected = True
            accounts_status.append(
                {
                    "id": acct_id,
                    "name": name,
                    "connected": c,
                    "authorized": c,
                }
            )
        logger.info(
            "status() clients=%d connected=%s forwarding=%s",
            len(self.clients),
            connected,
            self._forwarding,
        )
        return {
            "connected": connected,
            "authorized": authorized,
            "forwarding": self._forwarding,
            "pairs": pairs,
            "accounts": accounts_status,
        }

    # ── Authentication ────────────────────────────────────────────────────────

    def send_code(self, phone: str, account_id: Optional[int] = None) -> dict:
        return self._run(self._async_send_code(phone, account_id), timeout=30)

    async def _async_send_code(self, phone: str, account_id: Optional[int] = None):
        if not account_id:
            raise RuntimeError("account_id required")
        # Clear any stale auth state for THIS account so we start fresh.
        self._auth_state.pop(account_id, None)
        if account_id in self.clients:
            client = self.clients[account_id]
            is_temp = False
        else:
            # Create a temporary client for new account auth. This also
            # covers the multi-account case where account_id is set but the
            # account is not authorized yet (so not in self.clients) — the
            # temp client must survive into _async_sign_in, otherwise the
            # sign-in step raises "No client available for sign-in".
            session = StringSession()
            client = TelegramClient(
                session,
                config.TG_API_ID,
                config.TG_API_HASH,
                connection=ConnectionTcpAbridged, connection_retries=3,
                retry_delay=3,
            )
            # Bounded timeout — see _start_account_client for rationale.
            await asyncio.wait_for(client.connect(), timeout=20)
            is_temp = True

        try:
            if not client.is_connected():
                await asyncio.wait_for(client.connect(), timeout=20)
            result = await asyncio.wait_for(
                client.send_code_request(phone), timeout=20
            )
        except SendCodeUnavailableError:
            if is_temp:
                await client.disconnect()
            raise RuntimeError(
                "RATE_LIMIT: Please wait 5-10 minutes before requesting a new code."
            )
        except FloodWaitError as e:
            if is_temp:
                await client.disconnect()
            raise RuntimeError(
                f"RATE_LIMIT: Please wait {e.seconds} seconds before trying again."
            )
        except Exception:
            # Connect/send failed (e.g. Telegram unreachable). Drop the temp
            # client so it doesn't keep reconnecting in the background.
            if is_temp:
                try:
                    await client.disconnect()
                except Exception:
                    pass
            raise
        # Keep the temp client for _async_sign_in whenever we created one
        # (multi-account auth on a not-yet-authorized account). Re-auth of an
        # already-connected account reuses self.clients[...].
        self._auth_state[account_id] = {
            "phone": phone,
            "phone_code_hash": result.phone_code_hash,
            "account_id": account_id,
            "_temp_client": client if is_temp else None,
        }
        if account_id:
            # Persist the pending session + phone so sign_in can be completed
            # even if the Passenger worker recycles between the two requests
            # (the temp client is in-memory otherwise), and so the account card
            # shows the phone before auth finishes.
            updates = {
                "phone_code_hash": result.phone_code_hash,
                "phone": phone,
                "status": "pending",
            }
            if is_temp:
                session_str = client.session.save()
                # Only persist if the session is non-empty — an empty session
                # has no auth key and will cause AuthKeyUnregisteredError on
                # sign_in when the worker recycles.
                if session_str:
                    updates["pending_session"] = session_str
            await self._run_db(self.db.update_account, account_id, **updates)
        return {"status": "code_sent", "phone": phone}

    def sign_in(
        self, code: str, password: Optional[str] = None, account_id: Optional[int] = None
    ) -> dict:
        return self._run(
            self._async_sign_in(code, password, account_id), timeout=60
        )

    async def _async_sign_in(
        self,
        code: str,
        password: Optional[str] = None,
        account_id: Optional[int] = None,
    ):
        if not account_id:
            raise RuntimeError("account_id required")
        state = self._auth_state.get(account_id, {})
        temp_client = state.get("_temp_client")
        if account_id in self.clients:
            client = self.clients[account_id]
        elif temp_client:
            client = temp_client
        else:
            # Worker recycled between send_code and sign_in (Passenger), or the
            # in-memory temp client was otherwise lost. Rebuild the client from
            # the pending StringSession persisted right after send_code_request
            # — the login code is bound to that session, so this is the only
            # way to complete sign-in across a process boundary.
            pending = await self._run_db(
                self.db.get_account_pending_session, account_id
            )
            if not pending:
                raise RuntimeError(
                    "LOGIN_SESSION_EXPIRED — request a new code."
                )
            client = TelegramClient(
                StringSession(pending),
                config.TG_API_ID,
                config.TG_API_HASH,
                connection=ConnectionTcpAbridged, connection_retries=3,
                retry_delay=3,
            )
            # Bounded timeout — see _start_account_client for rationale.
            await asyncio.wait_for(client.connect(), timeout=20)

        phone = state.get("phone")
        phone_code_hash = state.get("phone_code_hash")
        if not phone or not phone_code_hash:
            acct = await self._run_db(self.db.get_account, account_id)
            phone = (acct.get("phone") if acct else None) or config.TG_PHONE
            phone_code_hash = (acct.get("phone_code_hash") if acct else None)
            if not phone or not phone_code_hash:
                raise RuntimeError("No pending auth request")

        try:
            if password:
                # The UI sends a password only on the second attempt, after the
                # first returned 2FA_PASSWORD_REQUIRED — the code is already
                # consumed then, so go straight to the 2FA step instead of
                # re-submitting the (now invalid) code.
                await client.sign_in(password=password)
            else:
                await client.sign_in(phone, code, phone_code_hash=phone_code_hash)
        except SessionPasswordNeededError:
            auto_password = password or config.TG_PASSWORD
            if not auto_password:
                raise RuntimeError("2FA_PASSWORD_REQUIRED")
            try:
                await client.sign_in(password=auto_password)
            except PasswordHashInvalidError:
                raise RuntimeError("2FA password is incorrect")
        except PhoneCodeInvalidError:
            raise RuntimeError("INVALID_CODE")
        except (PhoneCodeExpiredError, PasswordHashInvalidError):
            raise RuntimeError("Invalid or expired code/2FA password")
        except AuthKeyUnregisteredError:
            # The pending session's auth key was never registered (incomplete
            # prior auth or worker recycled with a stale session). Clear it
            # so the user can request a fresh code.
            if account_id:
                await self._run_db(
                    self.db.update_account, account_id,
                    pending_session="", phone_code_hash="",
                )
            raise RuntimeError(
                "LOGIN_SESSION_EXPIRED — request a new code."
            )

        me = await client.get_me()
        logger.info(
            "Signed in as user_id=%s (account_id=%s)",
            getattr(me, "id", None),
            account_id,
        )
        session_str = client.session.save()
        acct_phone = getattr(me, "phone", None) or phone or ""

        if account_id:
            await self._run_db(
                self.db.update_account,
                account_id,
                session_string=session_str,
                is_authorized=True,
                status="active",
                phone=acct_phone,
                phone_code_hash="",
                pending_session="",
            )
            if account_id not in self.clients:
                self.clients[account_id] = client
            await self._load_pairs()
            if account_id in self._accounts_needing_handlers():
                self._install_handlers(account_id)
        else:
            # New account: create DB record
            name = f"Account {acct_phone}"
            new_id = await self._run_db(
                self.db.save_account, name, acct_phone, is_primary=not self.clients
            )
            await self._run_db(
                self.db.update_account,
                new_id,
                session_string=session_str,
                is_authorized=True,
                status="active",
            )
            self.clients[new_id] = client
            await self._load_pairs()
            if new_id in self._accounts_needing_handlers():
                self._install_handlers(new_id)

        self._auth_state.pop(account_id, None)
        return {"status": "authenticated", "account_id": account_id}

    # ── QR Login ────────────────────────────────────────────────────────────

    def start_qr_login(self, account_id: Optional[int] = None) -> dict:
        return self._run(self._async_start_qr(account_id), timeout=30)

    async def _async_start_qr(self, account_id: Optional[int] = None):
        logger.info("QR: starting for account_id=%s", account_id)

        if not account_id:
            raise RuntimeError("account_id required")

        if not config.TG_API_ID or not config.TG_API_HASH:
            raise RuntimeError(
                "Telegram API ID/Hash not configured. Set them under "
                "Settings \u2192 Telegram API first."
            )

        # Cancel any previous QR session for THIS account.
        prev = self._qr_state.get(account_id)
        if prev:
            logger.info("QR: replacing previous active QR session")
            # Only disconnect if the login didn't already complete — once
            # finished the client was moved into self.clients and is now a
            # live account that must NOT be torn down.
            if not prev.get("finished"):
                try:
                    await prev["client"].disconnect()
                except Exception:
                    pass
            self._qr_state.pop(account_id, None)

        session = StringSession()
        client = TelegramClient(
            session,
            config.TG_API_ID,
            config.TG_API_HASH,
            connection=ConnectionTcpAbridged, connection_retries=3,
            retry_delay=3,
        )
        # Bounded timeout — see _start_account_client for rationale.
        await asyncio.wait_for(client.connect(), timeout=20)
        logger.info("QR: client connected (dc=%s)", client.session.dc_id)

        qr = QRLogin(client, ignored_ids=[])
        try:
            await qr.recreate()
            logger.info(
                "QR: recreate() OK, token=%s...",
                qr.url[:40] if qr.url else "NONE",
            )
        except Exception as e:
            logger.error("QR: recreate() failed: %s", e)
            await client.disconnect()
            raise RuntimeError(f"QR initialization failed: {e}")

        self._qr_state[account_id] = {
            "client": client,
            "qr": qr,
            "account_id": account_id,
            "done": False,
            "error": None,
            "user": None,
            "started_at": time.time(),
            "deadline": time.time() + config.QR_LOGIN_TIMEOUT,
            "stage": "waiting_for_scan",
            "qr_url": qr.url,
            "refresh": False,
            "finished": False,
            "result_account_id": None,
            "user_id": None,
        }

        # Background task: wait for QR scan, refreshing the token on expiry,
        # and completing the login in the background so the polling HTTP
        # requests stay cheap.
        self.loop.create_task(self._qr_background_wait(account_id))
        logger.info("QR: background wait task started, returning qr_ready")

        return {
            "status": "qr_ready",
            "qr_url": qr.url,
            "timeout": config.QR_LOGIN_TIMEOUT,
        }

    async def _qr_background_wait(self, account_id: int):
        state = self._qr_state.get(account_id)
        if not state:
            logger.warning("QR: _qr_background_wait called but no state for account %s", account_id)
            return
        logger.info("QR: background wait started")
        try:
            while time.time() < state["deadline"]:
                qr = state["qr"]
                # Server-reported token TTL; recreate before it expires.
                try:
                    ttl = (qr.expires - datetime.now(timezone.utc)).total_seconds()
                except Exception:
                    ttl = 30
                remaining = state["deadline"] - time.time()
                wait_for = max(1.0, min(ttl, remaining))
                try:
                    user = await qr.wait(timeout=wait_for)
                    state["user"] = user
                    state["stage"] = "scanned"
                    logger.info(
                        "QR: scanned user=%s, finishing login",
                        getattr(user, "id", None),
                    )
                    await self._finish_qr_login(state)
                    state["stage"] = "authenticated"
                    state["done"] = True
                    logger.info("QR: login finished in background")
                    return
                except SessionPasswordNeededError:
                    logger.info("QR: 2FA required, waiting for password")
                    state["stage"] = "2fa_required"
                    state["done"] = True
                    return
                except (asyncio.TimeoutError, AuthTokenExpiredError, AuthTokenInvalidError):
                    # Token expired/invalid (or wait window elapsed) before a
                    # scan. Telegram raises AuthTokenExpiredError (not a plain
                    # asyncio timeout) once the login token dies, so recreate a
                    # fresh token and keep the QR alive until the overall
                    # deadline — the poller re-renders the new code via the
                    # "refresh" status.
                    logger.info("QR: token expired/invalid, recreating")
                    try:
                        client = state["client"]
                        # recreate() sends a request, so the client must be
                        # connected (it can drop on a slow/shared host).
                        if not client.is_connected():
                            await asyncio.wait_for(client.connect(), timeout=15)
                        await qr.recreate()
                        state["qr_url"] = qr.url
                        state["refresh"] = True
                    except Exception as e:
                        logger.error("QR: recreate after expiry failed: %s", e)
                        state["error"] = f"QR refresh failed: {e}"
                        state["done"] = True
                        state["stage"] = "error"
                        return
                    continue
            # Overall deadline reached without a scan.
            state["error"] = "QR login timed out"
            state["done"] = True
            state["stage"] = "timeout"
            logger.warning("QR: overall timeout reached")
        except Exception as e:
            logger.error("QR: background wait error: %s (%s)", e, type(e).__name__)
            state["error"] = str(e)
            state["done"] = True
            state["stage"] = "error"

    def complete_qr_2fa(self, account_id: int, password: str) -> dict:
        return self._run(self._async_complete_qr_2fa(account_id, password), timeout=60)

    async def _async_complete_qr_2fa(self, account_id: int, password: str):
        state = self._qr_state.get(account_id)
        if not state:
            raise RuntimeError("No active QR login session for this account")

        # Already finished (e.g. background completed after 2FA prompt).
        if state.get("stage") == "authenticated" and state.get("finished"):
            return {
                "status": "authenticated",
                "account_id": state.get("result_account_id"),
                "user_id": state.get("user_id"),
            }

        client = state["client"]
        logger.info("QR: 2FA submit for account_id=%s", account_id)
        try:
            user = await client.sign_in(password=password)
        except SessionPasswordNeededError:
            raise RuntimeError("Incorrect 2FA password")
        except Exception as e:
            logger.error("QR: 2FA sign_in failed: %s (%s)", e, type(e).__name__)
            raise RuntimeError(f"2FA failed: {e}")

        state["user"] = user
        state["stage"] = "scanned"
        await self._finish_qr_login(state)
        state["stage"] = "authenticated"
        state["done"] = True
        return {
            "status": "authenticated",
            "account_id": state.get("result_account_id"),
            "user_id": state.get("user_id"),
        }

    async def _finish_qr_login(self, state: dict, account_id: Optional[int] = None):
        """Persist the session + (re)connect the client after a successful scan.

        Does NOT clear ``_qr_state`` — the caller (background task / 2FA path)
        is responsible for advancing the stage so the poller can observe the
        result. This keeps the heavy DB/pairs work off the HTTP request thread.
        """
        user = state["user"]
        if not user:
            raise RuntimeError("QR login failed: no user returned")
        client = state["client"]

        session_str = client.session.save()
        phone = getattr(user, "phone", "") or ""

        target_id = state.get("account_id") or account_id
        if target_id and target_id > 0:
            logger.info("QR: updating existing account %d, stopping old client", target_id)
            await self._stop_account_client(target_id)
            await self._run_db(
                self.db.update_account,
                target_id,
                session_string=session_str,
                is_authorized=True,
                status="active",
                phone=phone or "",
            )
            self.clients[target_id] = client
            logger.info("QR: account %d reconnected with new session", target_id)
        else:
            name = f"Account {phone or 'QR'}"
            logger.info("QR: creating new account: %s", name)
            new_id = await self._run_db(
                self.db.save_account, name, phone or "", is_primary=not self.clients
            )
            await self._run_db(
                self.db.update_account,
                new_id,
                session_string=session_str,
                is_authorized=True,
                status="active",
            )
            self.clients[new_id] = client
            target_id = new_id
            logger.info("QR: new account %d created", new_id)

        await self._load_pairs()
        if target_id in self._accounts_needing_handlers():
            self._install_handlers(target_id)
            logger.info("QR: handlers installed for account %d", target_id)

        state["finished"] = True
        state["result_account_id"] = target_id
        state["user_id"] = getattr(user, "id", None)
        return {
            "status": "authenticated",
            "account_id": target_id,
            "user_id": getattr(user, "id", None),
        }

    def cancel_qr_login(self, account_id: Optional[int] = None) -> dict:
        return self._run(self._async_cancel_qr(account_id), timeout=10)

    async def _async_cancel_qr(self, account_id: Optional[int] = None):
        if not account_id:
            return {"status": "no_session"}
        state = self._qr_state.pop(account_id, None)
        if not state:
            return {"status": "no_session"}
        logger.info("QR: cancelling login for account_id=%s", state.get("account_id"))
        # Don't disconnect a client that was already promoted into self.clients.
        if not state.get("finished"):
            try:
                await state["client"].disconnect()
            except Exception as e:
                logger.warning("QR: disconnect during cancel: %s", e)
        return {"status": "cancelled"}

    def check_qr_login(self, account_id: Optional[int] = None) -> dict:
        # The check is a pure status reader and returns immediately; a short
        # timeout keeps the HTTP request responsive if the event loop stalls.
        return self._run(self._async_check_qr(account_id), timeout=20)

    async def _async_check_qr(self, account_id: Optional[int] = None):
        if not account_id:
            raise RuntimeError("No active QR login session")
        state = self._qr_state.get(account_id)
        if not state:
            raise RuntimeError("No active QR login session for this account")

        elapsed = round(time.time() - state.get("started_at", time.time()))
        stage = state.get("stage", "unknown")
        logger.debug(
            "QR: check stage=%s done=%s elapsed=%ss error=%s",
            stage, state.get("done"), elapsed, state.get("error"),
        )

        if state.get("error"):
            err = state["error"]
            self._qr_state.pop(account_id, None)
            if not state.get("finished"):
                try:
                    await state["client"].disconnect()
                except Exception:
                    pass
            logger.error("QR: returning error to client: %s", err)
            raise RuntimeError(err)

        # A new QR token was generated after the previous one expired —
        # surface the fresh URL so the client can re-render the image.
        if state.get("refresh") and not state.get("done"):
            state["refresh"] = False
            return {
                "status": "refresh",
                "qr_url": state.get("qr_url"),
                "elapsed": elapsed,
            }

        if not state.get("done"):
            return {"status": "waiting", "stage": stage, "elapsed": elapsed}

        # 2FA required — user scanned QR but needs a password.
        if state.get("stage") == "2fa_required":
            return {"status": "2fa_required", "stage": "2fa_required"}

        # Login completed in the background.
        if state.get("stage") == "authenticated" and state.get("finished"):
            result = {
                "status": "authenticated",
                "account_id": state.get("result_account_id"),
                "user_id": state.get("user_id"),
            }
            self._qr_state.pop(account_id, None)
            logger.info(
                "QR: reporting authenticated account_id=%s",
                result["account_id"],
            )
            return result

        # done=True but no user / unknown stage — clean up and fail.
        self._qr_state.pop(account_id, None)
        if not state.get("finished"):
            try:
                await state["client"].disconnect()
            except Exception:
                pass
        logger.error("QR: done but no user returned (stage=%s)", stage)
        raise RuntimeError("QR login failed: no user returned")

    def test_telegram_api(self) -> dict:
        """Validate TG_API_ID / TG_API_HASH with a lightweight call."""
        return self._run(self._async_test_telegram_api(), timeout=25)

    async def _async_test_telegram_api(self):
        if not config.TG_API_ID or not config.TG_API_HASH:
            return {"ok": False, "error": "API ID/Hash not configured"}
        client = TelegramClient(
            StringSession(),
            config.TG_API_ID,
            config.TG_API_HASH,
            connection=ConnectionTcpAbridged, connection_retries=2,
            retry_delay=2,
        )
        try:
            await asyncio.wait_for(client.connect(), timeout=15)
            # help.getNearestDc works for unauthorized clients and validates
            # the api_id/api_hash against Telegram's servers.
            res = await client(functions.help.GetNearestDcRequest())
            return {
                "ok": True,
                "dc": client.session.dc_id,
                "nearest_dc": getattr(res, "nearest_dc", None),
                "country": getattr(res, "country", None),
            }
        except Exception as e:
            logger.warning("test_telegram_api failed: %s", e)
            return {"ok": False, "error": f"{type(e).__name__}: {e}"}
        finally:
            try:
                await client.disconnect()
            except Exception:
                pass

    def get_dialogs(self, account_id: Optional[int] = None) -> list:
        return self._run(
            self._async_get_dialogs(refresh=False, account_id=account_id), timeout=30
        )

    def refresh_dialogs(self, account_id: Optional[int] = None) -> dict:
        dialogs = self._run(
            self._async_get_dialogs(refresh=True, account_id=account_id), timeout=60
        )
        return {"status": "ok", "count": len(dialogs)}

    def _dialog_cache_path(self, account_id: Optional[int]) -> str:
        """Per-account cache path so one account's channels never shadow another's."""
        suffix = account_id if account_id is not None else "none"
        if config.DIALOGS_CACHE_FILE.endswith(".json"):
            return f"{config.DIALOGS_CACHE_FILE[:-5]}.{suffix}.json"
        return f"{config.DIALOGS_CACHE_FILE}.{suffix}.json"

    def _load_cached_dialogs(self, account_id: int) -> Optional[list]:
        try:
            path = self._dialog_cache_path(account_id)
            if not os.path.exists(path):
                return None
            ttl = config.DIALOGS_CACHE_TTL_SECONDS
            if ttl > 0 and time.time() - os.path.getmtime(path) > ttl:
                return None
            with open(path, "r", encoding="utf-8") as f:
                data = json.load(f)
            # The cache stores {"id","title","type","category"}. Older or
            # partially-saved caches may lack "category"; those still load so
            # the UI can search, and get re-saved (with category) on refresh.
            if isinstance(data, list):
                # Robust malformed-entry guard: skip anything unusable.
                data = [
                    d
                    for d in data
                    if isinstance(d, dict)
                    and d.get("id") is not None
                    and d.get("title")
                ]
                if data:
                    logger.info(
                        "Loaded %d dialog(s) from cache (account %s)",
                        len(data),
                        account_id,
                    )
                    return data
        except Exception as e:
            logger.warning("Failed to load cached dialogs for account %s: %s", account_id, e)
        return None

    def _save_cached_dialogs(self, dialogs: list, account_id: int):
        try:
            path = self._dialog_cache_path(account_id)
            with open(path, "w", encoding="utf-8") as f:
                json.dump(dialogs, f, ensure_ascii=False, indent=2)
            logger.info("Saved %d dialog(s) to cache (account %s)", len(dialogs), account_id)
        except Exception as e:
            logger.warning("Failed to save dialog cache for account %s: %s", account_id, e)

    async def _async_get_dialogs(
        self, refresh: bool = False, account_id: Optional[int] = None
    ):
        resolved_account_id = self._resolve_account_id(account_id)
        try:
            client = self._get_client(resolved_account_id)
        except RuntimeError:
            # Passive worker (the leader holds the core lock): no Telegram
            # client here — serve the shared on-disk dialog cache so the
            # pairs UI still works on whichever worker a request lands on.
            client = None
        if client is None:
            cached = self._load_cached_dialogs(resolved_account_id)
            if cached is not None:
                return cached
            raise RuntimeError("Telegram client not connected on this worker")
        if not await client.is_user_authorized():
            raise RuntimeError("Not authorized")
        if not refresh:
            cached = self._load_cached_dialogs(resolved_account_id)
            if cached is not None:
                return cached
        dialogs = []
        async for dialog in client.iter_dialogs():
            entity = dialog.entity
            if entity is None:
                continue
            title = (
                getattr(entity, "title", None)
                or f"{getattr(entity, 'first_name', '')} {getattr(entity, 'last_name', '')}".strip()
            )
            dialogs.append(
                {
                    "id": dialog.id,
                    "title": title or "Unnamed",
                    "type": type(entity).__name__,
                    "category": self._classify_dialog(entity),
                }
            )
        self._save_cached_dialogs(dialogs, resolved_account_id)
        return dialogs

    def _classify_dialog(self, entity) -> str:
        if getattr(entity, "bot", False):
            return "bot"
        if getattr(entity, "creator", False):
            return "own"
        if getattr(entity, "username", None):
            return "public"
        return "private"

    # ── Pair management ───────────────────────────────────────────────────────

    def get_pairs(self) -> List[dict]:
        return self._run_db_sync(self.db.get_pairs, only_enabled=False)

    def get_pair(self, pair_id: int) -> Optional[dict]:
        return self._run_db_sync(self.db.get_pair_by_id, pair_id)

    def create_pair(self, source_id: int, dest_id: int, **fields) -> dict:
        pair_id = self._run_db_sync(self.db.save_pair, source_id, dest_id, **fields)
        self._run(self._load_pairs(), timeout=10)
        for acct_id in self._accounts_needing_handlers():
            self._run(self._ensure_handlers(acct_id), timeout=10)
        return {"status": "ok", "id": pair_id}

    def update_pair(self, pair_id: int, **fields) -> dict:
        ok = self._run_db_sync(self.db.update_pair, pair_id, **fields)
        self._run(self._load_pairs(), timeout=10)
        return {"status": "ok" if ok else "not_found"}

    def delete_pair(self, pair_id: int) -> dict:
        ok = self._run_db_sync(self.db.delete_pair, pair_id)
        self._run(self._load_pairs(), timeout=10)
        return {"status": "ok" if ok else "not_found"}

    def toggle_pair(self, pair_id: int) -> dict:
        ok = self._run_db_sync(self.db.toggle_pair, pair_id)
        self._run(self._load_pairs(), timeout=10)
        return {"status": "ok" if ok else "not_found"}

    # ── LLM engine config ─────────────────────────────────────────────────────

    def reload_llm_config(self) -> dict:
        """Reload LLM config from DB (blocking) and return engine status."""
        self._run_db_sync(self._llm_engine.reload_from_db)
        status = self._llm_engine.status()
        logger.info("LLM engine reloaded: %s", status)
        return status

    def llm_status(self) -> dict:
        return self._llm_engine.status()

    def llm_list_models(self, base_url: str = "", api_key: str = "") -> dict:
        return self._llm_engine.list_models(base_url, api_key)

    def llm_test_connection(
        self, base_url: str = "", api_key: str = "", model: str = ""
    ) -> dict:
        return self._llm_engine.test_connection(base_url, api_key, model)

    def start_forwarding(self) -> dict:
        self._run(self._load_pairs(), timeout=10)
        for acct_id in self._accounts_needing_handlers():
            self._run(self._ensure_handlers(acct_id), timeout=10)
        if self._pairs:
            return {"status": "started", "pairs": sum(len(v) for v in self._pairs.values())}
        return {"status": "no_pair"}

    def stop_forwarding(self) -> dict:
        for acct_id in list(self._handlers_accounts):
            self._remove_handlers(acct_id)
        return {"status": "stopped"}

    # ── Test signal emission (delegated to MessageForwarder) ─────────────────

    def emit_test_signal(self, pair_id: int, text: str) -> dict:
        pair = self._run_db_sync(self.db.get_pair_by_id, pair_id)
        if not pair:
            return {"status": "error", "error": "Pair not found"}
        if not pair.get("enabled"):
            return {"status": "error", "error": "Pair is disabled"}
        if not self._forwarder:
            return {
                "status": "error",
                "error": "Forwarder not initialized; start forwarding first",
            }
        return self._run(
            self._forwarder.emit_test_signal(pair, text),
            timeout=30,
        )

    async def _ensure_handlers(self, account_id: int):
        client = self.clients.get(account_id)
        if client and await client.is_user_authorized():
            self._install_handlers(account_id)

    # ── Lifecycle ─────────────────────────────────────────────────────────────

    async def _disconnect_all(self):
        for acct_id, client in list(self.clients.items()):
            try:
                await client.disconnect()
            except Exception:
                pass

    def shutdown(self):
        logger.info("Shutting down forwarder core...")
        try:
            for acct_id in list(self._handlers_accounts):
                self._remove_handlers(acct_id)
        except Exception as e:
            logger.warning("Handler removal error: %s", e)
        try:
            if self.clients:
                self._run(self._disconnect_all(), timeout=15)
        except Exception as e:
            logger.warning("Client disconnect error: %s", e)
        try:
            if self._ctrader_started and self._ctrader_client:
                self._run(self._ctrader_client.stop(), timeout=15)
        except Exception as e:
            logger.warning("cTrader stop error: %s", e)
        for task in (
            self._memory_task,
            getattr(self, "_prune_task", None),
            getattr(self, "_pairs_sync_task", None),
            getattr(self, "_watchdog_task", None),
        ):
            if task and not task.done():
                task.cancel()
        if self._lock_conn:
            try:
                self.db.release_core_lock(self._lock_conn)
            except Exception as e:
                logger.warning("Core lock release error: %s", e)
            self._lock_conn = None
        try:
            self._run_db_sync(self.db.dispose)
        except Exception as e:
            logger.warning("DB pool dispose error: %s", e)
        try:
            self._executor.shutdown(wait=True, cancel_futures=True)
        except Exception as e:
            logger.warning("Executor shutdown error: %s", e)
        if self.loop.is_running():
            self.loop.call_soon_threadsafe(self.loop.stop)
        logger.info("Shutdown complete")
