diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..de1dde3 --- /dev/null +++ b/.env.example @@ -0,0 +1,40 @@ +# tgclient-mcp — настройка. Скопировать в .env рядом с сервисом, chmod 600, +# значения секретов продублировать в gnexus-creds (канон handbook secrets.md). +# Секреты в git попадать не должны — файл в .gitignore. + +# Порт HTTP-сервера (в compose уже проброшен; для голого uvicorn не используется) +TGCLIENT_PORT=8710 +# Путь к SQLite-базе (в контейнере — том /data) +TGCLIENT_DB_PATH=/data/tgclient.db + +# Статический супер-токен MCP-бэкдора (handbook mcp.md): `openssl rand -hex 32`. +# Пусто — работают только персональные ключи mcp_* (или cookie-сессии). +TGCLIENT_ADMIN_TOKEN= + +# Креды приложения Telegram (one set на деплой — api_id идентифицирует +# приложение, не аккаунт; получить на my.telegram.org) +TGCLIENT_API_ID= +TGCLIENT_API_HASH= + +# Лимиты +TGCLIENT_MAX_ACTIVE_ACCOUNTS=30 # максимум аккаунтов на пользователя +TGCLIENT_MEDIA_MAX_BYTES=20971520 # cap на download/upload через MCP, байт (20 МБ) +TGCLIENT_WRITE_LIMIT_N=20 # мутации через MCP: столько +TGCLIENT_WRITE_LIMIT_WINDOW=60 # ... за столько секунд, на аккаунт + +# gnexus-auth (SSO): задан client_id — браузерная авторизация через SSO; +# пусто — сервис открыт, все акки уходят локальному служебному юзеру +TGCLIENT_AUTH_BASE_URL= +TGCLIENT_AUTH_CLIENT_ID= +TGCLIENT_AUTH_CLIENT_SECRET= +TGCLIENT_AUTH_REDIRECT_URI= +TGCLIENT_AUTH_WEBHOOK_SECRET= +# пусто = любой залогиненный; через запятую — e-mail или user_id +TGCLIENT_AUTH_ALLOWLIST= +TGCLIENT_SESSION_TTL=604800 # TTL cookie-сессии, сек (7 суток) + +# Gnexus Synapse (исходящие уведомления): пустой api_key = интеграция выключена +TGCLIENT_SYNAPSE_URL= +TGCLIENT_SYNAPSE_API_KEY= +TGCLIENT_SYNAPSE_DEFAULT_SOURCE=tgclient-mcp +TGCLIENT_SYNAPSE_TIMEOUT=10 \ No newline at end of file diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..3bb3f8b --- /dev/null +++ b/.gitignore @@ -0,0 +1,17 @@ +# секреты и данные +.env +data/ +*.db +*.db-journal +*.db-wal +*.db-shm + +__pycache__/ +*.pyc +.venv/ +*.egg-info/ + +node_modules/ +frontend/dist/ +backend/static/ +.vite/ \ No newline at end of file diff --git a/README.md b/README.md new file mode 100644 index 0000000..23ba773 --- /dev/null +++ b/README.md @@ -0,0 +1,87 @@ +# tgclient-mcp + +Telegram-клиент формата MCP для экосистемы Gnexus (конвенции +[gnexus-handbook](https://git.gnexus.space/git/root/gnexus-handbook)): +мульти-аккаунтный MTProto юзер-клиент (Telethon), MCP-сервер для ИИ-агентов +(`/mcp-protocol/`, alias `/mcp`) и SPA-админка, закрытая SSO gnexus-auth, +с персональными MCP-токенами. + +## Возможности + +- **Мульти-аккаунтность**: каждый SSO-пользователь логинит свои Telegram-аккаунты + (phone → код → 2FA в SPA-визарде или текстом через MCP-тулы) и работает только со + своими; админ видит все. +- **MCP-каталог тулов** (`/mcp-protocol/`, streamable HTTP + Bearer `mcp_*`): + - ядро: `accounts_list`, `me_get`, `dialogs_list`, `messages_history`, + `messages_read`, `messages_search`, `message_send/reply/edit/delete`; + - медиа: `download_media` (base64), `upload_file`; **голосовые** (`upload_voice` + ogg-opus + waveform 63×5 бит) и **кружки** (`upload_round`, mp4-квадрат ≤60 с); + - справочники: `contacts_list`, `chat_info`, `chat_participants`; + - текстовый логин: `account_login_start/code/password/status/cancel`. +- **Ошибки — данные** (`{"error": code, "detail": ...}`), перс. ключи по канону + `mcp.md` (≤10 на юзера, снейпшот роли, plaintext один раз, ревок вместо архива), + блокировка юзера гасит ключи тем же 401, rate-limit мутаций и логинов, + Synapse-уведомления (логин/логаут/auth-lost/флуд). + +## Стек + +FastAPI + uvicorn, aiosqlite (WAL, одна БД SQLite в `/data`), Telethon +(только `StringSession`), `mcp` python-sdk (FastMCP), gnexus-gauth SDK (OAuth +PKCE + cookie-сессии + вебхуки), gnexus-synapse (нотификации), Vue 3 + Vite + +gnexus-ui-kit (SPA). + +## Запуск + +```bash +cp .env.example .env && chmod 600 .env +# заполнить: TGCLIENT_API_ID/_API_HASH (my.telegram.org), TGCLIENT_ADMIN_TOKEN, +# TGCLIENT_AUTH_* ( gnexus-auth приложение), TGCLIENT_SYNAPSE_* (опционально) +docker compose up -d --build +# http://:TGCLIENT_PORT — SPA; MCP: http://:TGCLIENT_PORT/mcp-protocol/ +``` + +### Подключение агента + +``` +claude mcp add --transport http tgclient https:///mcp-protocol/ \ + --header "Authorization: Bearer mcp_<персональный ключ>" +``` + +Ключ выпускается на странице «MCP-ключи» сервиса; показывается один раз. +Статический `TGCLIENT_ADMIN_TOKEN` — супер-бэкдор для скриптов (роль superadmin). + +## Разработка (без docker) + +```bash +cd backend && python3 -m venv .venv && .venv/bin/pip install -e . +TGCLIENT_DB_PATH=./data/tgclient.db .venv/bin/uvicorn app.main:app --reload +# SPA: cd frontend && npm ci && npm run dev (vite proxy /api,/auth → 8710) +``` + +- `TGCLIENT_AUTH_CLIENT_ID` пуст → **auth-off**: всё принадлежит служебному + юзеру `local`, MCP открыт без ключей (только для разработки!). +- `/api/v1/health` → `{status, version, accounts:{active,connected,pending_logins}, db}`. + +## Структура + +``` +backend/app/ + config.py db.py schema.sql security.py auth.py locales.py synapse_report.py errors.py + main.py # lifespan (telethon-пул, GC, persist), роуты, /mcp, SPA + api/ auth_routes accounts mcp_tokens admin + tg/ manager login_flow limits waveform + mcp/ server context tools serializers +frontend/src/ App.vue router.js api.js i18n/ pages/ +``` + +## Pitfalls (из плана) + +1. Telethon-сессии — только `StringSession` в БД (SQLiteSession конфликтует с aiosqlite). +2. TL-объекты сериализуются только через `mcp/serializers.py`, никогда `to_dict()`; + наружу — `id` в bot-API-маркировке (юзер >0, группа `-id`, канал `-100id`). +3. Кружок требует mp4 H.264/AAC квадрат ≤60 с, иначе уйдёт как обычное видео; + waveform голосового считаем сами (63 семпла × 5 бит). +4. base64 в JSON ≈ +33%: cap `TGCLIENT_MEDIA_MAX_BYTES` (по умолчанию 20 MB). +5. Повторный `send_code_request` на живую pending-сессию не дёргается; код/пароль + не логируются и не хранятся. +6. FloodWait ≤60 с Telethon спит сам; длинный — 429 как данные + Synapse-репорт. \ No newline at end of file diff --git a/backend/Dockerfile b/backend/Dockerfile new file mode 100644 index 0000000..b49a7e4 --- /dev/null +++ b/backend/Dockerfile @@ -0,0 +1,29 @@ +# tgclient — multi-stage: сборка SPA (node) → backend (python + static) + +# --- Stage 1: SPA (frontend/dist) --------------------------------------------- +FROM node:22-slim AS frontend +RUN apt-get update \ + && apt-get install -y --no-install-recommends git ca-certificates \ + && rm -rf /var/lib/apt/lists/* +WORKDIR /build +COPY frontend/package.json frontend/.npmrc ./ +RUN npm ci +COPY frontend/ ./ +RUN npm run build + +# --- Stage 2: сервис ------------------------------------------------------------ +FROM python:3.12-slim +RUN apt-get update \ + && apt-get install -y --no-install-recommends git ca-certificates \ + && rm -rf /var/lib/apt/lists/* +WORKDIR /srv +# git-зависимости (gnexus-gauth, gnexus-synapse) качаются на этапе сборки +COPY backend/pyproject.toml backend/README.md* ./ +COPY backend/app ./app +RUN pip install --no-cache-dir . +# SPA из stage 1 — обслуживается catch-all (app/main.py) +COPY --from=frontend /build/dist ./static +EXPOSE 8000 +HEALTHCHECK --interval=30s --timeout=5s --start-period=20s \ + CMD python -c "import os,urllib.request;urllib.request.urlopen('http://127.0.0.1:8000/api/v1/health',timeout=4)" || exit 1 +CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] \ No newline at end of file diff --git a/backend/app/__init__.py b/backend/app/__init__.py new file mode 100644 index 0000000..e69de29 --- /dev/null +++ b/backend/app/__init__.py diff --git a/backend/app/api/__init__.py b/backend/app/api/__init__.py new file mode 100644 index 0000000..e69de29 --- /dev/null +++ b/backend/app/api/__init__.py diff --git a/backend/app/api/accounts.py b/backend/app/api/accounts.py new file mode 100644 index 0000000..9d4b98c --- /dev/null +++ b/backend/app/api/accounts.py @@ -0,0 +1,209 @@ +"""Аккаунты 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"} \ No newline at end of file diff --git a/backend/app/api/admin.py b/backend/app/api/admin.py new file mode 100644 index 0000000..31d0f32 --- /dev/null +++ b/backend/app/api/admin.py @@ -0,0 +1,63 @@ +"""Админ API: обзор всех пользователей/аккаунтов/ключей + принудительный +disconnect аккаунта. Только role admin/superadmin (require_admin) — +мульти-юзер сервис: обычные пользователи ничего админского не видят. +""" + +from fastapi import APIRouter, Depends, HTTPException + +from app.config import get_settings +from app.db import get_db +from app.security import now_iso, require_admin + +router = APIRouter(prefix="/api/v1/admin", dependencies=[Depends(require_admin)]) + + +@router.get("/overview") +async def overview() -> dict: + """Сводка админки: пользователи, аккаунты, токены, счётчики.""" + db = get_db() + + users = await db.execute_fetchall( + "SELECT user_id, email, display_name, system_role, blocked, created_at FROM users" + " WHERE user_id != 'local' ORDER BY created_at DESC" + ) + accounts = await db.execute_fetchall( + """SELECT a.id, a.user_id, a.phone, a.label, a.status, a.error, a.tg_user_id, + a.username, a.display_name, a.created_at, a.last_used_at, u.email AS owner_email + FROM accounts a LEFT JOIN users u ON u.user_id = a.user_id ORDER BY a.id DESC""" + ) + tokens = await db.execute_fetchall( + "SELECT id, user_id, name, user_email, system_role, token_hint, created_at," + " last_used_at, revoked_at FROM mcp_tokens ORDER BY id DESC" + ) + cursor = await db.execute( + "SELECT COUNT(*) AS n FROM login_sessions" + ) + pending_logins = (await cursor.fetchone())["n"] + return { + "users": [dict(u) for u in users], + "accounts": [dict(a) for a in accounts], + "tokens": [dict(t) for t in tokens], + "counters": {"pending_logins": pending_logins, + "max_active_accounts": get_settings().max_active_accounts}, + } + + +@router.post("/accounts/{account_id}/revoke_session") +async def revoke_account_session(account_id: int) -> dict: + """Принудительно снять клиент (сессия остаётся в БД — админ не логаутит + пользователя Telegram; это только управленческий disconnect пула).""" + db = get_db() + cursor = await db.execute("SELECT id FROM accounts WHERE id = ?", (account_id,)) + if (await cursor.fetchone()) is None: + raise HTTPException(status_code=404, detail="account not found") + from app.main import get_account_manager + + manager = get_account_manager() + await manager.drop(account_id) + await db.execute( + "UPDATE accounts SET updated_at = ?, last_used_at = COALESCE(last_used_at, ?) WHERE id = ?", + (now_iso(), now_iso(), account_id), + ) + await db.commit() + return {"status": "ok", "connected": False} \ No newline at end of file diff --git a/backend/app/api/auth_routes.py b/backend/app/api/auth_routes.py new file mode 100644 index 0000000..d236502 --- /dev/null +++ b/backend/app/api/auth_routes.py @@ -0,0 +1,243 @@ +"""OAuth/session-роуты gnexus-auth — каркас скопирован с hard-panel +(паттерн gnexus-creds). + +GET /auth/login?return_to=/… — 302 на gnexus-auth (PKCE, state в oauth_states) +GET /auth/callback?code&state — обмен кода, сессия в cookie, redirect +POST /auth/logout — удалить сессию +GET /api/v1/me — профиль из сессии (гейт SPA и identity в дровере) +PATCH /api/v1/me — override языка интерфейса (только cookie-сессия) +POST /webhooks/gnexus-auth(+) — HMAC-верифицированные события gnexus-auth +""" + +import asyncio +import json +from datetime import datetime, timedelta, timezone + +from fastapi import APIRouter, Request, Response +from fastapi.responses import JSONResponse, RedirectResponse +from gnexus_gauth.exceptions import ( + StateValidationException, + TokenExchangeException, + TransportException, +) +from gnexus_gauth.oauth import AuthorizationUrlBuilder, PkceGenerator +from pydantic import BaseModel + +from app.auth import ( + SCOPES, + SESSION_COOKIE, + allowlisted, + auth_enabled, + create_session, + delete_session, + exchange_and_fetch, + get_ctx, + safe_return_to, +) +from app.config import get_settings +from app.db import get_db +from app.security import auth_mode_error + +router = APIRouter() +webhook_router = APIRouter() + + +def _now() -> datetime: + return datetime.now(timezone.utc) + + +def _now_iso() -> str: + return _now().isoformat() + + +@router.get("/auth/login") +async def login(request: Request, return_to: str = "/") -> RedirectResponse: + """Начало OAuth: state + PKCE → oauth_states, 302 на gnexus-auth.""" + settings = get_settings() + return_to = safe_return_to(return_to) + ctx = get_ctx() + state = PkceGenerator.generate_state() + verifier = PkceGenerator.generate_verifier() + challenge = PkceGenerator.generate_challenge(verifier) + db = get_db() + await db.execute( + "INSERT INTO oauth_states (state, pkce_verifier, return_to, expires_at) VALUES (?, ?, ?, ?)", + (state, verifier, return_to, + (_now() + timedelta(seconds=settings.session_ttl_seconds)).isoformat()), + ) + await db.commit() + url = AuthorizationUrlBuilder(ctx["config"]).build( + state=state, pkce_challenge=challenge, return_to=return_to, scopes=SCOPES, + ) + return RedirectResponse(url, status_code=302) + + +@router.get("/auth/callback") +async def callback(request: Request, code: str = "", state: str = ""): + """Обмен кода на токен, сессия в cookie, редирект на return_to.""" + if not state or not code: + raise auth_mode_error() + settings = get_settings() + db = get_db() + cursor = await db.execute( + "SELECT pkce_verifier, return_to, expires_at FROM oauth_states WHERE state = ?", + (state,), + ) + saved = await cursor.fetchone() + if saved is None or saved["expires_at"] < _now_iso(): + return JSONResponse( + {"detail": "Invalid or expired OAuth state."}, status_code=400 + ) + try: + _, user = await asyncio.to_thread(exchange_and_fetch, code, saved["pkce_verifier"]) + except (TokenExchangeException, StateValidationException, TransportException): + return JSONResponse( + {"detail": "Authorization server is unreachable or rejected the login."}, + status_code=502, + ) + if not allowlisted(user.email, user.user_id): + return JSONResponse({"detail": "User is not allowed here."}, status_code=403) + if user.status in {"disabled", "blocked", "deleted"}: + return JSONResponse({"detail": "User is disabled."}, status_code=403) + session_id = await create_session(user) + # Успешный логин — пользователь жив: снимаем флаг blocked вебхука + await db.execute("UPDATE users SET blocked = 0 WHERE user_id = ?", (user.user_id,)) + await db.commit() + await db.execute("DELETE FROM oauth_states WHERE state = ?", (state,)) + await db.commit() + response = RedirectResponse(safe_return_to(saved["return_to"]), status_code=302) + response.set_cookie( + SESSION_COOKIE, + session_id, + httponly=True, + samesite="lax", + max_age=settings.session_ttl_seconds, + secure=bool(settings.auth_redirect_uri.startswith("https://")), + ) + return response + + +@router.post("/auth/logout") +async def logout(request: Request) -> Response: + """Удалить свою сессию (cookie гасится); SSO-логаут MCP-ключи не трогает.""" + settings = get_settings() + await delete_session(request.cookies.get(SESSION_COOKIE)) + response = JSONResponse({"status": "ok"}) + response.delete_cookie( + SESSION_COOKIE, + httponly=True, + samesite="lax", + secure=bool(settings.auth_redirect_uri.startswith("https://")), + ) + return response + + +# --- /me для фронтенда --------------------------------------------------------- + +me_router = APIRouter(prefix="/api/v1") + + +class LocaleUpdate(BaseModel): + locale: str = "" + + +def _user_payload(session: dict, row) -> dict: + profile = json.loads(row["profile"]) if row and row["profile"] else {} + return { + "user_id": session["user_id"], + "email": session["email"], + "display_name": session["display_name"], + "avatar_url": session["avatar_url"], + "role": (row["system_role"] if row and row["system_role"] else None) or "user", + }, profile + + +@me_router.get("/me") +async def me(request: Request) -> dict: + """Профиль текущего пользователя + флаг режима авторизации. + + Auth выключен → {auth_enabled: false} — SPA молча открывает сервис. + Включён и сессии нет → 401 (SPA показывает экран входа). + """ + if not auth_enabled(): + return {"auth_enabled": False, "user": None, "locale_effective": None} + from app.auth import get_session, get_user_row + from app.locales import effective_locale + + session = await get_session(request.cookies.get(SESSION_COOKIE)) + if session is None: + raise auth_mode_error() + row = await get_user_row(session["user_id"]) + profile = json.loads(row["profile"]) if row and row["profile"] else {} + override = row["locale"] if row else None + role = (row["system_role"] if row and row["system_role"] else None) or "user" + return { + "auth_enabled": True, + "user": { + "user_id": session["user_id"], + "email": session["email"], + "display_name": session["display_name"], + "avatar_url": session["avatar_url"], + "role": role, + }, + "locale": override, + "locale_effective": effective_locale(override, profile), + } + + +@me_router.patch("/me") +async def update_me(request: Request, payload: LocaleUpdate) -> dict: + """Override языка (настройки → «Язык»); 'auto'/пусто сбрасывают. + + Только cookie-сессия: у Bearer-токена нет личности для настроек (401). + """ + if not auth_enabled(): + raise auth_mode_error() + from app.auth import get_session, get_user_row, set_user_locale + from app.locales import effective_locale, normalize_locale + + session = await get_session(request.cookies.get(SESSION_COOKIE)) + if session is None: + raise auth_mode_error() + + value = (payload.locale or "").strip().lower() + override = None if value in ("", "auto") else normalize_locale(value) + if value not in ("", "auto") and override is None: + return JSONResponse({"detail": "unsupported locale"}, status_code=422) + + await set_user_locale(session["user_id"], override, email=session["email"], + display_name=session["display_name"]) + row = await get_user_row(session["user_id"]) + profile = json.loads(row["profile"]) if row and row["profile"] else {} + return { + "auth_enabled": True, + "user": dict(_user_payload(session, row)[0]), + "locale": row["locale"] if row else None, + "locale_effective": effective_locale(row["locale"] if row else None, profile), + } + + +# --- Webhook gnexus-auth (HMAC) ------------------------------------------------ + +@webhook_router.post("/webhooks/gnexus-auth") +@webhook_router.post("/webhooks/gnexus-auth/") +async def gnexus_auth_webhook(request: Request) -> dict: + from app.auth import apply_webhook_update, get_ctx + + settings = get_settings() + raw = (await request.body()).decode("utf-8", errors="replace") + headers = dict(request.headers) + ctx = get_ctx() + try: + ctx["webhook"].verify(raw, headers, settings.auth_webhook_secret) + event = ctx["parser"].parse(raw) + except Exception: # noqa: BLE001 — чужое/несанкционированное всегда 401 + return JSONResponse({"detail": "invalid webhook signature"}, status_code=401) + subject = ( + event.target_identifiers.get("sub") + or event.target_identifiers.get("user_id") + or event.metadata.get("sub") + ) + if subject: + await apply_webhook_update(str(subject), event.event_type, event.metadata) + return {"status": "ok"} \ No newline at end of file diff --git a/backend/app/api/mcp_tokens.py b/backend/app/api/mcp_tokens.py new file mode 100644 index 0000000..257da33 --- /dev/null +++ b/backend/app/api/mcp_tokens.py @@ -0,0 +1,174 @@ +"""Персональные MCP-ключи (канон handbook 10-platform/mcp.md, референс +gnexus-synapse): выдача, списки, ревок. + +Выпуск — только из cookie-сессии: у Bearer-админ-токена нет личности, которой +можно выдать ключ, а ключ — персональные креды. Персональный mcp_* на /mcp +принимает гард require_mcp (app/security.py), статический +TGCLIENT_ADMIN_TOKEN остаётся супер-токеном. Plaintext существует ровно один +раз — в ответе 201 создания; в БД только sha256-хэш и хвост-хинт. +""" + +import time + +from pydantic import BaseModel, Field + +from fastapi import APIRouter, Depends, HTTPException, Request + +from app.auth import SESSION_COOKIE, get_session, get_user_row # noqa: F401 (SESSION_COOKIE: /me-паттерн) +from app.db import get_db +from app.security import ( + generate_mcp_token, + hash_key, + mcp_token_hint, + now_iso, + require_admin, + require_session, +) + +ACTIVE_LIMIT = 10 + +router = APIRouter(prefix="/api/v1") + +# Простейший rate-limit на выдачу (канон mcp.md: чувствительные операции): +# окно в памяти процесса — после рестарта окно чистое, для защиты от спама хватает. +_CREATE_WINDOW = 600.0 +_CREATE_MAX = 5 +_create_log: dict[str, list[float]] = {} + + +class TokenCreate(BaseModel): + name: str = Field(default="", max_length=120) + + +def _row_out(row) -> dict: + return { + "id": row["id"], + "name": row["name"], + "token_hint": row["token_hint"], + "system_role": row["system_role"], + "user_email": row["user_email"], + "created_at": row["created_at"], + "last_used_at": row["last_used_at"], + "revoked_at": row["revoked_at"], + } + + +async def _issue(session: dict, name: str) -> dict: + user_id = session["user_id"] + db = get_db() + cursor = await db.execute( + "SELECT COUNT(*) AS n FROM mcp_tokens WHERE user_id = ? AND revoked_at IS NULL", + (user_id,), + ) + active = (await cursor.fetchone())["n"] + if active >= ACTIVE_LIMIT: + raise HTTPException( + status_code=409, + detail=f"too many active tokens ({active}/{ACTIVE_LIMIT}): revoke one first", + ) + # снейпшот роли — из users (пишется при логине); свежие данные важнее + user_row = await get_user_row(user_id) + role = (user_row["system_role"] if user_row and user_row["system_role"] else "") or "user" + token = generate_mcp_token() + now = now_iso() + cursor = await db.execute( + "INSERT INTO mcp_tokens (user_id, name, user_email, system_role, token_hash," + " token_hint, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)", + ( + user_id, + name or "agent", + session.get("email", ""), + role, + hash_key(token), + mcp_token_hint(token), + now, + ), + ) + await db.commit() + token_id = cursor.lastrowid + return { + "id": token_id, + "name": name or "agent", + "user_email": session.get("email", ""), + "system_role": role, + "token": token, # plaintext — ровно один раз + "token_hint": mcp_token_hint(token), + "created_at": now, + "last_used_at": None, + "revoked_at": None, + } + + +@router.get("/me/mcp_tokens") +async def list_my_mcp_tokens(request: Request) -> dict: + """Свои ключи (включая отозванные — так видно историю).""" + session = await require_session(request) + db = get_db() + cursor = await db.execute( + "SELECT * FROM mcp_tokens WHERE user_id = ? ORDER BY id DESC", + (session["user_id"],), + ) + rows = await cursor.fetchall() + active = sum(1 for r in rows if not r["revoked_at"]) + return {"tokens": [_row_out(r) for r in rows], "active_count": active, "limit": ACTIVE_LIMIT} + + +@router.post("/me/mcp_tokens", status_code=201) +async def create_mcp_token(request: Request, payload: TokenCreate) -> dict: + """Выдать персональный ключ; в ответе — plaintext единственный раз.""" + session = await require_session(request) + now = time.monotonic() + log = [t for t in _create_log.get(session["user_id"], []) if now - t < _CREATE_WINDOW] + if len(log) >= _CREATE_MAX: + raise HTTPException(status_code=429, detail="too many tokens issued, retry later") + log.append(now) + _create_log[session["user_id"]] = log + return await _issue(session, payload.name.strip()) + + +@router.post("/me/mcp_tokens/{token_id}/revoke") +async def revoke_mcp_token(token_id: int, request: Request) -> dict: + """Отозвать свой ключ (не удаляем — строка остаётся в списках с хвост-хинтом).""" + session = await require_session(request) + db = get_db() + cursor = await db.execute( + "SELECT id, revoked_at FROM mcp_tokens WHERE id = ? AND user_id = ?", + (token_id, session["user_id"]), + ) + row = await cursor.fetchone() + if row is None: + raise HTTPException(status_code=404, detail="token not found") + if not row["revoked_at"]: + await db.execute( + "UPDATE mcp_tokens SET revoked_at = ? WHERE id = ?", (now_iso(), token_id) + ) + await db.commit() + return {"status": "ok"} + + +# --- Админ: все ключи всех пользователей -------------------------------------- + +admin_router = APIRouter(prefix="/api/v1/admin", dependencies=[Depends(require_admin)]) + + +@admin_router.get("/mcp_tokens") +async def list_all_mcp_tokens() -> list[dict]: + """Все ключи всех пользователей (включая отозванные) — новые сверху.""" + db = get_db() + cursor = await db.execute("SELECT * FROM mcp_tokens ORDER BY id DESC") + return [_row_out(r) for r in await cursor.fetchall()] + + +@admin_router.post("/mcp_tokens/{token_id}/revoke") +async def admin_revoke_mcp_token(token_id: int) -> dict: + db = get_db() + cursor = await db.execute("SELECT id, revoked_at FROM mcp_tokens WHERE id = ?", (token_id,)) + row = await cursor.fetchone() + if row is None: + raise HTTPException(status_code=404, detail="token not found") + if not row["revoked_at"]: + await db.execute( + "UPDATE mcp_tokens SET revoked_at = ? WHERE id = ?", (now_iso(), token_id) + ) + await db.commit() + return {"status": "ok"} \ No newline at end of file diff --git a/backend/app/auth.py b/backend/app/auth.py new file mode 100644 index 0000000..7b64141 --- /dev/null +++ b/backend/app/auth.py @@ -0,0 +1,276 @@ +"""gnexus-auth (OAuth PKCE) интеграция: конфиг клиента, обмен кода, сессии. + +Каркас скопирован с hard-panel (паттерн gnexus-creds): state/PKCE пишем напрямую +в БД (наша — aiosqlite), у SDK берём детали (HttpTokenEndpoint, +HttpRuntimeUserProvider, webhook-верификатор). Синхронные httpx-запросы к +gnexus-auth выполняются в потоке (asyncio.to_thread), чтобы медленный/опавший +auth-сервер не замораживал event loop. + +Режим: TGCLIENT_AUTH_CLIENT_ID пуст → SPA открыт без логина, личности всех +запросов — служебный LOCAL_USER_ID (app/security.py); задан → браузер шлюзится +через gnexus-auth, MCP — по персональным mcp_* или супер-токену. + +Рядом с sessions живёт таблица users (снейпшот роли и настроек, переживает +сессии): upsert профиля из auth никогда не трогает колонку locale. +""" + +import json +import uuid +from datetime import datetime, timedelta, timezone + +import httpx + +from app.config import get_settings +from app.db import get_db + +SESSION_COOKIE = "tgclient_session" + +SCOPES = ["openid", "email", "profile", "roles", "permissions"] + +# Вебхуки, означающие блокировку пользователя на gnexus-auth (канон mcp.md): +# гасят все персональные ключи; разгильдяйство с тем же 401 — на гварде. +BLOCK_EVENTS = {"user.blocked", "user.deleted", "user.archived"} +UNBLOCK_EVENTS = {"user.unblocked", "user.restored"} +# разблокировка: успешный логин снимает флаг (см. auth_routes.callback) + +# Лениво собираемый синхронный SDK-контекст (httpx-клиент живёт один на процесс) +_ctx: dict = {} + + +def _now() -> datetime: + return datetime.now(timezone.utc) + + +def get_ctx() -> dict: + global _ctx + if not _ctx: + settings = get_settings() + from gnexus_gauth.config import GAuthConfig + from gnexus_gauth.oauth import HttpTokenEndpoint + from gnexus_gauth.runtime import HttpRuntimeUserProvider + from gnexus_gauth.webhook import HmacWebhookVerifier, JsonWebhookParser + + gconf = GAuthConfig( + base_url=settings.auth_base_url, + client_id=settings.auth_client_id, + client_secret=settings.auth_client_secret, + redirect_uri=settings.auth_redirect_uri, + user_agent="tgclient-mcp", + ) + # gnexus-auth живёт внутри сети — TLS не проверяем (по умолчанию http) + http = httpx.Client(timeout=15.0, verify=False) + _ctx = { + "config": gconf, + "http": http, + "token": HttpTokenEndpoint(gconf, http), + "runtime": HttpRuntimeUserProvider(gconf, http), + "webhook": HmacWebhookVerifier(gconf), + "parser": JsonWebhookParser(), + } + return _ctx + + +def auth_enabled() -> bool: + return bool(get_settings().auth_client_id) + + +def safe_return_to(value: str | None) -> str: + """Только внутренние пути: '/path', не '//evil' и не 'https://…'.""" + if not value or not value.startswith("/") or value.startswith("//") or "\\" in value: + return "/" + return value + + +def allowlisted(email: str, user_id: str) -> bool: + """TGCLIENT_AUTH_ALLOWLIST: пусто = любой залогиненный; иначе через + запятую e-mail или user_id.""" + raw = get_settings().auth_allowlist.strip() + if not raw: + return True + allowed = {item.strip().lower() for item in raw.split(",") if item.strip()} + return email.lower() in allowed or str(user_id).lower() in allowed + + +def exchange_and_fetch(code: str, pkce_verifier: str) -> tuple: + """Синхронный обмен кода → (TokenSet, AuthenticatedUser). Выполнять в треде.""" + ctx = get_ctx() + token_set = ctx["token"].exchange_authorization_code(code, pkce_verifier) + user = ctx["runtime"].fetch_user(token_set.access_token) + return token_set, user + + +# --- Сессии в SQLite + users (снейпшоты, переживающие сессии) ------------------- + +def _display_name_from_profile(profile: dict) -> str: + return profile.get("display_name") or profile.get("name") or "" + + +async def upsert_user_row(user, user_id: str | None = None, email: str | None = None) -> None: + """users(upsert): профиль gnexus-auth verbatim. Колонку locale НЕ трогаем — + она меняется только из настроек сервиса, override переживает sync профиля. + system_role пишем только когда auth его отдал (None из вебхук-профиля — + не повод затирать снейпшот, снятый при логине).""" + profile = (user.profile or {}) if hasattr(user, "profile") else (user if isinstance(user, dict) else {}) + if not isinstance(profile, dict): + profile = {} + row_id = user_id or getattr(user, "user_id", None) + row_email = email or getattr(user, "email", None) or "" + role = (getattr(user, "system_role", None) or "").strip() + now = _now().isoformat() + db = get_db() + await db.execute( + "INSERT INTO users (user_id, email, display_name, profile, system_role, created_at, updated_at)" + " VALUES (?, ?, ?, ?, ?, ?, ?)" + " ON CONFLICT(user_id) DO UPDATE SET email=excluded.email," + " display_name=excluded.display_name, profile=excluded.profile," + " system_role=CASE WHEN excluded.system_role != '' THEN excluded.system_role" + " ELSE users.system_role END," + " updated_at=excluded.updated_at", + ( + row_id, + row_email, + _display_name_from_profile(profile) or row_email, + json.dumps(profile, ensure_ascii=False), + role, + now, + now, + ), + ) + await db.commit() + + +async def create_session(user) -> str: + """user — gnexus_gauth.dto.AuthenticatedUser. Возвращает id сессии.""" + db = get_db() + session_id = uuid.uuid4().hex + profile = user.profile or {} + settings = get_settings() + await db.execute( + "INSERT INTO sessions (id, user_id, email, display_name, avatar_url, expires_at, created_at)" + " VALUES (?, ?, ?, ?, ?, ?, ?)", + ( + session_id, + user.user_id, + user.email, + profile.get("display_name") or profile.get("name") or user.email, + profile.get("avatar_url") or "", + (_now() + timedelta(seconds=settings.session_ttl_seconds)).isoformat(), + _now().isoformat(), + ), + ) + await db.commit() + await upsert_user_row(user) + return session_id + + +async def get_session(session_id: str | None): + """Свежая сессия по cookie-значению или None.""" + if not session_id: + return None + db = get_db() + cursor = await db.execute( + "SELECT id, user_id, email, display_name, avatar_url" + " FROM sessions WHERE id = ? AND expires_at > ?", + (session_id, _now().isoformat()), + ) + row = await cursor.fetchone() + return dict(row) if row else None + + +async def delete_session(session_id: str | None) -> None: + if not session_id: + return + db = get_db() + await db.execute("DELETE FROM sessions WHERE id = ?", (session_id,)) + await db.commit() + + +async def purge_expired() -> None: + """Севшие сессии, state и pending-логины — подчистить (фоновый GC-цикл).""" + db = get_db() + now = _now().isoformat() + await db.execute("DELETE FROM sessions WHERE expires_at < ?", (now,)) + await db.execute("DELETE FROM oauth_states WHERE expires_at < ?", (now,)) + await db.execute("DELETE FROM login_sessions WHERE expires_at < ?", (now,)) + await db.commit() + + +async def apply_webhook_update(user_id: str, event_type: str, metadata: dict) -> None: + """Webhook gnexus-auth — отражаем на нашей стороне (канон auth.md): + + auth.global_logout — у пользователя заканчиваются все его сессии (SPA-cookie + перестаёт работать); персональные MCP-ключи НЕ трогаем — самостоятельные + креды (mcp.md). user.blocked/deleted/archived — блокировка: флаг в users + гасит все персональные ключи на гварде require_mcp; остальное трактуем + как user.updated (обновление имени/аватара). + """ + db = get_db() + cutoff = _now().isoformat() + if event_type == "auth.global_logout": + await db.execute( + "DELETE FROM sessions WHERE user_id = ? AND expires_at > ?", (user_id, cutoff) + ) + elif event_type in BLOCK_EVENTS: + await db.execute("UPDATE users SET blocked = 1 WHERE user_id = ?", (user_id,)) + elif event_type in UNBLOCK_EVENTS: + await db.execute("UPDATE users SET blocked = 0 WHERE user_id = ?", (user_id,)) + else: + profile = (metadata.get("profile") or {}) if isinstance(metadata, dict) else {} + display_name = profile.get("display_name") or profile.get("name") + avatar = profile.get("avatar_url") + if display_name: + await db.execute( + "UPDATE sessions SET display_name = ? WHERE user_id = ? AND expires_at > ?", + (display_name, user_id, cutoff), + ) + if profile: + await upsert_user_row( + profile, user_id=user_id, + email=metadata.get("email") or profile.get("email"), + ) + await db.commit() + + +async def get_user_row(user_id: str) -> dict | None: + """Строка users для /me и снейпшота роли; None, если не входил.""" + db = get_db() + cursor = await db.execute( + "SELECT user_id, email, display_name, profile, locale, system_role, blocked" + " FROM users WHERE user_id = ?", + (user_id,), + ) + row = await cursor.fetchone() + return dict(row) if row else None + + +async def set_user_locale(user_id: str, locale: str | None, email: str, display_name: str) -> None: + """Записать override языка (None = авто). Строки users может не быть — + пользователь мог не залогиниться повторно после добавления таблицы.""" + db = get_db() + now = _now().isoformat() + await db.execute( + "INSERT INTO users (user_id, email, display_name, profile, locale, created_at, updated_at)" + " VALUES (?, ?, ?, '{}', ?, ?, ?)" + " ON CONFLICT(user_id) DO UPDATE SET" + " email = CASE WHEN excluded.email != '' THEN excluded.email ELSE users.email END," + " display_name = CASE WHEN excluded.display_name != '' THEN excluded.display_name" + " ELSE users.display_name END," + " locale = excluded.locale, updated_at = excluded.updated_at", + (user_id, email, display_name, locale, now, now), + ) + await db.commit() + + +async def ensure_local_user_async() -> None: + """Асинхронный вариант ensure_local_user (lifespan, только auth-off).""" + db = get_db() + now = _now().isoformat() + from app.security import LOCAL_USER_ID + + await db.execute( + "INSERT INTO users (user_id, email, display_name, profile, created_at, updated_at)" + " VALUES (?, 'local@tgclient.local', 'local', '{}', ?, ?)" + " ON CONFLICT(user_id) DO NOTHING", + (LOCAL_USER_ID, now, now), + ) + await db.commit() \ No newline at end of file diff --git a/backend/app/config.py b/backend/app/config.py new file mode 100644 index 0000000..1d69344 --- /dev/null +++ b/backend/app/config.py @@ -0,0 +1,60 @@ +from functools import lru_cache +from pathlib import Path + +from pydantic import AliasChoices, Field +from pydantic_settings import BaseSettings, SettingsConfigDict + + +class Settings(BaseSettings): + """Настройки tgclient-mcp, всё через env с префиксом TGCLIENT_ (см. .env.example).""" + + model_config = SettingsConfigDict(env_prefix="TGCLIENT_", env_file=".env", extra="ignore") + + # алиасы: краткий TGCLIENT_DB_PATH / TGCLIENT_SESSION_TTL (как в .env.example); + # pydantic-settings не применяет env_prefix к AliasChoices — имена полные + database_path: Path = Field( + default=Path("/data/tgclient.db"), + validation_alias=AliasChoices( + "TGCLIENT_DB_PATH", "TGCLIENT_DATABASE_PATH", + "tgclient_db_path", "tgclient_database_path", + ), + ) + # пусто = MCP доступен только персональным ключам / cookie-сессиям; + # задан — супер-токен для скриптов, MCP, curl (handbook mcp.md: бэкдор) + admin_token: str = "" + # Креды приложения Telegram: один набор на деплой (api_id — приложение); + # пусто — добавление аккаунтов недоступно, тулы отвечают 503-подобной ошибкой. + api_id: int = 0 + api_hash: str = "" + + # Лимиты + max_active_accounts: int = 30 # аккаунтов на пользователя + media_max_bytes: int = 20 * 1024 * 1024 # cap download/upload через MCP + write_limit_n: int = 20 # мутаций на аккаунт за окно + write_limit_window: float = 60.0 + + # gnexus-auth (SSO): задан client_id — браузерная авторизация через SSO + auth_base_url: str = "" + auth_client_id: str = "" + auth_client_secret: str = "" + auth_redirect_uri: str = "" + auth_webhook_secret: str = "" + auth_allowlist: str = "" # пусто = любой залогиненный + session_ttl_seconds: int = Field( + default=7 * 24 * 3600, + validation_alias=AliasChoices( + "TGCLIENT_SESSION_TTL", "TGCLIENT_SESSION_TTL_SECONDS", + "tgclient_session_ttl", "tgclient_session_ttl_seconds", + ), + ) + + # Gnexus Synapse: пустой api_key = интеграция выключена + synapse_url: str = "" + synapse_api_key: str = "" + synapse_default_source: str = "tgclient-mcp" + synapse_timeout: float = 10.0 + + +@lru_cache +def get_settings() -> Settings: + return Settings() \ No newline at end of file diff --git a/backend/app/db.py b/backend/app/db.py new file mode 100644 index 0000000..87da6d7 --- /dev/null +++ b/backend/app/db.py @@ -0,0 +1,43 @@ +import json +from pathlib import Path + +import aiosqlite + +from app.config import get_settings + +_conn: aiosqlite.Connection | None = None + + +async def init_db() -> None: + """Открыть БД, включить WAL, применить schema.sql. + + Одно соединение на процесс: Telethon-сессии пишутся той же транзакцией, + что и остальной сервис — параллельных писателей нет (WAL снимает блокировки). + """ + global _conn + settings = get_settings() + settings.database_path.parent.mkdir(parents=True, exist_ok=True) + _conn = await aiosqlite.connect(settings.database_path) + _conn.row_factory = aiosqlite.Row + await _conn.execute("PRAGMA journal_mode=WAL") + await _conn.execute("PRAGMA foreign_keys=ON") + schema = (Path(__file__).with_name("schema.sql")).read_text(encoding="utf-8") + await _conn.executescript(schema) + await _conn.commit() + + +async def close_db() -> None: + global _conn + if _conn is not None: + await _conn.close() + _conn = None + + +def get_db() -> aiosqlite.Connection: + assert _conn is not None, "database not initialized" + return _conn + + +def j(value) -> str: + """Сериализация JSON-колонок.""" + return json.dumps(value, ensure_ascii=False) \ No newline at end of file diff --git a/backend/app/errors.py b/backend/app/errors.py new file mode 100644 index 0000000..864225f --- /dev/null +++ b/backend/app/errors.py @@ -0,0 +1,13 @@ +"""Доменные ошибки: единый carrier для SPA-роутов (→ HTTPException) и +MCP-тулов (→ данные {"error": code, "detail"} — канон handbook mcp.md).""" + + +class DomainError(Exception): + def __init__(self, code: int, detail: str) -> None: + super().__init__(detail) + self.code = code + self.detail = detail + + +def as_dict(exc: DomainError) -> dict: + return {"error": exc.code, "detail": exc.detail} \ No newline at end of file diff --git a/backend/app/locales.py b/backend/app/locales.py new file mode 100644 index 0000000..e761e98 --- /dev/null +++ b/backend/app/locales.py @@ -0,0 +1,23 @@ +"""Языковые теги сервиса: нормализация и выбор действующего языка. + +Зеркало фронтовых i18n-слагов: override из настроек сервиса → язык аккаунта +gnexus-auth (users.profile.locale) → 'en'. +""" + +SUPPORTED_LOCALES = ("en", "uk", "ru") + + +def normalize_locale(tag) -> str | None: + """'en-US' -> 'en'; всё вне SUPPORTED_LOCALES -> None.""" + text = str(tag or "").strip().lower() + prefix = text.split("-", 1)[0].split("_", 1)[0] + return prefix if prefix in SUPPORTED_LOCALES else None + + +def effective_locale(override, profile) -> str: + """override -> язык аккаунта -> 'en' (совпадает с фолбэком фронта).""" + return ( + normalize_locale(override) + or normalize_locale((profile or {}).get("locale")) + or "en" + ) \ No newline at end of file diff --git a/backend/app/main.py b/backend/app/main.py new file mode 100644 index 0000000..ddaf24b --- /dev/null +++ b/backend/app/main.py @@ -0,0 +1,167 @@ +"""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 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(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 \ No newline at end of file diff --git a/backend/app/mcp/__init__.py b/backend/app/mcp/__init__.py new file mode 100644 index 0000000..e69de29 --- /dev/null +++ b/backend/app/mcp/__init__.py diff --git a/backend/app/mcp/context.py b/backend/app/mcp/context.py new file mode 100644 index 0000000..9649267 --- /dev/null +++ b/backend/app/mcp/context.py @@ -0,0 +1,45 @@ +"""Контекст MCP-запроса: какой пользователь за Bearer-ключом. + +Личность резолвит гард require_mcp (app/security.py) и кладёт в +request.state.mcp_user; mcp_endpoint (app/mcp/server.py) переносит её в +ContextVar — тулы читают current_mcp_user(). Роль действует снейпшотом +выпуска ключа (канон mcp.md) ролью админа/юзера на момент выдачи. + +MCP_USER — синтетический superadmin бэкдора статического TGCLIENT_ADMIN_TOKEN. +""" + +from contextvars import ContextVar + +from app.security import SUPERADMIN + +_mcp_user: ContextVar[dict | None] = ContextVar("mcp_user", default=None) + + +def set_mcp_user(user: dict) -> object: + """Гарду: установить пользователя на время запроса, вернуть token для reset().""" + return _mcp_user.set(user) + + +def reset_mcp_user(token: object) -> None: + _mcp_user.reset(token) # type: ignore[arg-type] + + +def current_mcp_user() -> dict: + """{user_id, email, system_role} текущего MCP-вызова; бэкдор для прямых + вызовов вне запроса (тесты) — возвращаем супер-юзера.""" + found = _mcp_user.get() + if found is not None: + return found + return dict(SUPERADMIN) + + +def current_user_id() -> str: + return current_mcp_user()["user_id"] + + +def user_role() -> str: + return (current_mcp_user().get("system_role") or "user").lower() + + +def is_admin() -> bool: + return user_role() in ("admin", "superadmin") \ No newline at end of file diff --git a/backend/app/mcp/serializers.py b/backend/app/mcp/serializers.py new file mode 100644 index 0000000..980438d --- /dev/null +++ b/backend/app/mcp/serializers.py @@ -0,0 +1,166 @@ +"""Сериализация Telethon-объектов в JSON-данные тулов. + +Правила (pitfall): tl-объекты никогда не сериализуются через to_dict() — +только явные поля здесь. Entity наружу — id в bot-API-маркировке +(telethon.utils.get_peer_id: user > 0, chat > 0 без минуса по-telethon-ски +… см. get_peer_id), type, username, title; access_hash не светим — повторный +доступ к чату через client.get_input_entity(dialog_id) (резолвит из кеша +StringSession, вот зачем нужен persist). +""" + +from telethon.utils import get_peer_id as _tg_peer_id + +from app.tg.waveform import unpack_waveform + + +def marked_peer_id(peer) -> int: + """Entity/Peer/int → bot-API-маркированный int id (Telethon.utils). + Юзер >0, малая группа -id, канал/супергруппа -100id; именно эта + маркировка принимается назад через get_input_entity. None-safe.""" + if peer is None: + return 0 + try: + return _tg_peer_id(peer) + except (ValueError, AttributeError, TypeError): + return 0 + + +def entity_dict(entity, *, short: bool = False) -> dict: + """User/Chat/Channel → плоский dict. id — маркированный int. + Меняются через get_input_entity(dialog_id): юзер >0, группа -id, + канал -100id (всё единообразно — маркировка Telethon).""" + if entity is None: + return {"id": 0, "type": "unknown"} + kind = type(entity).__name__ # User / Chat / Channel / ChatForbidden... + if kind == "User": + out = { + "id": marked_peer_id(entity), # сами объекты маркируются правильно + "type": "user", + "title": " ".join(p for p in (getattr(entity, "first_name", ""), + getattr(entity, "last_name", "")) if p), + "username": getattr(entity, "username", "") or "", + "is_bot": bool(getattr(entity, "bot", False)), + "is_self": bool(getattr(entity, "is_self", False)), + } + elif kind == "Channel": + out = { + "id": marked_peer_id(entity), + "type": "channel", + "title": getattr(entity, "title", "") or "", + "username": getattr(entity, "username", "") or "", + "is_broadcast": bool(getattr(entity, "broadcast", False)), + } + else: # Chat / ChatForbidden / ChatEmpty + out = { + "id": marked_peer_id(entity) if kind != "ChatForbidden" else getattr(entity, "id", 0), + "type": "chat", + "title": getattr(entity, "title", "") or "", + } + if short: + return {"id": out["id"], "type": out["type"], "title": out.get("title", ""), + "username": out.get("username", "")} + return out + + +def dialog_dict(dialog, entity) -> dict: + """Диалог: entity + счётчики + превью последнего сообщения.""" + d = entity_dict(entity) + d["unread_count"] = getattr(dialog, "unread_count", 0) + d["pinned"] = bool(getattr(dialog, "pinned", False)) + d["is_muted"] = bool(getattr(dialog, "muted", False)) or bool( + getattr(getattr(dialog, "notify_settings", None), "mute_until", 0) + ) + msg = getattr(dialog, "message", None) + if msg is not None: + d["last_message"] = message_dict(msg) + return d + + +def media_meta(message) -> dict: + """Описание медиа без скачивания: тип, имя, размер, duration, dimensions. + + voice = DocumentAttributeAudio(voice=True); round = DocumentAttributeVideo( + round_message=True) — оба Document; это важно и для download_media и для UI. + """ + media = getattr(message, "media", None) + if media is None: + return {} + kind = type(media).__name__ + out = {"kind": kind} + doc = getattr(media, "document", None) + photo = getattr(media, "photo", None) + if kind == "MessageMediaDocument" and doc is not None: + attrs = list(getattr(doc, "attributes", []) or []) + for attr in attrs: + aname = type(attr).__name__ + if aname == "DocumentAttributeAudio" and getattr(attr, "voice", False): + out["voice"] = True + out["duration"] = getattr(attr, "duration", None) + out["waveform"] = unpack_waveform(getattr(attr, "waveform", None) or b"") + elif aname == "DocumentAttributeAudio": + out["audio"] = True + out["duration"] = getattr(attr, "duration", None) + out["title"] = getattr(attr, "title", "") or "" + out["performer"] = getattr(attr, "performer", "") or "" + elif aname == "DocumentAttributeVideo": + if getattr(attr, "round_message", False): + out["round"] = True + else: + out["video"] = True + out["duration"] = getattr(attr, "duration", None) + w = getattr(attr, "w", 0) or 0 + h = getattr(attr, "h", 0) or 0 + if w or h: + out["size_hint"] = {"w": w, "h": h} + elif aname == "DocumentAttributeSticker": + out["sticker"] = True + out["emoji"] = getattr(attr, "alt", "") or "" + elif aname == "DocumentAttributeImageSize": + out["size_hint"] = {"w": getattr(attr, "w", 0), "h": getattr(attr, "h", 0)} + elif aname == "DocumentAttributeFilename": + out["file_name"] = getattr(attr, "file_name", "") or "" + out["mime_type"] = getattr(doc, "mime_type", "") or "" + out["size"] = getattr(doc, "size", 0) or 0 + elif kind == "MessageMediaPhoto" and photo is not None: + out["photo"] = True + # размер фото — по самому большому size у photo.sizes (обычный int) + sizes = [s.size for s in (getattr(photo, "sizes", None) or []) + if isinstance(getattr(s, "size", None), int)] + out["size"] = max(sizes) if sizes else None + elif kind == "MessageMediaWebPage": + wp = getattr(media, "webpage", None) + out["webpage_url"] = getattr(wp, "url", "") or "" + out["webpage_title"] = getattr(wp, "title", "") or "" + elif kind == "MessageMediaGeo": + geo = getattr(media, "geo", None) + out["geo"] = {"lat": getattr(geo, "lat", None), "long": getattr(geo, "long", None)} + return out + + +def message_dict(message, *, with_media: bool = True) -> dict: + """Message → плоский dict. MessageService → экшен-строка (join/leave).""" + if type(message).__name__ == "MessageService": + return { + "id": message.id, + "sender_id": marked_peer_id(getattr(message, "from_id", None) or getattr(message, "peer_id", None)), + "action": type(getattr(message, "action", None)).__name__, + } + out = { + "id": message.id, + "date": message.date.isoformat() if getattr(message, "date", None) else None, + "sender_id": marked_peer_id(message.from_id), + "out": bool(message.out), + "text": message.message or "", + "reply_to": getattr(message.reply_to, "reply_to_msg_id", None) + if message.reply_to is not None else None, + "edited": getattr(message, "edit_date", None).isoformat() + if getattr(message, "edit_date", None) else None, + "views": getattr(message, "views", None), + } + if with_media: + out["media"] = media_meta(message) + sender = getattr(message, "sender", None) + if sender is not None: + ent = entity_dict(sender, short=True) + out["sender"] = {"title": ent.get("title"), "username": ent.get("username", "")} + return out \ No newline at end of file diff --git a/backend/app/mcp/server.py b/backend/app/mcp/server.py new file mode 100644 index 0000000..6527fdb --- /dev/null +++ b/backend/app/mcp/server.py @@ -0,0 +1,157 @@ +"""Сборка MCP-сервера tgclient (FastMCP, streamable HTTP, stateless). + +Эндпоинты: /mcp (основной) и /mcp-protocol/ (канон handbook mcp.md — слэш в +конце значим для реверс-прокси). Транспорт stateless + json_response; процесс +тот же, что API деплоя. Авторизация — гард require_mcp на APIRoute +(app/security.py); личность в request.state.mcp_user, здесь переносится в +ContextVar (app/mcp/context.py), тулы читают current_mcp_user(). + +Инструкции агенту (AGENT_INSTRUCTIONS) отдаются и как instructions, и как +prompt agent_guide (паттерн hard-panel / gnexus-creds). +""" + +from fastapi import Request +from mcp.server.fastmcp import FastMCP +from starlette.responses import Response + +from app.mcp.context import set_mcp_user, reset_mcp_user +from app.mcp.tools import register_tools + +AGENT_INSTRUCTIONS = """\ +# tgclient — инструкция для ИИ-агента + +## Назначение +Этот сервер — MCP-интерфейс мульти-аккаунтного Telegram-клиента (MTProto +юзер-аккаунты, не боты). Ты действуешь от лица владельца MCP-ключа: у тебя +доступ только к **его** Telegram-аккаунтам (админ-ключ видит все, но писать +обязан всё равно в чужой чат только по явной просьбе этого человека). + +## Данные и аккаунты +- **Аккаунты** (`accounts_list`): каждый имеет `id`, phone, статус + (`pending|active|logged_out|error`) и отметку живости клиента. +- Почти все тулы принимают `account_id`; если у владельца один активный + аккаунт, указывать его не обязательно (иначе — 422 с просьбой уточнить). +- `dialog_id` — маркированный int из тулов (юзер >0, группа -id, + канал -100id). Не угадывай id: бери из `dialogs_list` / `messages_*`. +- Имя канала/юзера в запросе («@username») — сначала реши через + `chat_info`/`dialogs_list`, чтобы получить id. + +## Карта инструментов — что вызывать +| Вопрос пользователя | Инструмент | +|---|---| +| «какие у меня акки?» | `accounts_list` | +| «кто я в этом акке?» | `me_get` | +| «список чатов / диалогов» | `dialogs_list` | +| «что в этом чате» (последние сообщения) | `messages_history` | +| «что за сообщение N» | `messages_read` | +| «найди сообщения про X» | `messages_search` | +| «напиши в этот чат …» | `message_send` (или `message_reply`) | +| «поправь сообщение N» | `message_edit` (только свои) | +| «удали сообщения …» | `message_delete` | +| «сходи по контактам / кто в группе» | `contacts_list` / `chat_participants` | +| «что за файл тут?» | `messages_read` (media-метаданные) | +| «скачай/файл есть?» | `download_media` (возвращает base64) | +| «отправь файл/фото» | `upload_file` | +| «отправь голосовое» | `upload_voice` (ogg-opus, duration+waveform) | +| «отправь кружок» | `upload_round` (mp4 квадрат ≤60с) | +| «добавь мой акк / залогинь» | `account_login_start` → `code` → [`password`] | + +## Добавление аккаунта (текст-режим — интерактив с человеком) +Телевижн-код и пароль 2FA знает только человек; ты не догадываешься и не +выдумываешь их. Сценарий: +1. Попроси номер телефона и подтверди вслух («логиним +7…123, верно?»). +2. `account_login_start(phone)` → попроси у пользователя код из Telegram. +3. `account_login_code(login_id, code)` — промах не фатален: вернётся + `attempts_left` — вежливо попроси код ещё раз. +4. Если сервер спросит пароль (`step: "awaiting_password"`) — это + Cloud-пароль (2FA): попроси его у пользователя, `account_login_password`. +5. Статус/отмена — `account_login_status` / `account_login_cancel`. +Лимиты: ≤3 параллельных логинов, ограничен поток кодов; код живёт ~10 минут. + +## Правила +- **Мутации (send/edit/delete/upload и т.п.) — только по явной просьбе + пользователя.** Удаление и логаут аккаунта — необратимы. +- Ошибки — данные: `{"error": 403|404|409|413|422|429, "detail": "…"}` — + прочти detail и скажи человеку нормально (или исправь вызов сам). +- 403 по аккаунту = не твой аккаунт: не перебирай чужие id. +- 413 — файл больше лимита; 429 — rate-limit/флуд, скажи сколько ждать + (деталь содержит секунды). +- Голосовой файл (voice) и кружок (round) в медиа-метаданных видны сразу + (`is_voice`/`is_round`, duration и waveform) — не нужен download для + понимания, что это голосовое. +- Личные данные (peer-ids, тексты) — конфиденциальны: не выкладывай наружу. + +## Подключение +``` +claude mcp add --transport http tgclient https:///mcp-protocol/ \\ + --header "Authorization: Bearer mcp_<персональный ключ>" +``` +`/mcp` — historический alias того же сервера. Ключи выпускаются на странице +«MCP-ключи» сервиса (≤10 активных на пользователя, отзыв — там же); роль +ключа — снейпшот роли владельца на момент выпуска. +""" + + +def build_mcp_server() -> FastMCP: + mcp = FastMCP( + "tgclient", + instructions=AGENT_INSTRUCTIONS, + stateless_http=True, + json_response=True, + streamable_http_path="/", + ) + register_tools(mcp) + + @mcp.prompt() + def agent_guide() -> str: + """📘 Гайд: как работать с Telegram через tgclient (аккаунты, чаты, медиа).""" + return AGENT_INSTRUCTIONS + + return mcp + + +mcp = build_mcp_server() +mcp_app = mcp.streamable_http_app() + +# ASGI-приложениеstreamable-http: FastAPI валидирует сигнатуру endpoint-а, +# поэтому тело в receive() вручную и ответ собираем из send() (stateless + +# json_response — ответ всегда один JSON, SSE-моста нет). +mcp_asgi = mcp_app.router.routes[0].app + + +async def mcp_endpoint(request: Request) -> "Response": + """Мост FastAPI → ASGI streamable-http-app. + + Личность из require_mcp (request.state.mcp_user) переносится в ContextVar + на время выполнения запроса — тулы видят current_mcp_user(). + """ + user = getattr(request.state, "mcp_user", None) + body = await request.body() + received = False + start: dict = {} + chunks: list[bytes] = [] + + async def receive(): + nonlocal received + if received: + return {"type": "http.disconnect"} + received = True + return {"type": "http.request", "body": body, "more_body": False} + + async def send(message) -> None: + if message["type"] == "http.response.start": + start.update(message) + else: + chunks.append(message.get("body", b"")) + + ctx_token = set_mcp_user(user) if user else None + try: + await mcp_asgi(request.scope, receive, send) + finally: + if ctx_token is not None: + reset_mcp_user(ctx_token) + response = Response( + content=b"".join(chunks) if chunks else b"", status_code=start["status"] + ) + response.raw_headers.extend(start.get("headers") or []) + return response \ No newline at end of file diff --git a/backend/app/mcp/tools.py b/backend/app/mcp/tools.py new file mode 100644 index 0000000..fa11d7d --- /dev/null +++ b/backend/app/mcp/tools.py @@ -0,0 +1,707 @@ +"""Каталог 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, + marked_peer_id, + 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 +PARTICIPANTS_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}"} + + +def _data_error(code: int, detail: str) -> dict: + return {"error": code, "detail": detail} + + +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): + code = 422 + raise DomainError(422, "dialog_id обязателен (маркированный int из dialogs_list)") + client, row = await _resolve_client(db, account_id, writable=writable) + try: + entity = await client.get_input_entity(dialog_id) + 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: int = 0) -> dict: + """ℹ️ Инфо о чате/канале/юзере: полное описание, счётчики участников. + dialog_id — маркированный int из dialogs_list или @username.""" + 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) -> list: + """👥 Участники чата/канала (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") + # глобальный поиск: по всем диалогам (GetSearchRequest global=True) + client, _row = await _resolve_client(db, account_id) + from_user_id = await _resolve_from_user(client, from_user) if from_user else None + result = await client(functions.messages.SearchGlobalRequest( + q=query, from_id=from_user_id, 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, expect_mime="audio/ogg") + 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, expect_mime="video/mp4") + 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, *, expect_mime: 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 \ No newline at end of file diff --git a/backend/app/schema.sql b/backend/app/schema.sql new file mode 100644 index 0000000..67ea025 --- /dev/null +++ b/backend/app/schema.sql @@ -0,0 +1,100 @@ +-- tgclient-mcp: схема SQLite. +-- Время — UTC ISO-8601 (строки, лексикографически сортируются корректно). + +-- Незавершённые OAuth-потоки gnexus-auth (state → pkce_verifier, 10 минут) +CREATE TABLE IF NOT EXISTS oauth_states ( + state TEXT PRIMARY KEY, + pkce_verifier TEXT NOT NULL, + return_to TEXT NOT NULL DEFAULT '/', + expires_at TEXT NOT NULL +); + +-- Браузерные сессии gnexus-auth (cookie ↔ профиль; TTL по expires_at) +CREATE TABLE IF NOT EXISTS sessions ( + id TEXT PRIMARY KEY, + user_id TEXT NOT NULL, + email TEXT NOT NULL, + display_name TEXT NOT NULL DEFAULT '', + avatar_url TEXT NOT NULL DEFAULT '', + expires_at TEXT NOT NULL, + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_sessions_user_id ON sessions(user_id); + +-- Пользователи gnexus-auth: пер-сервисные настройки + снейпшоты для гейтов. +-- profile = словарь профиля verbatim (язык аккаунта — profile.locale); +-- locale = override из настроек сервиса (NULL = следовать аккаунту). +CREATE TABLE IF NOT EXISTS users ( + user_id TEXT PRIMARY KEY, -- sub из gnexus-auth (= sessions.user_id) + email TEXT NOT NULL, + display_name TEXT NOT NULL DEFAULT '', + profile TEXT NOT NULL DEFAULT '{}', -- JSON профиля gnexus-auth verbatim + locale TEXT, -- override 'en'|'uk'|'ru'; NULL = авто + system_role TEXT, -- снейпшот system_role gnexus-auth + blocked INTEGER NOT NULL DEFAULT 0, -- вебхук user.blocked/deleted/archived + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); + +-- Персональные MCP-ключи (handbook 10-platform/mcp.md, канон = gnexus-synapse): +-- plaintext показывается один раз, в БД — sha256-хэш (unique) + хвост-хинт; +-- снейпшот роли владельца на момент выпуска; вместо архива — ревок. +CREATE TABLE IF NOT EXISTS mcp_tokens ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id TEXT NOT NULL, -- sub gnexus-auth (выдаёт сам себе) + name TEXT NOT NULL DEFAULT '', -- «агент Navi rei» + user_email TEXT NOT NULL DEFAULT '', -- снейпшот для списков + system_role TEXT NOT NULL DEFAULT 'user',-- снейпшот роли выпуска + token_hash TEXT NOT NULL UNIQUE, -- sha256 plaintext-ключа + token_hint TEXT NOT NULL DEFAULT '', -- хвост 8 символов для опознания в списке + created_at TEXT NOT NULL, + last_used_at TEXT, + revoked_at TEXT +); + +CREATE INDEX IF NOT EXISTS idx_mcp_tokens_user ON mcp_tokens(user_id); + +-- Telegram-аккаунты (MTProto юзер-клиенты). Владелец = sub gnexus-auth +-- (режим без SSO — служебный 'local'). Сессия Telethon — только StringSession, +-- сериализуется в session_data (SQLiteSession запретён: файл-конфликт с БД). +CREATE TABLE IF NOT EXISTS accounts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id TEXT NOT NULL REFERENCES users(user_id), + phone TEXT NOT NULL, -- E.164, telethon.utils.parse_phone + label TEXT NOT NULL DEFAULT '', -- заметка владельца («личный») + tg_user_id INTEGER, -- get_me().id после логина + username TEXT NOT NULL DEFAULT '', + display_name TEXT NOT NULL DEFAULT '', + session_data TEXT, -- StringSession.serialize(); NULL = не залогинен + api_id INTEGER, -- NULL → дефолт из env (задел под per-account creds) + api_hash TEXT, + status TEXT NOT NULL DEFAULT 'pending', -- pending|active|logged_out|error + error TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + last_used_at TEXT +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_accounts_owner_phone ON accounts(user_id, phone); +CREATE UNIQUE INDEX IF NOT EXISTS idx_accounts_owner_tgid ON accounts(user_id, tg_user_id) WHERE tg_user_id IS NOT NULL; +CREATE INDEX IF NOT EXISTS idx_accounts_owner ON accounts(user_id); + +-- Незавершённые логины Telegram — общее состояние для SPA и MCP-тулов +-- (человек может начать добавление в UI и продолжить агентом, и наоборот). +-- Живой TelegramClient pending-логина живёт в RAM (AccountManager._pending); +-- код и пароль 2FA в БД никогда не пишутся — сразу в sign_in. +CREATE TABLE IF NOT EXISTS login_sessions ( + id TEXT PRIMARY KEY, -- uuid hex + user_id TEXT NOT NULL REFERENCES users(user_id), + phone TEXT NOT NULL, -- нормализованный E.164 + label TEXT NOT NULL DEFAULT '', -- заметка из UI (доедет на аккаунт) + phone_code_hash TEXT NOT NULL DEFAULT '', -- из send_code_request; стирается + step TEXT NOT NULL DEFAULT 'awaiting_code', -- awaiting_code|awaiting_password + attempts INTEGER NOT NULL DEFAULT 0, -- неверный код; >=3 → отмена + error TEXT NOT NULL DEFAULT '', -- текст последней ошибки шага + expires_at TEXT NOT NULL, -- +15 мин; чистит общий GC-цикл + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_login_sessions_user ON login_sessions(user_id); \ No newline at end of file diff --git a/backend/app/security.py b/backend/app/security.py new file mode 100644 index 0000000..1638ae7 --- /dev/null +++ b/backend/app/security.py @@ -0,0 +1,192 @@ +import hashlib +import secrets +from datetime import datetime, timezone + +from fastapi import Depends, HTTPException, Request, status +from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer + +from app.config import get_settings + +_bearer = HTTPBearer(auto_error=False) + +# Служебный владелец в режиме auth-off (SSO не настроен): у деплоя без +# gnexus-auth личности нет, все аккаунты уходят этому юзеру; MCP — только +# по статическому админ-токену (у открытого режима нет mcp-пути). +LOCAL_USER_ID = "local" + +ADMIN_ROLES = ("admin", "superadmin") + + +def generate_mcp_token() -> str: + """Plaintext персонального MCP-ключа (канон handbook mcp.md, префикс mcp_). + Показывается один раз на странице «MCP-ключи»; в БД — хэш hash_key().""" + return "mcp_" + secrets.token_urlsafe(24) + + +def mcp_token_hint(token: str) -> str: + """Хвост для опознания ключа в списке (plaintext никогда не возвращается).""" + return token[-8:] + + +def hash_key(key: str) -> str: + return hashlib.sha256(key.encode()).hexdigest() + + +def now_iso() -> str: + return datetime.now(timezone.utc).isoformat() + + +def auth_mode_error() -> HTTPException: + """401: браузер не залогинен в gnexus-auth (TGCLIENT_AUTH_CLIENT_ID задан).""" + return HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="not authenticated") + + +def mcp_user(user_id: str, email: str, system_role: str) -> dict: + """Личность mcp-вызова (Role — снейпшот выпуска для персональных ключей).""" + return {"user_id": user_id, "email": email, "system_role": system_role} + + +SUPERADMIN = {"user_id": "mcp", "email": "", "system_role": "superadmin"} + + +async def _cookie_session(request: Request): + """Cookie-сессия или None (auth включён).""" + from app.auth import SESSION_COOKIE, get_session + + if not get_settings().auth_client_id: + return None + return await get_session(request.cookies.get(SESSION_COOKIE)) + + +async def _load_user_row(user_id: str) -> dict | None: + from app.auth import get_user_row + + return await get_user_row(user_id) + + +async def is_session_admin(request: Request) -> bool: + session = await _cookie_session(request) + if session is None: + return False + row = await _load_user_row(session["user_id"]) + return ((row["system_role"] if row and row["system_role"] else "") or "user") in ADMIN_ROLES + + +async def require_admin( + request: Request, + creds: HTTPAuthorizationCredentials | None = Depends(_bearer), +) -> dict: + """Только админ gnexus-auth (мульти-юзер сервис): админка видит всех. + + Порядок: + 1. Bearer TGCLIENT_ADMIN_TOKEN (задан) — статический супер-токен. + 2. Cookie-сессия gnexus-auth c users.system_role ∈ (admin, superadmin). + 3. auth выключен → открытый режим (Bearer-токен не требуется). + Возвращает dict личности (для ответов роутов). + """ + settings = get_settings() + bearer = creds.credentials if creds is not None and creds.scheme.lower() == "bearer" else "" + if settings.admin_token and bearer == settings.admin_token: + return dict(SUPERADMIN) + if settings.auth_client_id: + session = await _cookie_session(request) + if session is None: + raise auth_mode_error() + row = await _load_user_row(session["user_id"]) + role = (row["system_role"] if row and row["system_role"] else "") or "user" + if role not in ADMIN_ROLES: + raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="admin role required") + return mcp_user(session["user_id"], session["email"], role) + if not settings.admin_token: + return dict(SUPERADMIN) + raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="invalid admin token") + + +async def require_session(request: Request) -> dict: + """Только cookie-сессия (личные вещи: выпуск MCP-ключей, PATCH /me, + управление своими аккаунтами): у Bearer-токена нет личности.""" + settings = get_settings() + if settings.auth_client_id: + session = await _cookie_session(request) + if session is None: + raise auth_mode_error() + return session + if settings.admin_token: + raise HTTPException( + status_code=401, detail="personal actions require gnexus-auth sign-in" + ) + # открытый режим: единственная личность деплоя + return {"user_id": LOCAL_USER_ID, "email": "", "display_name": "local"} + + +def _mcp_401() -> HTTPException: + # Причина (отозван / заблокирован / неверный) не раскрывается — канон mcp.md + return HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail="invalid MCP token", + headers={"www-authenticate": "Bearer"}, + ) + + +async def require_mcp( + request: Request, + creds: HTTPAuthorizationCredentials | None = Depends(_bearer), +) -> None: + """Гард /mcp и /mcp-protocol/: персональный mcp_* | статический админ-токен | + cookie-сессия. Это dependency, а не ASGI-гард как в Synapse: /mcp здесь один + APIRoute, dependency решает то же самое без middleware-обёртки. + + Порядок: + 1. Bearer mcp_… — hash-lookup: отозванный или ключ владельца с users.blocked + → тот же 401 (канон mcp.md); валидный → touch last_used_at. + 2. Bearer TGCLIENT_ADMIN_TOKEN — статический супер-токен (бэкдор; роль + superadmin, user_id mcp: им можно всё, владение аккаунта — «mcp»). + 3. Cookie-сессия gnexus-auth — как в require_admin (роль из users). + 4. auth выключен и admin_token пуст → открытый режим. + + Личность кладётся в request.state.mcp_user; mcp_endpoint (app/mcp/server.py) + переносит её в ContextVar, чтобы тулы читали current_mcp_user(). + """ + settings = get_settings() + bearer = creds.credentials if creds is not None and creds.scheme.lower() == "bearer" else "" + if not settings.auth_client_id and not settings.admin_token: + # открытый режим (dev): единственная личность деплоя, любой bearer принят + request.state.mcp_user = mcp_user(LOCAL_USER_ID, "", "superadmin") + return + if bearer.startswith("mcp_"): + from app.db import get_db + + db = get_db() + cursor = await db.execute( + "SELECT id, user_id, user_email, system_role, revoked_at FROM mcp_tokens" + " WHERE token_hash = ?", + (hash_key(bearer),), + ) + row = await cursor.fetchone() + if row is None or row["revoked_at"]: + raise _mcp_401() + user_row = await _load_user_row(row["user_id"]) + if user_row and user_row["blocked"]: + # флаг от вебхуков gnexus-auth; гард молчит тем же 401 — + # не раскрывать статус блокировки (канон mcp.md) + raise _mcp_401() + await db.execute( + "UPDATE mcp_tokens SET last_used_at = ? WHERE id = ?", (now_iso(), row["id"]) + ) + await db.commit() + request.state.mcp_user = mcp_user( + row["user_id"], row["user_email"] or "", row["system_role"] or "user" + ) + return + if settings.admin_token and bearer == settings.admin_token: + request.state.mcp_user = dict(SUPERADMIN) + return + if settings.auth_client_id: + session = await _cookie_session(request) + if session is None: + raise _mcp_401() + row = await _load_user_row(session["user_id"]) + role = (row["system_role"] if row and row["system_role"] else "") or "user" + request.state.mcp_user = mcp_user(session["user_id"], session["email"], role) + return + raise _mcp_401() \ No newline at end of file diff --git a/backend/app/synapse_report.py b/backend/app/synapse_report.py new file mode 100644 index 0000000..dc5bc23 --- /dev/null +++ b/backend/app/synapse_report.py @@ -0,0 +1,96 @@ +"""Исходящие уведомления через Gnexus Synapse (handbook notifications.md). + +Тонкий fire-and-forget репортер (референс: hard-panel events.py, клиент +gn-synapse-client-py): свой ретрай не строим — доставку маршрутизирует Synapse. +Конфиг только env (TGCLIENT_SYNAPSE_*); пустой api_key — интеграция выключена, +код ничего не блокирует и не логирует спамом. +""" + +import asyncio +from dataclasses import dataclass +from datetime import datetime, timezone + +from app.config import get_settings + +TTL_SECONDS = 24 * 3600 + + +@dataclass(frozen=True) +class Sink: + subject: str + action: str + priority: str # low | normal | high | critical + dedup: bool # повторные открытия — dedup по типу+сущность+UTC-день + + +SINKS: dict[str, Sink] = { + "tg_account_logged_in": Sink("account", "logged_in", "high", True), + "tg_account_logged_out": Sink("account", "logged_out", "normal", False), + "tg_account_auth_lost": Sink("account", "auth_lost", "high", True), + "tg_login_failed": Sink("account", "login_failed", "normal", True), + "tg_flood_wait": Sink("api", "flood_wait", "high", True), +} + +_client = None +_warned_disabled = False + + +def synapse_client(): + global _client, _warned_disabled + if _client is not None: + return _client + settings = get_settings() + if not settings.synapse_api_key or not settings.synapse_url: + if not _warned_disabled: + print("synapse reporter disabled (TGCLIENT_SYNAPSE_* unset)", flush=True) + _warned_disabled = True + return None + from gnexus_synapse import AsyncSynapseClient + + _client = AsyncSynapseClient( + settings.synapse_url, + settings.synapse_api_key, + timeout=settings.synapse_timeout, + default_source=settings.synapse_default_source or "tgclient-mcp", + ) + return _client + + +def _utc_date() -> str: + return datetime.now(timezone.utc).strftime("%Y%m%d") + + +def report(event_type: str, payload: dict) -> None: + """Конверт v1 в Synapse, fire-and-forget: падения доставки не роняют код. + + Приоритет high/critical уходит send-ом (ошибки видны в логе), остальное — + тихим emit. Отчёт по имени из SINKS; неизвестное имя — программная ошибка. + """ + sink = SINKS.get(event_type) + if sink is None: + raise ValueError(f"unknown synapse sink: {event_type}") + client = synapse_client() + if client is None: + return + dedup = f"{event_type}-{payload.get('entity') or payload.get('user_id')}-{_utc_date()}" if sink.dedup else None + send = client.send if sink.priority in ("high", "critical") else client.emit + _ship( + send( + subject=sink.subject, + action=sink.action, + priority=sink.priority, + payload=payload, + dedup_key=dedup, + ttl_seconds=TTL_SECONDS, + ) + ) + + +def _ship(coro) -> None: + task = asyncio.create_task(coro) + task.add_done_callback(_ship_done) + + +def _ship_done(task: asyncio.Task) -> None: + if not task.cancelled() and task.exception() is not None: + print(f"synapse report failed: {task.exception()}", flush=True) \ No newline at end of file diff --git a/backend/app/tg/__init__.py b/backend/app/tg/__init__.py new file mode 100644 index 0000000..e69de29 --- /dev/null +++ b/backend/app/tg/__init__.py diff --git a/backend/app/tg/limits.py b/backend/app/tg/limits.py new file mode 100644 index 0000000..34702f7 --- /dev/null +++ b/backend/app/tg/limits.py @@ -0,0 +1,25 @@ +"""Per-account rate limiter на мутации через MCP (канон mcp.md: чувствительные +операции). Token bucket в памяти процесса; после рестарт окно чистое — для +защиты от спама агента хватает (Telegram сам флудит глубже). + +Ключ — account_id: у каждого аккаунта свой Telegram-флуд, лимит честнее +per-user. Плюс отдельное окно — на выдачу кодов логина (login_flow сам). +""" + +import time +from collections import defaultdict, deque + +_windows: dict[int, deque] = defaultdict(deque) + + +def check_write(account_id: int, limit_n: int, window_s: float) -> dict | None: + """None — проход; dict — {"error": 429, "detail": …} для возврата как данных.""" + now = time.monotonic() + win = _windows[account_id] + while win and now - win[0] > window_s: + win.popleft() + if len(win) >= limit_n: + retry = int(window_s - (now - win[0])) + 1 + return {"error": 429, "detail": f"rate limit: {limit_n} writes / {int(window_s)}s per account, retry in ~{retry}s"} + win.append(now) + return None \ No newline at end of file diff --git a/backend/app/tg/login_flow.py b/backend/app/tg/login_flow.py new file mode 100644 index 0000000..836bb5c --- /dev/null +++ b/backend/app/tg/login_flow.py @@ -0,0 +1,334 @@ +"""Стейт-машина добавления 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() \ No newline at end of file diff --git a/backend/app/tg/manager.py b/backend/app/tg/manager.py new file mode 100644 index 0000000..20e2ad1 --- /dev/null +++ b/backend/app/tg/manager.py @@ -0,0 +1,215 @@ +"""AccountManager: пул живых TelegramClient по accounts.id (один процесс). + +- Клиенты создаются лениво (`get_client`), коннект кешируется; Telethon сам + переподключается (auto_reconnect=True). `warmup()` — фоновый eager-connect + активных аккаунтов при старте (не блокирует подъём сервиса). +- Сессии — только StringSession, сериализуются в accounts.session_data + (aiosqlite WAL, одно соединение — см. db.py). SQLiteSession запретён: + файл SQLite внутри контейнера конфликтует с нашей БД. +- Persist: после логина и после мутирующих вызовов (`save_session`) и + периодический дифф-пersist (task в lifespan) — StringSession хранит + кеш entity-хэшей, без этого get_input_entity ломается после рестарта. +- FloodWait: flood_sleep_threshold=60 — Telethon сам спит на коротких; + длинные ловятся в тулах (429 как данные). +""" + +import asyncio +from contextlib import suppress + +from telethon import errors +from telethon.sessions import StringSession + +from app.config import get_settings +from app.db import get_db +from app.security import now_iso +from app.synapse_report import report + + +class AccountManager: + def __init__(self) -> None: + self.clients: dict[int, TelegramClient] = {} + self._pending: dict[str, TelegramClient] = {} # login_session_id → client + + # --- Живые аккаунты ----------------------------------------------------------- + + async def get_client(self, account_id: int) -> TelegramClient: + """Клиент по account_id: connect если нужен, результат кешируется. + + Строка accounts читается прямо из БД (aiosqlite); AuthKeyUnregistered / + UserDeactivated → статус logged_out, клиент отбывается — DomainError. + """ + existing = self.clients.get(account_id) + if existing is not None and existing.is_connected(): + return existing + row = await self._load_account(account_id) + if row is None: + from app.errors import DomainError + + raise DomainError(404, f"account #{account_id} not found") + if row["status"] == "logged_out": + raise DomainError(409, f"account #{account_id} is logged out — login again") + client = TelegramClient( + StringSession(row["session_data"] or ""), + row["api_id"] or get_settings().api_id, + row["api_hash"] or get_settings().api_hash, + auto_reconnect=True, + retry_delay=3, + connection_retries=None, + request_retries=3, + flood_sleep_threshold=60, + catch_up=False, + sequential_updates=False, + ) + await client.connect() + authorized = await client.is_user_authorized() + if not authorized: + await client.disconnect() + await self._mark_status(account_id, "logged_out") + raise DomainError(409, f"account #{account_id} is not authorized — login again") + self.clients[account_id] = client + await self._touch(account_id) + return client + + async def drop(self, account_id: int) -> None: + """Забыть клиент (disconnect): рестарт-логин, логаут, ручной disconnect.""" + client = self.clients.pop(account_id, None) + if client is not None: + with suppress(Exception): + await client.disconnect() + + async def warmup(self) -> None: + """eager-connect всех активных аккаунтов (фоновая задача; ошибки не + блокируют старт и остальные аккаунты).""" + db = get_db() + cursor = await db.execute( + "SELECT id FROM accounts WHERE status = 'active' ORDER BY id" + ) + for row in await cursor.fetchall(): + try: + await self.get_client(row["id"]) + except Exception as exc: # noqa: BLE001 — один упавший не валит остальных + print(f"warmup: account #{row['id']} failed: {exc}", flush=True) + + # --- Pending-логины ------------------------------------------------------------ + + def pending_client(self, login_id: str) -> TelegramClient | None: + return self._pending.get(login_id) + + def put_pending(self, login_id: str, client: TelegramClient) -> None: + self._pending[login_id] = client + + async def drop_pending(self, login_id: str, *, disconnect: bool = True) -> None: + client = self._pending.pop(login_id, None) + if client is not None and disconnect: + with suppress(Exception): + await client.disconnect() + + def pending_ids(self) -> list[str]: + return list(self._pending) + + # --- Persist сессий -------------------------------------------------------------- + + async def save_session(self, account_id: int) -> None: + """session.save() → accounts.session_data (после логина и мутаций: + StringSession хранит entity-кеш, важно для get_input_entity).""" + client = self.clients.get(account_id) + if client is None: + return + data = client.session.save() + db = get_db() + await db.execute( + "UPDATE accounts SET session_data = ?, updated_at = ? WHERE id = ?", + (data, now_iso(), account_id), + ) + await db.commit() + + async def persist_loop(self) -> None: + """Периодический дифф-пersist всех живых клиентов (изменённые сессии).""" + while True: + try: + for account_id, client in list(self.clients.items()): + stored = await self._stored_session(account_id) + current = client.session.save() + if stored != current: + await self.save_session(account_id) + except Exception as exc: # не роняем цикл + print(f"session persist error: {exc}", flush=True) + await asyncio.sleep(60) + + # --- Внутреннее ------------------------------------------------------------------- + + async def _load_account(self, account_id: int): + db = get_db() + cursor = await db.execute( + "SELECT * FROM accounts WHERE id = ?", (account_id,) + ) + return await cursor.fetchone() + + async def _stored_session(self, account_id: int) -> str | None: + cursor = await get_db().execute( + "SELECT session_data FROM accounts WHERE id = ?", (account_id,) + ) + row = await cursor.fetchone() + return row["session_data"] if row else None + + async def _mark_status(self, account_id: int, status: str, error: str = "") -> None: + db = get_db() + await db.execute( + "UPDATE accounts SET status = ?, error = ?, updated_at = ? WHERE id = ?", + (status, error[:500], now_iso(), account_id), + ) + await db.commit() + + async def note_auth_lost(self, account_id: int, detail: str) -> None: + """AuthKeyUnregistered/аналог: аккаунт мёртв, клиент выкинут, Synapse + уведомлён (ttl/dedup делает репортер).""" + await self.drop(account_id) + await self._mark_status(account_id, "logged_out", detail) + db = get_db() + cursor = await db.execute("SELECT user_id FROM accounts WHERE id = ?", (account_id,)) + row = await cursor.fetchone() + report( + "tg_account_auth_lost", + {"entity": f"account-{account_id}", "user_id": row["user_id"] if row else None, + "account_id": account_id, "detail": detail}, + ) + + async def shutdown(self) -> None: + for account_id in list(self.clients): + self.clients[account_id].disconnect() # fire-and-forget is fine + self.clients.pop(account_id, None) + for login_id in list(self._pending): + await self.drop_pending(login_id) + + async def count(self) -> dict: + """Счётчики для /health: active/connected/pending_logins.""" + db = get_db() + cursor = await db.execute( + "SELECT COUNT(*) AS n FROM accounts WHERE status = 'active'" + ) + active = (await cursor.fetchone())["n"] + return {"active": active, "connected": len(self.clients), "pending_logins": len(self._pending)} + + +async def call_with_flood_guard(manager, account_id: int, coro_fn): + """Обёртка вызовов Telethon: длинный FloodWaitError → 429 как данные. + + Короткие (≤60 с) Telethon переспит сам (flood_sleep_threshold) и короутина + вернётся. Длинные не ждём: агент получит причину и решит сам. + """ + try: + return await coro_fn() + except errors.FloodWaitError as exc: + report( + "tg_flood_wait", + {"entity": f"account-{account_id}", "account_id": account_id, "seconds": exc.seconds}, + ) + from app.errors import DomainError + + raise DomainError(429, f"telegram flood wait: {exc.seconds}s — retry later") from exc + except errors.AuthKeyUnregisteredError as exc: + await manager.note_auth_lost(account_id, str(exc)) + raise DomainError(409, "telegram session is dead (auth key unregistered)") from exc + except errors.UserDeactivatedError as exc: + await manager.note_auth_lost(account_id, str(exc)) + raise DomainError(409, "telegram user deactivated") from exc \ No newline at end of file diff --git a/backend/app/tg/waveform.py b/backend/app/tg/waveform.py new file mode 100644 index 0000000..4603f83 --- /dev/null +++ b/backend/app/tg/waveform.py @@ -0,0 +1,69 @@ +"""Голосовые: duration из ogg-opus и waveform 5-бит. + +- Duration ogg-opus: последний page заголовков содержит гранула-позицию + (granule-position): по Ogg-спеке (финальный гранул = длительность в + сэмплах 48 кГц Opus). Парсим финальную гранулу по последнему page + (заголовок 'OggS' — байты 6..14 little-endian). Скачивать/декодировать + PCM не нужно. +- Waveform: Telegram голосовое носит DocumentAttributeAudio(voice=True, + waveform=bytes) — 63 семпла по 5 бит = 315 бит → 40 байт. Принимаем + семплы от агента (список 0..100), квантуем и пакуем; без семплов — + прямой бар высоты по умолчанию (клиент рисует плоский). +""" + +import base64 +import struct + +WAVEFORM_SAMPLES = 63 # Telegram: ровно 63 семпла голосовой волны + + +def ogg_opus_duration(data: bytes) -> int | None: + """Длительность ogg-opus в секундах (final granule / 48000), None если + формат не узнан. Работает по хвосту файла — не декодирует PCM.""" + # последний page 'OggS' в хвосте файла (страницы до 64К; ищем с конца) + idx = data.rfind(b"OggS") + if idx < 0 or idx + 18 > len(data): + return None + granule = int.from_bytes(data[idx + 6 : idx + 14], "little", signed=True) + if granule <= 0: + return None + return max(1, round(granule / 48000)) + + +def pack_waveform(samples: list[int]) -> bytes: + """63 семпла 0..100 → 40 байт 5-бит (старшие биты первыми — как у TG).""" + if not samples: + return b"" + vals = samples[:WAVEFORM_SAMPLES] + if len(vals) < WAVEFORM_SAMPLES: + vals = vals + [0] * (WAVEFORM_SAMPLES - len(vals)) + bits = "" + for v in vals: + q = max(0, min(31, round(v / 100 * 31))) + bits += format(q, "05b") + out = bytearray((len(bits) + 7) // 8) + for i in range(len(bits)): # старший бит первого байта — первый семпл + if bits[i] == "1": + out[i // 8] |= 0x80 >> (i % 8) + return bytes(out) + + +def unpack_waveform(blob: bytes) -> list[int]: + """40 байт 5-бит → 63 числа 0..100 (для сериализации сообщений).""" + if not blob: + return [] + bits = "".join(format(b, "08b") for b in blob) + vals = [] + for i in range(WAVEFORM_SAMPLES): + off = i * 5 + if off + 5 > len(bits): + break + q = int(bits[off : off + 5], 2) + vals.append(round(q / 31 * 100)) + return vals + + +def base64_or_none(data: str | None) -> bytes: + if not data: + return b"" + return base64.b64decode(data, validate=False) \ No newline at end of file diff --git a/backend/pyproject.toml b/backend/pyproject.toml new file mode 100644 index 0000000..4b3490c --- /dev/null +++ b/backend/pyproject.toml @@ -0,0 +1,25 @@ +[project] +name = "tgclient-mcp" +version = "0.1.0" +description = "tgclient-mcp — мульти-аккаунтный Telegram-клиент (Telethon) с MCP-сервером для ИИ-агентов и SPA-админкой (экосистема Gnexus)" +requires-python = ">=3.12" +dependencies = [ + "fastapi>=0.115", + "uvicorn[standard]>=0.30", + "aiosqlite>=0.20", + "pydantic>=2.7", + "pydantic-settings>=2.3", + "mcp>=1.10,<2", + "telethon>=1.36", + "cryptg>=0.4", # быстрый AES для MTProto (иначе Telethon качает/шлёт заметно медленнее) + "gnexus-gauth @ git+https://git.gnexus.space/git/root/gnexus-auth-client-py.git", + # установка по свежему тегу (handbook 10-platform/notifications.md) + "gnexus-synapse @ git+https://git.gnexus.space/git/root/gn-synapse-client-py.git@v0.1.2", +] + +[build-system] +requires = ["setuptools>=68"] +build-backend = "setuptools.build_meta" + +[tool.setuptools.packages.find] +include = ["app*"] \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..652fdcf --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,24 @@ +# tgclient-mcp — docker-compose (один сервис + volume data:/data) +# Секреты — в .env рядом (chmod 600), в git не попадают (.gitignore). +name: tgclient + +services: + tgclient: + build: + context: . + dockerfile: backend/Dockerfile + restart: unless-stopped + ports: + - "${TGCLIENT_PORT:-8710}:8000" + env_file: + - .env + environment: + # внутри контейнера фиксированный порт; наружу — TGCLIENT_PORT хоста + UVICORN_LOG_LEVEL: info + volumes: + - data:/data + # БД — единственное состояние сервиса (aiosqlite WAL); /data обязателен, + # TGCLIENT_DB_PATH=/data/tgclient.db задаётся в .env + +volumes: + data: \ No newline at end of file