Newer
Older
tgclient-mcp / backend / app / tg / login_flow.py
"""Стейт-машина добавления Telegram-аккаунта — общий слой для SPA и MCP-тулов.

Пользователь может начать добавление в UI и продолжить агентом (или наоборот):
состояние одно — таблица login_sessions (переживает рестарт) + живой
TelegramClient незавершённого логина в RAM (AccountManager._pending).

Правила безопасности:
- phone нормализуется telethon.utils.parse_phone (E.164), не «как есть»;
- код и пароль 2FA ни в БД, ни в логи не попадают — сразу в sign_in;
- повторный start того же phone в TTL возвращает ту же login-сессию
  (второй send_code_request не дёргается — Telethon/Telegram флудят);
- 3 неверных кода / 2 неверных пароля — login-сессия отменяется;
- ошибки — DomainError(code, detail): SPA-роуты превращают в HTTPException,
  MCP-тулы — в данные {"error": ..., "detail": ...}.
"""

import uuid
from datetime import datetime, timedelta, timezone

from telethon import errors
from telethon.sessions import StringSession

from app.config import get_settings
from app.db import get_db
from app.errors import DomainError
from app.security import now_iso
from app.synapse_report import report
from app.tg.manager import AccountManager

LOGIN_TTL_MIN = 15
MAX_CODE_ATTEMPTS = 3
MAX_PASSWORD_ATTEMPTS = 2


def _now() -> datetime:
    return datetime.now(timezone.utc)


def _expires_iso() -> str:
    return (_now() + timedelta(minutes=LOGIN_TTL_MIN)).isoformat()


def _manager() -> AccountManager:
    from app.main import get_account_manager

    return get_account_manager()


def _phone_mask(phone: str) -> str:
    return phone[:3] + "***" + phone[-3:]


async def _require_creds() -> tuple[int, str]:
    """Креды приложения: у незавершённого логина без них делать нечего."""
    settings = get_settings()
    if not settings.api_id or not settings.api_hash:
        raise DomainError(503, "TGCLIENT_API_ID / TGCLIENT_API_HASH are not configured")
    return settings.api_id, settings.api_hash


async def login_start(phone: str, owner_user_id: str, label: str = "") -> dict:
    """Отправить код (или переиспользовать неистёкшую pending-сессию).

    Лимиты: ≤3 активных pending-логинов на владельца; ошибки — DomainError
    (422 — телефон, 409 — лимит, 429 — flood, 503 — нет API_ID).
    label — заметка из UI, доезжает на аккаунт при финализации.
    """
    from telethon.utils import parse_phone

    api_id, api_hash = await _require_creds()
    manager = _manager()
    db = get_db()

    try:
        phone = str(parse_phone(phone))
    except ValueError as exc:
        raise DomainError(422, "invalid phone number (E.164, e.g. +15551234567)") from exc

    cursor = await db.execute(
        "SELECT COUNT(*) AS n FROM login_sessions WHERE user_id = ?", (owner_user_id,)
    )
    if (await cursor.fetchone())["n"] >= 3:
        raise DomainError(409, "too many pending logins (max 3) — cancel or wait for expiry")

    # неистёкшая pending на тот же user+phone — вернуть без второго send_code
    cursor = await db.execute(
        "SELECT id, step, expires_at FROM login_sessions"
        " WHERE user_id = ? AND phone = ? AND expires_at > ?",
        (owner_user_id, phone, _now().isoformat()),
    )
    row = await cursor.fetchone()
    if row is not None and manager.pending_client(row["id"]) is not None:
        return {"login_id": row["id"], "step": row["step"], "expires_at": row["expires_at"],
                "reused": True}
    if row is not None:
        await db.execute("DELETE FROM login_sessions WHERE id = ?", (row["id"],))
        await db.commit()

    client = TelegramClientWithSession(api_id, api_hash)
    try:
        await client.connect()
        sent = await client.send_code_request(phone)
    except errors.PhoneNumberInvalidError as exc:
        await client.disconnect()
        raise DomainError(422, "phone number is invalid") from exc
    except errors.FloodWaitError as exc:
        await client.disconnect()
        raise DomainError(429, f"telegram flood wait: {exc.seconds}s — try later") from exc
    except errors.ApiIdInvalidError as exc:
        await client.disconnect()
        raise DomainError(503, "TGCLIENT_API_ID/API_HASH rejected by Telegram") from exc

    login_id = uuid.uuid4().hex
    expires = _expires_iso()
    await db.execute(
        "INSERT INTO login_sessions (id, user_id, phone, label, phone_code_hash, step,"
        " expires_at, created_at) VALUES (?, ?, ?, ?, ?, 'awaiting_code', ?, ?)",
        (login_id, owner_user_id, phone, label[:60], sent.phone_code_hash or "", expires, now_iso()),
    )
    await db.commit()
    manager.put_pending(login_id, client)
    return {
        "login_id": login_id,
        "phone_masked": _phone_mask(phone),
        "step": "awaiting_code",
        "expires_at": expires,
    }


async def _load_login(db, login_id: str, owner_user_id: str):
    cursor = await db.execute("SELECT * FROM login_sessions WHERE id = ?", (login_id,))
    row = await cursor.fetchone()
    if row is None:
        raise DomainError(404, "login session not found (expired or never existed)")
    if row["user_id"] != owner_user_id:
        # чужой login_id не раскрываем как существующий
        raise DomainError(404, "login session not found (expired or never existed)")
    if row["expires_at"] < _now().isoformat():
        await _cleanup(db, login_id)
        raise DomainError(410, "login session expired — start again")
    return row


async def _cleanup(db, login_id: str, *, with_client: bool = True) -> None:
    await db.execute("DELETE FROM login_sessions WHERE id = ?", (login_id,))
    await db.commit()
    await _manager().drop_pending(login_id, disconnect=with_client)


async def login_code(login_id: str, code: str, owner_user_id: str) -> dict:
    """Подтвердить код из Telegram. Возвращает текущий step логина."""
    manager = _manager()
    db = get_db()
    row = await _load_login(db, login_id, owner_user_id)
    if row["step"] != "awaiting_code":
        raise DomainError(409, f"login is at step '{row['step']}'")
    client = manager.pending_client(login_id)
    if client is None:
        # рестарт сервиса убил RAM-клиент: продолжить нечем
        await _cleanup(db, login_id, with_client=False)
        raise DomainError(503, "login client lost (service restart) — start again")

    phone = row["phone"]
    phone_code_hash = row["phone_code_hash"]
    try:
        await client.sign_in(phone=phone, code=code.strip(), phone_code_hash=phone_code_hash or None)
    except errors.PhoneCodeInvalidError:
        attempts = row["attempts"] + 1
        if attempts >= MAX_CODE_ATTEMPTS:
            await _cleanup(db, login_id)
            report("tg_login_failed", {"entity": f"login-{login_id}", "user_id": owner_user_id,
                                       "phone": _phone_mask(phone), "reason": "code_attempts"})
            raise DomainError(410, f"invalid code {attempts} times — login cancelled, start again")
        await db.execute(
            "UPDATE login_sessions SET attempts = ?, error = ? WHERE id = ?",
            (attempts, "invalid code", login_id),
        )
        await db.commit()
        return {"login_id": login_id, "step": "awaiting_code",
                "attempts_left": MAX_CODE_ATTEMPTS - attempts, "error": "invalid code"}
    except errors.PhoneCodeExpiredError:
        await _cleanup(db, login_id)
        raise DomainError(410, "code expired — start again (code is ~5-10 min valid)")
    except errors.SessionPasswordNeededError:
        await db.execute(
            "UPDATE login_sessions SET step = 'awaiting_password', error = '', attempts = 0,"
            " expires_at = ? WHERE id = ?",
            (_expires_iso(), login_id),
        )
        await db.commit()
        return {"login_id": login_id, "step": "awaiting_password"}
    except errors.FloodWaitError as exc:
        raise DomainError(429, f"telegram flood wait: {exc.seconds}s — try later") from exc
    except errors.PhoneNumberUnoccupiedError:
        # номер вообще не зарегистрирован в Telegram — на этом этапе маловероятно
        await _cleanup(db, login_id)
        raise DomainError(404, "this phone number is not registered on Telegram")

    await _finalize(db, manager, login_id, row, client)
    return {"login_id": login_id, "step": "done"}


async def login_password(login_id: str, password: str, owner_user_id: str) -> dict:
    manager = _manager()
    db = get_db()
    row = await _load_login(db, login_id, owner_user_id)
    if row["step"] != "awaiting_password":
        raise DomainError(409, "no 2FA password step expected — confirm code first")
    client = manager.pending_client(login_id)
    if client is None:
        await _cleanup(db, login_id, with_client=False)
        raise DomainError(503, "login client lost (service restart) — start again")
    try:
        await client.sign_in(password=password)
    except errors.PasswordHashInvalidError:
        attempts = row["attempts"] + 1
        if attempts >= MAX_PASSWORD_ATTEMPTS:
            await _cleanup(db, login_id)
            report("tg_login_failed", {"entity": f"login-{login_id}", "user_id": owner_user_id,
                                       "phone": _phone_mask(row["phone"]), "reason": "password_attempts"})
            raise DomainError(410, "invalid password twice — login cancelled, start again")
        await db.execute(
            "UPDATE login_sessions SET attempts = ?, error = 'invalid password' WHERE id = ?",
            (attempts, login_id),
        )
        await db.commit()
        return {"login_id": login_id, "step": "awaiting_password",
                "attempts_left": MAX_PASSWORD_ATTEMPTS - attempts, "error": "invalid password"}
    except errors.FloodWaitError as exc:
        raise DomainError(429, f"telegram flood wait: {exc.seconds}s — try later") from exc
    await _finalize(db, manager, login_id, row, client)
    return {"login_id": login_id, "step": "done"}


async def login_status(login_id: str, owner_user_id: str) -> dict:
    """Текущий шаг логина (SPA poll / MCP account_login_status)."""
    db = get_db()
    row = await _load_login(db, login_id, owner_user_id)
    expires_in = max(0, int((datetime.fromisoformat(row["expires_at"]) - _now()).total_seconds()))
    out = {
        "login_id": login_id,
        "phone_masked": _phone_mask(row["phone"]),
        "step": row["step"],
        "expires_in_sec": expires_in,
    }
    if row["error"]:
        out["last_error"] = row["error"]
    return out


async def login_cancel(login_id: str, owner_user_id: str) -> None:
    db = get_db()
    await _load_login(db, login_id, owner_user_id)  # владение + существование
    await _cleanup(db, login_id)


async def _finalize(db, manager: AccountManager, login_id: str, row, client) -> None:
    """Успешный логин: get_me → session_data → upsert accounts → cleanup."""
    me = await client.get_me()
    session_data = client.session.save()
    now = now_iso()
    tg_user_id = me.id
    username = getattr(me, "username", "") or ""
    display_name = " ".join(
        p for p in (getattr(me, "first_name", ""), getattr(me, "last_name", "")) if p
    ) or username or str(tg_user_id)
    cursor = await db.execute(
        "SELECT id FROM accounts WHERE user_id = ? AND phone = ?",
        (row["user_id"], row["phone"]),
    )
    existing = await cursor.fetchone()
    if existing is not None:
        await db.execute(
            "UPDATE accounts SET tg_user_id = ?, username = ?, display_name = ?, label = COALESCE(NULLIF(?, ''), label),"
            " session_data = ?, status = 'active', error = '', updated_at = ? WHERE id = ?",
            (tg_user_id, username, display_name, row["label"], session_data, now, existing["id"]),
        )
        account_id = existing["id"]
    else:
        cursor = await db.execute(
            "INSERT INTO accounts (user_id, phone, label, tg_user_id, username, display_name,"
            " session_data, status, created_at, updated_at)"
            " VALUES (?, ?, ?, ?, ?, ?, ?, 'active', ?, ?)",
            (row["user_id"], row["phone"], row["label"], tg_user_id, username, display_name,
             session_data, now, now),
        )
        account_id = cursor.lastrowid
    await db.commit()
    # pending-клиент свою работу сделал: disconnect и выкинуть из RAM;
    # аккаунт дальше ходит полноценным клиентом пула manager
    await manager.drop_pending(login_id)
    await manager.drop(account_id)
    # код уже был передан в sign_in (нигде не логируем); phone_code_hash
    # вместе со строкой login_sessions удаляется
    report("tg_account_logged_in", {
        "entity": f"account-{account_id}",
        "user_id": row["user_id"], "account_id": account_id,
        "phone_masked": _phone_mask(row["phone"]), "tg_user_id": tg_user_id,
    })


class TelegramClientWithSession:
    """TelegramClient логина: свежий StringSession, не кешируется в пуле.

    Отдельный класс нужен только читаемости; после finalize клиент
    disconnect-ится, аккаунт ходит полноценным клиентом пула.
    """

    def __init__(self, api_id: int, api_hash: str) -> None:
        from telethon import TelegramClient

        self._client = TelegramClient(
            StringSession(), api_id, api_hash,
            flood_sleep_threshold=60, request_retries=3,
        )

    async def connect(self):  # noqa: ANN201
        await self._client.connect()
        return self._client

    async def send_code_request(self, phone: str):  # noqa: ANN201
        return await self._client.send_code_request(phone)

    async def sign_in(self, *args, **kwargs):  # noqa: ANN002, ANN003
        return await self._client.sign_in(*args, **kwargs)

    async def get_me(self):  # noqa: ANN201
        return await self._client.get_me()

    async def disconnect(self) -> None:
        from contextlib import suppress

        with suppress(Exception):
            await self._client.disconnect()