"""Каталог MCP-тулов tgclient (полный: core + media/voice/rounds + справочники
+ логин-тулы текст-режима).

Правила:
- Ошибки — данные: DomainError и неожиданные исключения ловит `_tool` и
  возвращает {"error": code, "detail": …} (канон mcp.md) вместо падения.
- Scoping: строка accounts сверяется с current_user_id() (контекст из
  require_mcp); админ-роль видит все аккаунты, обычный ключ — только свои.
- Почти все тулы опционально принимают account_id; None → ровно один
  активный аккаунт владельца, иначе 422 с перечнем id/label.
- Мутации проходят check_write (per-account token bucket) и call_with_flood_guard;
  после мутаций save_session (StringSession держит entity-кеш).
- dialog_id — маркированный int (юзер >0, малая группа -id, канал -100id);
  резолв через client.get_input_entity(dialog_id) из кеша сессии.
"""

import base64
from contextlib import suppress
from pathlib import Path

from telethon import functions, errors
from telethon.tl.types import DocumentAttributeAudio, DocumentAttributeVideo

from app.config import get_settings
from app.db import get_db
from app.errors import DomainError
from app.mcp.context import current_user_id, is_admin
from app.mcp.serializers import (
    dialog_dict,
    entity_dict,
    media_meta,
    message_dict,
)
from app.tg.limits import check_write
from app.tg.login_flow import (
    login_cancel,
    login_code,
    login_password,
    login_start,
    login_status,
)
from app.tg.waveform import (
    WAVEFORM_SAMPLES,
    base64_or_none,
    ogg_opus_duration,
    pack_waveform,
)

HISTORY_CAP = 200


# --------------------------------------------------------------------------
# Инфраструктура тулов
# --------------------------------------------------------------------------

def _manager():
    from app.main import get_account_manager

    return get_account_manager()


async def _tool(coro_fn):
    """Единый error-boundary: DomainError → {"error", "detail"}, прочее → 500."""
    try:
        return await coro_fn()
    except DomainError as exc:
        return {"error": exc.code, "detail": exc.detail}
    except errors.FloodWaitError as exc:
        return {"error": 429, "detail": f"telegram flood wait: {exc.seconds}s — retry later"}
    except errors.EntityNotFoundError:
        return {"error": 404, "detail": "chat/dialog not found for this account (resolve id via dialogs_list)"}
    except Exception as exc:  # noqa: BLE001 — ошибка тулу как данные, не как падение
        return {"error": 500, "detail": f"{type(exc).__name__}: {exc}"}


async def _own_account_rows(db) -> list:
    """Строки accounts текущего mcp-юзера; админ — все (с пометкой владельца)."""
    if is_admin():
        cursor = await db.execute(
            "SELECT a.*, u.email AS owner_email FROM accounts a"
            " LEFT JOIN users u ON u.user_id = a.user_id ORDER BY a.id"
        )
    else:
        cursor = await db.execute(
            "SELECT a.*, NULL AS owner_email FROM accounts a WHERE user_id = ? ORDER BY a.id",
            (current_user_id(),),
        )
    return await cursor.fetchall()


async def _resolve_client(db, account_id, *, writable: bool = False):
    """account_id → (client, row). None → авто-выбор единственного активного.

    writable=True — только свои аккаунты (админ-ключ пишет в чужие не должен,
    даже если видит их: мутации по явной просьбе владельца — не нам решать).
    """
    rows = await _own_account_rows(db)
    if account_id is None:
        own = [r for r in rows if r["user_id"] == current_user_id()]
        pool = own if writable else rows
        if not pool:
            raise DomainError(404, "нет аккаунтов — добавьте через account_login_start(phone)")
        active = [r for r in pool if r["status"] == "active"]
        if not active:
            if writable:
                raise DomainError(409, "нет активного своего аккаунта — добавьте через account_login_*")
            active = pool[:1]
        if len(active) == 1:
            row = active[0]
        else:
            listed = "; ".join(f"{r['id']}={r['label'] or r['username'] or r['phone']}" for r in active[:10])
            raise DomainError(
                422, f"укажите account_id — активных несколько: {listed}"
            )
    else:
        matches = [r for r in rows if r["id"] == account_id]
        if not matches:
            raise DomainError(403, f"account #{account_id} не ваш (или не существует)")
        row = matches[0]
    if row is None:
        raise DomainError(404, f"account #{account_id} not found")
    if writable and row["user_id"] != current_user_id():
        raise DomainError(403, f"account #{account_id} принадлежит другому пользователю — писать в него нельзя")
    if row["status"] == "logged_out":
        raise DomainError(409, f"account #{account_id} logged out — логиньте заново через account_login_*")
    client = await _manager().get_client(account_id)
    return client, row


async def _client_for_dialog(db, account_id, dialog_id, *, writable: bool = False):
    if dialog_id in (None, "", "0", 0):
        raise DomainError(422, "dialog_id обязателен (маркированный int из dialogs_list)")
    client, row = await _resolve_client(db, account_id, writable=writable)
    if not isinstance(dialog_id, str):
        dialog_id = str(dialog_id)
    try:
        if dialog_id.lstrip("-").isdigit():
            entity = await client.get_input_entity(int(dialog_id))
        else:  # username-ссылка — сетевой резолв, попадёт в кеш сессии
            entity = await client.get_entity(dialog_id.removeprefix("@"))
    except Exception as exc:  # noqa: BLE001 — не найден в кеше/у этого акка
        raise DomainError(404, f"dialog {dialog_id} недоступен на этом аккаунте"
                               f" (id не найден в сессии — проверьте dialogs_list)") from exc
    return client, row, entity


def _limit(value: int | None, default: int, cap: int) -> int:
    """клампим пользовательский limit: 1..cap, default без значения."""
    if value is None:
        return default
    return max(1, min(int(value), cap))


def _write_guard(account_id: int) -> dict | None:
    settings = get_settings()
    return check_write(account_id, settings.write_limit_n, settings.write_limit_window)


async def _saved(client, account_id: int, message) -> dict:
    from app.main import get_account_manager

    await get_account_manager().save_session(account_id)
    return message_dict(message)


# --------------------------------------------------------------------------
# Core: аккаунты, me, диалоги, сообщения
# --------------------------------------------------------------------------

def register_tools(mcp) -> None:
    """Регистрация всего каталога на FastMCP-инстансе (app/mcp/server.py)."""

    # --- Аккаунты ---------------------------------------------------------

    @mcp.tool()
    async def accounts_list() -> dict:
        """📋 Телеграм-аккаунты владельца ключа: id, телефон, username, статус
        (active/logged_out/error), label, connected. account_id берите отсюда.

        Admin-ключ видит все аккаунты сервиса (поле owner_email — чей)."""
        return await _tool(_accounts_list)


    async def _accounts_list() -> dict:
        db = get_db()
        manager = _manager()
        out = []
        for r in await _own_account_rows(db):
            out.append({
                "id": r["id"],
                "phone": r["phone"],
                "label": r["label"],
                "tg_user_id": r["tg_user_id"],
                "username": r["username"],
                "display_name": r["display_name"],
                "status": r["status"],
                "error": r["error"],
                "connected": r["id"] in manager.clients,
                "owner_email": r["owner_email"],
                "last_used_at": r["last_used_at"],
            })
        return {"accounts": out}


    @mcp.tool()
    async def me_get(account_id: int | None = None) -> dict:
        """👤 Кто я на аккаунте account_id: имя, username, телефон, id (маркированный), premium, бот или нет."""
        return await _tool(lambda: _me_get(account_id))


    async def _me_get(account_id: int | None) -> dict:
        client, row = await _resolve_client(get_db(), account_id)
        me = await client.get_me()
        out = entity_dict(me)
        out["phone"] = getattr(me, "phone", "") or row["phone"]
        out["account_id"] = row["id"]
        out["label"] = row["label"]
        out["premium"] = bool(getattr(me, "premium", False))
        return out


    # --- Диалоги ----------------------------------------------------------

    @mcp.tool()
    async def dialogs_list(account_id: int | None = None, limit: int = 30) -> dict:
        """💬 Список диалогов (чаты): id, type (user/chat/channel), title,
        username, unread_count, pinned, last_message. limit ≤ 200."""
        return await _tool(lambda: _dialogs_list(account_id, limit))


    async def _dialogs_list(account_id, limit) -> dict:
        client, _row = await _resolve_client(get_db(), account_id)
        limit = _limit(limit, 30, HISTORY_CAP)
        out = []
        async for dialog in client.iter_dialogs(limit=limit):
            out.append(dialog_dict(dialog, dialog.entity))
        return {"dialogs": out}


    @mcp.tool()
    async def chat_info(account_id: int | None = None, dialog_id: str = "") -> dict:
        """ℹ️ Инфо о чате/канале/юзере: полное описание, счётчики участников.
        dialog_id — маркированный int из dialogs_list или username («@name»)."""
        return await _tool(lambda: _chat_info(account_id, dialog_id))


    async def _chat_info(account_id, dialog_id) -> dict:
        client, _row, entity = await _client_for_dialog(get_db(), account_id, dialog_id)
        full = await client.get_entity(entity)
        out = entity_dict(full)
        with suppress(Exception):
            if type(full).__name__ == "Channel":
                fch = await client(functions.channels.GetFullChannelRequest(channel=entity))
                out["description"] = getattr(fch.full_chat, "about", "") or ""
                out["participants_count"] = getattr(fch.full_chat, "participants_count", None)
            elif type(full).__name__ == "User":
                out["bio"] = getattr(full, "about", "") or ""
                out["phone"] = getattr(full, "phone", "") or ""
            elif type(full).__name__ == "Chat":
                fch = await client(functions.messages.GetFullChatRequest(chat_id=full.id))
                out["description"] = getattr(getattr(fch, "full_chat", None), "about", "") or ""
                out["participants_count"] = len(getattr(fch, "users", []) or [])
        return out


    @mcp.tool()
    async def chat_participants(account_id: int | None = None, dialog_id: int = 0,
                                limit: int = 50) -> dict:
        """👥 Участники чата/канала (limit ≤ 200): id, type, title/имя, username."""
        return await _tool(lambda: _chat_participants(account_id, dialog_id, limit))


    async def _chat_participants(account_id, dialog_id, limit) -> dict:
        client, _row, entity = await _client_for_dialog(get_db(), account_id, dialog_id)
        limit = _limit(limit, 50, HISTORY_CAP)
        kind = type(entity).__name__  # InputPeerChannel / InputPeerChat / InputPeerUser…
        out = []
        if kind == "InputPeerChannel":  # канал/супергруппа
            participants = await client.get_participants(entity, limit=limit)
            for user in participants:
                out.append(entity_dict(user, short=True))
        elif kind == "InputPeerChat":  # базовая малая группа
            chat_full = await client(
                functions.messages.GetFullChatRequest(chat_id=entity.chat_id)
            )
            for user in (chat_full.users or [])[:limit]:
                out.append(entity_dict(user, short=True))
        else:  # юзер — сам участник
            full = await client.get_entity(entity)
            out.append(entity_dict(full, short=True))
        return {"participants": out}


    # --- История и поиск ----------------------------------------------------

    @mcp.tool()
    async def messages_history(account_id: int | None = None, dialog_id: int = 0,
                               limit: int = 30, offset_id: int = 0,
                               min_id: int = 0) -> dict:
        """📜 История сообщений чата (новые первыми, limit ≤ 200).
        offset_id — id последнего сообщения страницы для пагинации вглубь;
        min_id — нижняя граница id (только новые после last seen)."""
        return await _tool(
            lambda: _messages_history(account_id, dialog_id, limit, offset_id, min_id)
        )


    async def _messages_history(account_id, dialog_id, limit, offset_id, min_id) -> dict:
        client, _row, entity = await _client_for_dialog(get_db(), account_id, dialog_id)
        limit = _limit(limit, 30, HISTORY_CAP)
        msgs = await client.get_messages(
            entity, limit=limit, offset_id=offset_id or None, min_id=min_id or None,
        )
        return {"messages": [message_dict(m) for m in msgs]}


    @mcp.tool()
    async def messages_read(account_id: int | None = None, dialog_id: int = 0,
                            message_id: int = 0) -> dict:
        """🔍 Одно сообщение (включая медиа-метаданные: is_voice/is_round/duration/waveform)."""
        return await _tool(lambda: _messages_read(account_id, dialog_id, message_id))


    async def _messages_read(account_id, dialog_id, message_id) -> dict:
        client, _row, entity = await _client_for_dialog(get_db(), account_id, dialog_id)
        msgs = await client.get_messages(entity, ids=[int(message_id)])
        if not msgs or msgs[0] is None:
            raise DomainError(404, f"message #{message_id} не найден в чате")
        return message_dict(msgs[0])


    @mcp.tool()
    async def messages_search(account_id: int | None = None, dialog_id: int = 0,
                              query: str = "", from_user: str = "", limit: int = 20) -> dict:
        """🔎 Поиск сообщений: в одном чате (dialog_id) или во всех диалогах (без dialog_id).
        from_user — username или маркированный id."""
        return await _tool(
            lambda: _messages_search(account_id, dialog_id, query, from_user, limit)
        )


    async def _messages_search(account_id, dialog_id, query, from_user, limit) -> dict:
        db = get_db()
        limit = _limit(limit, 20, HISTORY_CAP)
        out = []
        if dialog_id in (None, 0):
            if not query:
                raise DomainError(422, "нужен query или dialog_id")
            if from_user:
                raise DomainError(
                    422, "from_user доступен только при заданном dialog_id (глобальный поиск Telegram не фильтрует по автору)"
                )
            # глобальный поиск: по всем диалогам (SearchGlobal, InputPeerEmpty)
            client, _row = await _resolve_client(db, account_id)
            from telethon.tl.types import InputMessagesFilterEmpty, InputPeerEmpty

            result = await client(functions.messages.SearchGlobalRequest(
                q=query,
                filter=InputMessagesFilterEmpty(),
                min_date=None, max_date=None,
                offset_rate=0, offset_peer=InputPeerEmpty(),
                offset_id=0, limit=limit,
            ))
            for message in getattr(result, "messages", []) or []:
                out.append(message_dict(message))
            return {"messages": out}
        client, _row, entity = await _client_for_dialog(db, account_id, dialog_id)
        from_user_id = await _resolve_from_user(client, from_user) if from_user else None
        msgs = await client.get_messages(
            entity, limit=limit, search=query or None, from_user=from_user_id,
        )
        return {"messages": [message_dict(m) for m in msgs]}


    async def _resolve_from_user(client, from_user: str):
        """'@username' / маркированный id → InputUser для фильтра поиска."""
        with suppress(Exception):
            return await client.get_input_entity(from_user)
        raise DomainError(404, f"from_user '{from_user}' не найден на этом аккаунте")


    # --- Отправка / правка / удаление --------------------------------------

    @mcp.tool()
    async def message_send(account_id: int | None = None, dialog_id: int = 0,
                           text: str = "", reply_to: int = 0, silent: bool = False) -> dict:
        """✍️ Отправить текст в чат (reply_to — id сообщения для ответа,
        silent — без уведомления). Мутация: только по явной просьбе пользователя."""
        return await _tool(lambda: _message_send(account_id, dialog_id, text, reply_to, silent))


    async def _message_send(account_id, dialog_id, text, reply_to, silent) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        if not text.strip():
            raise DomainError(422, "text пустой")
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        message = await call_guard(_manager(), row["id"],
                                   lambda: client.send_message(
                                       entity, text.strip(), reply_to=reply_to or None,
                                       silent=silent))
        return await _saved(client, row["id"], message)


    @mcp.tool()
    async def message_reply(account_id: int | None = None, dialog_id: int = 0,
                            message_id: int = 0, text: str = "") -> dict:
        """↩️ Ответить на конкретное сообщение (то же, что message_send с reply_to)."""
        return await _tool(lambda: _message_reply(account_id, dialog_id, message_id, text))


    async def _message_reply(account_id, dialog_id, message_id, text) -> dict:
        return await _message_send(account_id, dialog_id, text, message_id, False)


    @mcp.tool()
    async def message_edit(account_id: int | None = None, dialog_id: int = 0,
                           message_id: int = 0, text: str = "") -> dict:
        """✏️ Правка своего сообщения (текст). Мутация; чужие править нельзя."""
        return await _tool(lambda: _message_edit(account_id, dialog_id, message_id, text))


    async def _message_edit(account_id, dialog_id, message_id, text) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        if not text.strip():
            raise DomainError(422, "text пустой")
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        try:
            message = await call_guard(_manager(), row["id"],
                                       lambda: client.edit_message(
                                           entity, int(message_id), text.strip()))
        except errors.MessageAuthorRequiredError as exc:
            raise DomainError(403, "править можно только свои сообщения") from exc
        except errors.MessageNotModifiedError as exc:
            raise DomainError(409, "текст сообщения не изменился") from exc
        return await _saved(client, row["id"], message)


    @mcp.tool()
    async def message_delete(account_id: int | None = None, dialog_id: int = 0,
                             message_ids: list[int] | None = None,
                             revoke: bool = True) -> dict:
        """🗑️ Удалить сообщения (список id). revoke=True — у всех участников
        (иначе только у себя). Необратимо — только по явной просьбе."""
        return await _tool(
            lambda: _message_delete(account_id, dialog_id, message_ids, revoke)
        )


    async def _message_delete(account_id, dialog_id, message_ids, revoke) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        ids = [int(i) for i in (message_ids or []) if int(i) > 0]
        if not ids:
            raise DomainError(422, "message_ids пустой")
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        await call_guard(_manager(), row["id"],
                         lambda: client.delete_messages(entity, ids, revoke=revoke))
        from app.main import get_account_manager

        await get_account_manager().save_session(row["id"])
        return {"status": "ok", "deleted": ids, "revoke": revoke}


    # --- Медиа --------------------------------------------------------------

    @mcp.tool()
    async def download_media(account_id: int | None = None, dialog_id: int = 0,
                             message_id: int = 0, max_bytes: int = 0) -> dict:
        """⬇️ Скачать медиа сообщения (фото/файл/голосовое/кружок) → base64.
        max_bytes — лимит (≤ TGCLIENT_MEDIA_MAX_BYTES, по умолчанию он же).
        Метаданные (voice/round/duration/waveform) уже в messages_read."""
        return await _tool(
            lambda: _download_media(account_id, dialog_id, message_id, max_bytes)
        )


    async def _download_media(account_id, dialog_id, message_id, max_bytes) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id)
        settings = get_settings()
        cap = min(int(max_bytes) or settings.media_max_bytes, settings.media_max_bytes)
        msgs = await client.get_messages(entity, ids=[int(message_id)])
        if not msgs or msgs[0] is None or msgs[0].media is None:
            raise DomainError(404, f"#{message_id} не найден или без медиа")
        meta = media_meta(msgs[0])
        size = meta.get("size")
        if size and size > cap:
            raise DomainError(
                413, f"файл {size} байт > лимита {cap} (поднимите max_bytes ≤ {settings.media_max_bytes})"
            )
        data = await client.download_media(msgs[0], bytes)
        if data is None:
            raise DomainError(404, "медиа не скачалось (удалено из Telegram?)")
        if len(data) > cap:
            raise DomainError(413, f"скачано {len(data)} байт > лимита {cap}")
        name = meta.get("file_name") or ""
        if not name:
            mime = meta.get("mime_type", "")
            ext = {"audio/ogg": "ogg", "video/mp4": "mp4", "image/jpeg": "jpg"}.get(mime, "")
            name = f"media-{message_id}" + (f".{ext}" if ext else "")
        return {
            "dialog_id": dialog_id,
            "message_id": message_id,
            "file_name": name,
            "mime_type": meta.get("mime_type") or "application/octet-stream",
            "size": len(data),
            "is_voice": bool(meta.get("voice")),
            "is_round": bool(meta.get("round")),
            "duration": meta.get("duration"),
            "waveform": meta.get("waveform"),
            "data_base64": base64.b64encode(data).decode("ascii"),
        }


    @mcp.tool()
    async def upload_file(account_id: int | None = None, dialog_id: int = 0,
                          data_base64: str = "", file_path: str = "",
                          file_name: str = "file",
                          caption: str = "", force_document: bool = False) -> dict:
        """📎 Отправить файл/фото: data_base64 ИЛИ file_path (имя файла внутри
        TGCLIENT_UPLOAD_DIR на сервере — без путей). force_document=True —
        фото без превью, как документ."""
        return await _tool(
            lambda: _upload_file(account_id, dialog_id, data_base64, file_path, file_name, caption, force_document)
        )


    async def _upload_file(account_id, dialog_id, data_base64, file_path, file_name, caption, force_document) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        data = _upload_bytes(data_base64, file_path)
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        message = await call_guard(_manager(), row["id"],
                                   lambda: client.send_file(
                                       entity, data, file_name=file_name,
                                       caption=caption or None,
                                       force_document=force_document))
        return await _saved(client, row["id"], message)


    @mcp.tool()
    async def upload_voice(account_id: int | None = None, dialog_id: int = 0,
                           data_base64: str = "", file_path: str = "",
                           duration: int = 0,
                           waveform_samples: list[int] | None = None) -> dict:
        """🎤 Отправить голосовое: ogg-opus в base64 (data_base64) ИЛИ файлом с
        диска сервера (file_path — имя файла внутри TGCLIENT_UPLOAD_DIR,
        без путей). duration (сек) и waveform_samples (до 63 значений 0..100) —
        если не переданы, duration вычисляется из ogg-opus (для mp3/wav укажите
        duration сами), волна — плоская. Без duration и без парсинга — 422."""
        return await _tool(
            lambda: _upload_voice(account_id, dialog_id, data_base64, file_path, duration, waveform_samples)
        )


    async def _upload_voice(account_id, dialog_id, data_base64, file_path, duration, waveform_samples) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        data = _upload_bytes(data_base64, file_path)
        if duration <= 0:
            duration = ogg_opus_duration(data) or 0
        if duration <= 0:
            raise DomainError(422, "не удалось определить duration из ogg-opus — передайте duration (сек)")
        samples = list(waveform_samples or [])[:WAVEFORM_SAMPLES]
        if not samples:
            # плоская волна по длительности — ТГ отображает прямую полоску
            samples = [50] * WAVEFORM_SAMPLES
        wave_bytes = pack_waveform(samples)
        attributes = [DocumentAttributeAudio(duration=duration, voice=True, waveform=wave_bytes)]
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        message = await call_guard(_manager(), row["id"],
                                   lambda: client.send_file(
                                       entity, data, file_name="voice.ogg",
                                       attributes=attributes, mime_type="audio/ogg"))
        return await _saved(client, row["id"], message)


    @mcp.tool()
    async def upload_round(account_id: int | None = None, dialog_id: int = 0,
                           data_base64: str = "", file_path: str = "",
                           duration: int = 0, size: int = 0) -> dict:
        """⭕ Отправить кружок (video note): mp4 H.264/AAC, квадрат, ≤60 сек —
        base64 ИЛИ файлом с диска сервера (file_path внутри TGCLIENT_UPLOAD_DIR)
        + duration (сек) + размер квадрата (size, 240..640). Контейнер
        не тот — Telegram покажет обычное видео, предупреждайте пользователя."""
        return await _tool(
            lambda: _upload_round(account_id, dialog_id, data_base64, file_path, duration, size)
        )


    async def _upload_round(account_id, dialog_id, data_base64, file_path, duration, size) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        data = _upload_bytes(data_base64, file_path)
        if duration <= 0:
            raise DomainError(422, "duration (сек) обязателен для кружка")
        if duration > 60:
            raise DomainError(422, f"кружок ≤60 сек — получено {duration}")
        side = max(240, min(int(size) or 512, 640))
        attributes = [
            DocumentAttributeVideo(duration=duration, w=side, h=side, round_message=True),
        ]
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        message = await call_guard(_manager(), row["id"],
                                   lambda: client.send_file(
                                       entity, data, file_name="round.mp4",
                                       attributes=attributes, mime_type="video/mp4"))
        return await _saved(client, row["id"], message)


    # --- Звонки (сигналинг; см. app/tg/calls.py) ------------------------------

    @mcp.tool()
    async def call_start(account_id: int | None = None, dialog_id: int = 0,
                         video: bool = False) -> dict:
        """📞 Позвонить в Telegram: настоящий входящий вызов у абонента (ring).
        dialog_id — юзер (>0; каналам/группам звонить нельзя → 422). Аудио
        не передаётся (Telethon без tgcalls): после ответа абонента вызов
        корректно подтверждается и сразу завершается; отклонил → declined.
        Статус — call_status(call_id)."""
        return await _tool(lambda: _call_start(account_id, dialog_id, video))


    async def _call_start(account_id, dialog_id, video) -> dict:
        db = get_db()
        client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
        kind = type(entity).__name__  # InputPeerUser / InputPeerChat / InputPeerChannel
        if kind != "InputPeerUser":
            raise DomainError(422, "звонить можно только юзерам (dialog_id > 0 из dialogs_list)")
        if (limited := _write_guard(row["id"])) is not None:
            return limited
        from app.tg.calls import call_start

        out = await call_guard(_manager(), row["id"],
                               lambda: call_start(client, row["id"], entity, video=video))
        return out


    @mcp.tool()
    async def call_status(call_id: int) -> dict:
        """❓ Статус звонка по call_id: ringing → answered → hangup_after_answer,
        missed / declined / busy / cancelled. None после рестарта сервиса."""
        return await _tool(lambda: _call_status(call_id))


    async def _call_status(call_id) -> dict:
        from app.tg.calls import call_status

        out = await call_status(int(call_id))
        if out is None:
            raise DomainError(404, "звонок неизвестен (рестарт сервиса сбросил состояния; после таймаута Telegram сам завершает ringing)")
        return out


    @mcp.tool()
    async def call_discard(call_id: int) -> dict:
        """🚫 Отменить свой звонок до ответа абонента (hangup). Мутация."""
        return await _tool(lambda: _call_discard(call_id))


    async def _call_discard(call_id) -> dict:
        db = get_db()
        # call_id → состояние (RAM): аккаунт, к которому привязан звонок
        from app.tg.calls import _calls, call_discard

        st = _calls.get(int(call_id))
        if st is None:
            raise DomainError(404, "звонок неизвестен — после рестарта сервиса состояния в RAM теряются")
        account_id = st["account_id"]
        client, row = await _resolve_client(db, account_id, writable=True)
        await call_guard(_manager(), row["id"],
                         lambda: call_discard(client, int(call_id)))
        return {"status": "ok", "call_id": call_id, "state": "cancelled"}


    # --- Справочники ---------------------------------------------------------

    @mcp.tool()
    async def contacts_list(account_id: int | None = None, limit: int = 200) -> dict:
        """📇 Контакты аккаунта (id, имя, username)."""
        return await _tool(lambda: _contacts_list(account_id, limit))


    async def _contacts_list(account_id, limit) -> dict:
        client, _row = await _resolve_client(get_db(), account_id)
        result = await client(functions.contacts.GetContactsRequest(hash=0))
        users = list(getattr(result, "users", []) or [])
        return {"contacts": [entity_dict(u, short=True) for u in users[:_limit(limit, 200, 1000)]]}


    # --- Логин-тулы (текст-режим; owner-scoped, админ тоже только свои) ------

    @mcp.tool()
    async def account_login_start(phone: str, label: str = "") -> dict:
        """➕ Начать добавление Telegram-аккаунта: отправить код на phone
        (E.164 «+15551234567»). Дальше — код из ТГ (account_login_code),
        при 2FA — пароль (account_login_password). Лимит 3 параллельных логинов."""
        return await _tool(lambda: _account_login_start(phone, label))


    async def _account_login_start(phone: str, label: str) -> dict:
        # 5 стартов в час на mcp-юзера: агент, зацыклившийся на логине, не флудит
        if (limited := _login_guard()) is not None:
            return limited
        return dict(await login_start(phone.strip(), current_user_id(), label.strip()))


    @mcp.tool()
    async def account_login_code(login_id: str, code: str) -> dict:
        """🔢 Ввести код из Telegram (или от ошибки: вернётся attempts_left).
        step 'awaiting_password' — аккаунт под Cloud-паролем 2FA."""
        return await _tool(lambda: login_code(login_id.strip(), code.strip(), current_user_id()))


    @mcp.tool()
    async def account_login_password(login_id: str, password: str) -> dict:
        """🔒 Ввести Cloud-пароль 2FA (после step 'awaiting_password')."""
        return await _tool(lambda: login_password(login_id.strip(), password.strip(), current_user_id()))


    @mcp.tool()
    async def account_login_status(login_id: str) -> dict:
        """❓ Текущий шаг логина: awaiting_code / awaiting_password / done,
        expires_in_sec, last_error. phone_code_hash наружу не отдаётся."""
        return await _tool(lambda: login_status(login_id.strip(), current_user_id()))


    @mcp.tool()
    async def account_login_cancel(login_id: str) -> dict:
        """🚫 Отменить незавершённый логин (клиент disconnect, строка стирается)."""
        return await _tool(lambda: _account_login_cancel(login_id))


    async def _account_login_cancel(login_id: str) -> dict:
        await login_cancel(login_id.strip(), current_user_id())
        return {"status": "ok", "login_id": login_id, "step": "cancelled"}


    # --- Памятка агентам — prompt agent_guide в server.py ---------------------


_login_windows: dict[str, list[float]] = {}


def _login_guard() -> dict | None:
    """5 логинов в час на mcp-юзера (защита от зациклившегося агента)."""
    import time

    now = time.monotonic()
    win = _login_windows.setdefault(current_user_id(), [])
    while win and now - win[0] > 3600:
        win.pop(0)
    if len(win) >= 5:
        return {"error": 429, "detail": "login limit: 5 starts per hour — cancel existing or wait"}
    win.append(now)
    return None


def call_guard(manager, account_id: int, coro_fn):
    """Обёртка мутаций/чтений: длинный FloodWait → 429-данные, мёртвый ключ → 409+репорт."""
    from app.tg.manager import call_with_flood_guard

    return call_with_flood_guard(manager, account_id, coro_fn)


def _decode_payload(data_base64: str) -> bytes:
    """base64-payload тула → bytes; лимит TGCLIENT_MEDIA_MAX_BYTES → 413-данные."""
    if not data_base64:
        raise DomainError(422, "data_base64 пустой")
    try:
        data = base64_or_none(data_base64)
    except Exception as exc:  # noqa: BLE001 — invalid base64 → 422
        raise DomainError(422, f"data_base64 не декодируется: {exc}") from exc
    cap = get_settings().media_max_bytes
    if len(data) > cap:
        raise DomainError(413, f"payload {len(data)} байт > лимита {cap}")
    return data


def _read_local_file(file_path: str) -> bytes:
    """Аудио/файл с диска сервера для upload_* тулов (file_path вместо
    data_base64): ТОЛЬКО имя файла внутри TGCLIENT_UPLOAD_DIR — подкаталоги и
    пути запрещены, иначе персональный MCP-ключ мог бы прочитать произвольный
    файл хоста (.env с секретами). Не настроено — 503, нет файла — 404."""
    if not file_path or "/" in file_path or "\\" in file_path or ".." in file_path:
        raise DomainError(
            422, "file_path — только имя файла внутри TGCLIENT_UPLOAD_DIR (без путей)")
    base_dir = get_settings().upload_dir
    if not base_dir:
        raise DomainError(
            503, "отправка файлов с диска отключена — задайте TGCLIENT_UPLOAD_DIR на сервере")
    path = Path(base_dir) / file_path
    if not path.is_file():
        raise DomainError(404, f"файл не найден в upload-каталоге: {file_path}")
    data = path.read_bytes()
    cap = get_settings().media_max_bytes
    if len(data) > cap:
        raise DomainError(413, f"файл {len(data)} байт > лимита {cap}")
    return data


def _upload_bytes(data_base64: str, file_path: str) -> bytes:
    """Payload upload_* тулов: file_path (диск сервера) или data_base64;
    оба/ни один — 422."""
    if file_path and data_base64:
        raise DomainError(422, "укажите что-то одно: file_path или data_base64")
    if file_path:
        return _read_local_file(file_path)
    return _decode_payload(data_base64)