"""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()