Newer
Older
tgclient-mcp / backend / app / main.py
"""FastAPI-приложение tgclient: SPA + SSO + MCP-протокол + Telethon-пул.

Порядок роутов значим (FastAPI match-first): auth_routes → accounts →
mcp_tokens (личные + admin) → admin → /api/v1/health → /mcp и /mcp-protocol/
(APIRoute с гардом require_mcp) → StaticFiles /assets → SPA catch-all.

/mcp — тот же streamable-http ASGI (app/mcp/server.py) под APIRoute: гард
require_mcp выполняется ДО внутренностей MCP, личность (request.state.mcp_user)
мост переносит в ContextVar для тулов. Слэш в конце /mcp-protocol/ значим
(handbook mcp.md): клиенты ходят строго в него, /mcp — alias.
"""

import asyncio
from contextlib import asynccontextmanager, suppress
from pathlib import Path

import anyio
from fastapi import Depends, FastAPI, Request
from fastapi.exceptions import RequestValidationError
from fastapi.responses import FileResponse, JSONResponse
from starlette.responses import Response
from starlette.staticfiles import StaticFiles

from app import auth as auth_module
from app.auth import ensure_local_user_async
from app.db import close_db, get_db, init_db
from app.tg.manager import AccountManager

_manager: AccountManager | None = None


def get_account_manager() -> AccountManager:
    """Пул TelegramClient (инстант на процесс; им пользуются API-роуты и MCP-тулы)."""
    global _manager
    if _manager is None:
        _manager = AccountManager()
    return _manager


async def gc_loop() -> None:
    """Очистка истёкших auth-сессий/oauth_states/login_sessions (60 с)."""
    while True:
        try:
            await auth_module.purge_expired()
        except Exception as exc:  # noqa: BLE001 — не роняем цикл
            print(f"gc loop error: {exc}", flush=True)
        await asyncio.sleep(60)


@asynccontextmanager
async def lifespan(_application: FastAPI):
    await init_db()
    if not auth_module.auth_enabled():
        # auth-off (нет TGCLIENT_AUTH_*): всё принадлежит служебному local-юзеру
        await ensure_local_user_async()
        print("AUTH DISABLED — service-open mode (owner 'local')", flush=True)
    manager = get_account_manager()
    tasks = [
        asyncio.create_task(manager.persist_loop()),
        asyncio.create_task(gc_loop()),
        asyncio.create_task(manager.warmup()),
    ]
    try:
        # streamable-http session manager MCP-сервера живёт на lifespan
        from app.mcp.server import mcp

        async with anyio.create_task_group() as tg:
            async with mcp.session_manager.run():
                yield
            tg.cancel_scope.cancel()
    finally:
        for task in tasks:
            task.cancel()
            with suppress(asyncio.CancelledError):
                await task
        await manager.shutdown()
        await close_db()


_app = FastAPI(
    title="tgclient",
    version="0.1.0",
    description="MCP Telegram-клиент экосистемы Gnexus",
    lifespan=lifespan,
    docs_url=None,
    redoc_url=None,
    openapi_url=None,
)

# --- API --------------------------------------------------------------------

from app.api.accounts import router as accounts_router  # noqa: E402
from app.api.admin import router as admin_router  # noqa: E402
from app.api.auth_routes import me_router as me_router  # noqa: E402
from app.api.auth_routes import router as auth_router  # noqa: E402
from app.api.mcp_tokens import admin_router as tokens_admin_router  # noqa: E402
from app.api.mcp_tokens import router as tokens_router  # noqa: E402
from app.mcp.server import mcp_endpoint  # noqa: E402
from app.security import require_mcp  # noqa: E402

_app.include_router(auth_router)
_app.include_router(me_router)
_app.include_router(accounts_router)
_app.include_router(tokens_router)
_app.include_router(tokens_admin_router)
_app.include_router(admin_router)


@_app.get("/api/v1/health")
async def health() -> dict:
    """Контракт health хендбука: {status, version, accounts:{...}}, db:"ok"."""
    await get_db().execute("SELECT 1")
    return {
        "status": "ok",
        "version": _app.version,
        "accounts": await get_account_manager().count(),
        "db": "ok",
    }


# --- MCP: тот же mcp_app (streamable HTTP, stateless), гард до ASGI ----------

for _mcp_path, _methods in (
    ("/mcp", ["GET", "POST", "DELETE"]),
    ("/mcp-protocol/", ["GET", "POST", "DELETE"]),
):
    _app.add_api_route(
        _mcp_path, mcp_endpoint, methods=_methods,
        dependencies=[Depends(require_mcp)],
        include_in_schema=False,
    )

# --- SPA ----------------------------------------------------------------------

_STATIC_DIR = Path(__file__).resolve().parent.parent / "static"

if (_STATIC_DIR / "assets").is_dir():
    _app.mount("/assets", StaticFiles(directory=_STATIC_DIR / "assets"), name="assets")


@_app.get("/{path:path}", include_in_schema=False, name="spa")
async def spa(path: str) -> Response:
    """SPA catch-all (no-store — index свежий после каждого релиза);
    неизвестные /api/* возвращают 404 данными, не index; без собранного
    frontend-дистрибутива (dev-бэкенд) — честный 404, не 500."""
    if path.startswith(("api/", "auth/", "webhooks/")):
        return JSONResponse(status_code=404, content={"error": 404, "detail": "not found"})
    candidate = _STATIC_DIR / path
    if (
        path
        and candidate.is_file()
        and candidate.resolve().is_relative_to(_STATIC_DIR.resolve())
    ):
        return FileResponse(candidate)
    if not (_STATIC_DIR / "index.html").is_file():
        return JSONResponse(status_code=404, content={"error": 404, "detail": "no SPA build"})
    return FileResponse(
        _STATIC_DIR / "index.html", headers={"Cache-Control": "no-store"}
    )


@_app.exception_handler(RequestValidationError)
async def validation_error(_request: Request, exc: RequestValidationError) -> JSONResponse:
    return JSONResponse(
        status_code=422, content={"error": 422, "detail": str(exc.errors()[:1])}
    )


app = _app