"""Каталог 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 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_name: str = "file",
caption: str = "", force_document: bool = False) -> dict:
"""📎 Отправить файл/фото (base64). force_document=True — фото без превью,
как документ."""
return await _tool(
lambda: _upload_file(account_id, dialog_id, data_base64, file_name, caption, force_document)
)
async def _upload_file(account_id, dialog_id, data_base64, file_name, caption, force_document) -> dict:
db = get_db()
client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
data = _decode_payload(data_base64)
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 = "", duration: int = 0,
waveform_samples: list[int] | None = None) -> dict:
"""🎤 Отправить голосовое: ogg-opus в base64. duration (сек) и
waveform_samples (до 63 значений 0..100) — если не переданы,
вычисляются из файла (гранула последней Ogg-страницы / плоская волна).
Без duration и без парсинга — 422."""
return await _tool(
lambda: _upload_voice(account_id, dialog_id, data_base64, duration, waveform_samples)
)
async def _upload_voice(account_id, dialog_id, data_base64, duration, waveform_samples) -> dict:
db = get_db()
client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
data = _decode_payload(data_base64)
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 = "", duration: int = 0,
size: int = 0) -> dict:
"""⭕ Отправить кружок (video note): mp4 H.264/AAC, квадрат, ≤60 сек —
base64 + duration (сек) + размер квадрата (size, 240..640). Контейнер
не тот — Telegram покажет обычное видео, предупреждайте пользователя."""
return await _tool(
lambda: _upload_round(account_id, dialog_id, data_base64, duration, size)
)
async def _upload_round(account_id, dialog_id, data_base64, duration, size) -> dict:
db = get_db()
client, row, entity = await _client_for_dialog(db, account_id, dialog_id, writable=True)
data = _decode_payload(data_base64)
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)
# --- Справочники ---------------------------------------------------------
@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