Newer
Older
tgclient-mcp / backend / app / api / accounts.py
"""Аккаунты Telegram: список, добавление (общий login flow — SPA и MCP-тулы
работают с одним состоянием), логаут, ручной connect/disconnect.

Личность — cookie-сессия (require_session): Bearer-админ-токен личности не
имеет и акки не выдаёт; владелец строки accounts сверяется на каждом шаге.
"""

import contextlib

from fastapi import APIRouter, HTTPException, Request
from pydantic import BaseModel, Field

from app.config import get_settings
from app.db import get_db
from app.errors import DomainError
from app.security import require_session
from app.tg.login_flow import (
    login_cancel,
    login_code,
    login_password,
    login_start,
    login_status,
)

router = APIRouter(prefix="/api/v1")


def get_manager():
    from app.main import get_account_manager

    return get_account_manager()


def _account_out(row, connected: bool) -> dict:
    return {
        "id": row["id"],
        "label": row["label"],
        "phone": row["phone"],
        "tg_user_id": row["tg_user_id"],
        "username": row["username"],
        "display_name": row["display_name"],
        "status": row["status"],
        "error": row["error"],
        "connected": connected,
        "created_at": row["created_at"],
        "last_used_at": row["last_used_at"],
    }


async def _own_account(db, account_id: int, user_id: str):
    cursor = await db.execute("SELECT * FROM accounts WHERE id = ?", (account_id,))
    row = await cursor.fetchone()
    if row is None or row["user_id"] != user_id:
        raise HTTPException(status_code=404, detail="account not found")
    return row


class LoginStart(BaseModel):
    phone: str = Field(min_length=4, max_length=32)
    label: str = Field(default="", max_length=60)


class LoginCode(BaseModel):
    code: str = Field(min_length=1, max_length=16)


class LoginPassword(BaseModel):
    password: str = Field(min_length=1, max_length=256)


@router.get("/accounts")
async def list_accounts(request: Request) -> dict:
    """Свои аккаунты (+живость клиента из пула); гейт карточек UI."""
    session = await require_session(request)
    db = get_db()
    manager = get_manager()
    cursor = await db.execute(
        "SELECT * FROM accounts WHERE user_id = ? ORDER BY id", (session["user_id"],)
    )
    rows = await cursor.fetchall()
    return {
        "accounts": [_account_out(r, r["id"] in manager.clients) for r in rows],
        "max_accounts": get_settings().max_active_accounts,
    }


@router.post("/accounts/logins")
async def start_login(request: Request, payload: LoginStart) -> dict:
    """Отправить код Telegram: login_sessions строка в шаге awaiting_code."""
    session = await require_session(request)
    db = get_db()
    cursor = await db.execute(
        "SELECT COUNT(*) AS n FROM accounts WHERE user_id = ?", (session["user_id"],)
    )
    n = (await cursor.fetchone())["n"]
    if n >= get_settings().max_active_accounts:
        raise HTTPException(
            status_code=409,
            detail=f"too many accounts ({n}/{get_settings().max_active_accounts})",
        )
    try:
        return await login_start(payload.phone, session["user_id"], payload.label.strip())
    except DomainError as exc:
        raise HTTPException(status_code=exc.code, detail=exc.detail) from exc


@router.get("/accounts/logins/{login_id}")
async def login_state(login_id: str, request: Request) -> dict:
    """Poll текущего шага (SPA LoginWizard; MCP account_login_status)."""
    session = await require_session(request)
    try:
        return await login_status(login_id, session["user_id"])
    except DomainError as exc:
        raise HTTPException(status_code=exc.code, detail=exc.detail) from exc


@router.post("/accounts/logins/{login_id}/code")
async def submit_code(login_id: str, request: Request, payload: LoginCode) -> dict:
    session = await require_session(request)
    try:
        return await login_code(login_id, payload.code, session["user_id"])
    except DomainError as exc:
        raise HTTPException(status_code=exc.code, detail=exc.detail) from exc


@router.post("/accounts/logins/{login_id}/password")
async def submit_password(login_id: str, request: Request, payload: LoginPassword) -> dict:
    session = await require_session(request)
    try:
        return await login_password(login_id, payload.password, session["user_id"])
    except DomainError as exc:
        raise HTTPException(status_code=exc.code, detail=exc.detail) from exc


@router.delete("/accounts/logins/{login_id}")
async def cancel_login(login_id: str, request: Request) -> dict:
    session = await require_session(request)
    try:
        await login_cancel(login_id, session["user_id"])
    except DomainError as exc:
        raise HTTPException(status_code=exc.code, detail=exc.detail) from exc
    return {"status": "ok"}


@router.post("/accounts/{account_id}/connect")
async def connect_account(account_id: int, request: Request) -> dict:
    """Ручная попытка поднять клиент (warmup делает это же фоном)."""
    session = await require_session(request)
    db = get_db()
    row = await _own_account(db, account_id, session["user_id"])
    if row["status"] != "active":
        raise HTTPException(status_code=409, detail=f"account is {row['status']} — login again")
    manager = get_manager()
    await manager.drop(account_id)
    try:
        await manager.get_client(account_id)
    except DomainError as exc:
        raise HTTPException(status_code=exc.code, detail=exc.detail) from exc
    return {"status": "ok", "connected": True}


@router.post("/accounts/{account_id}/disconnect")
async def disconnect_account(account_id: int, request: Request) -> dict:
    """Снять клиент из пула (сессия сохранена, аккаунт остаётся активным)."""
    session = await require_session(request)
    db = get_db()
    await _own_account(db, account_id, session["user_id"])
    await get_manager().drop(account_id)
    return {"status": "ok", "connected": False}


@router.delete("/accounts/{account_id}")
async def logout_account(account_id: int, request: Request) -> dict:
    """Логаут Telegram: client.log_out() отзывает MTProto-ключ, строка
    аккаунта удаляется. Необратимо — UI подтверждает перед вызовом."""
    session = await require_session(request)
    db = get_db()
    row = await _own_account(db, account_id, session["user_id"])
    manager = get_manager()
    client = manager.clients.get(account_id)
    if client is not None:
        with contextlib.suppress(Exception):
            await client.log_out()
    elif row["session_data"]:
        # клиент мёртв, сессия осталась — короткое подключение только для log_out
        settings = get_settings()
        if settings.api_id and settings.api_hash:
            with contextlib.suppress(Exception):
                from telethon import TelegramClient
                from telethon.sessions import StringSession

                tclient = TelegramClient(StringSession(row["session_data"]),
                                         settings.api_id, settings.api_hash)
                await tclient.connect()
                if await tclient.is_user_authorized():
                    await tclient.log_out()
                await tclient.disconnect()
    await manager.drop(account_id)
    await db.execute("DELETE FROM accounts WHERE id = ?", (account_id,))
    await db.commit()
    from app.synapse_report import report

    report("tg_account_logged_out", {
        "entity": f"account-{account_id}",
        "user_id": session["user_id"],
        "account_id": account_id,
        "phone_masked": row["phone"][:3] + "***" + row["phone"][-3:],
    })
    return {"status": "ok"}