From 6d341a9f572733afe7d1aa03e2cf7b24299b37bd Mon Sep 17 00:00:00 2001 From: Vladless Date: Wed, 18 Mar 2026 01:17:44 +0300 Subject: [PATCH] audit of client actions/ protection of admin TG v2/ fix of cash accounting status/ update of dependencies --- .gitignore | 2 + api/depends.py | 13 +- api/main.py | 71 ++- api/v2/router.py | 3 + api/v2/routes/auth.py | 79 ++- api/v2/routes/management.py | 63 +- api/v2/routes/payment_links.py | 12 +- api/v2/routes/root.py | 11 + api/v2/routes/tariffs.py | 69 ++- api/v2/routes/web.py | 142 +++++ api/v2/schemas/__init__.py | 1 + api/v2/schemas/audit.py | 28 + api/v2/schemas/identities.py | 19 +- api/v2/schemas/tariffs.py | 20 + api/v2/schemas/web.py | 36 ++ audit.py | 822 ++++++++++++++++++++++++++ core/cache_config.py | 49 ++ core/redis_cache.py | 65 ++ database/models.py | 53 ++ handlers/admin/users/__init__.py | 3 +- handlers/admin/users/keyboard.py | 98 ++- handlers/admin/users/users_audit.py | 490 +++++++++++++++ handlers/admin/users/users_balance.py | 1 + handlers/admin/users/users_manage.py | 17 +- handlers/payments/pay.py | 1 - handlers/refferal.py | 1 - handlers/start.py | 2 - middlewares/__init__.py | 3 +- middlewares/loggings.py | 80 ++- requirements.txt | 1 + 30 files changed, 2219 insertions(+), 36 deletions(-) create mode 100644 api/v2/routes/web.py create mode 100644 api/v2/schemas/audit.py create mode 100644 api/v2/schemas/tariffs.py create mode 100644 api/v2/schemas/web.py create mode 100644 audit.py create mode 100644 handlers/admin/users/users_audit.py diff --git a/.gitignore b/.gitignore index cb26abc2..3b25f6e9 100644 --- a/.gitignore +++ b/.gitignore @@ -59,5 +59,7 @@ setup.py .github/workflows/ modules/ storage/ +static/web_uploads/ +/web-app/ .license_state \ No newline at end of file diff --git a/api/depends.py b/api/depends.py index 15aa2ed2..9764683e 100644 --- a/api/depends.py +++ b/api/depends.py @@ -2,10 +2,11 @@ import hashlib from collections.abc import AsyncGenerator -from fastapi import Depends, HTTPException, Header, Query +from fastapi import Depends, HTTPException, Header, Query, Request from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession +from audit import set_api_actor from database import async_session_maker, identities as idb from database.models import Admin @@ -27,6 +28,7 @@ def hash_token(token: str) -> str: async def verify_admin_token( admin_id: int = Query(..., alias="tg_id"), token: str = Header(..., alias="X-Token"), + request: Request = None, session: AsyncSession = Depends(get_session), ) -> Admin: hashed = hash_token(token) @@ -34,24 +36,28 @@ async def verify_admin_token( admin = result.scalar_one_or_none() if not admin: raise HTTPException(status_code=401, detail="Unauthorized") + set_api_actor(request, tg_id=admin.tg_id) return admin async def verify_identity_token( x_identity_id: str = Header(..., alias="X-Identity-Id"), token: str = Header(..., alias="X-Token"), + request: Request = None, session: AsyncSession = Depends(get_session), ): """Проверяет пару identity_id + token; возвращает Identity. Для использования в API v2.""" identity = await idb.verify_identity_token(session, x_identity_id, token) if not identity: raise HTTPException(status_code=401, detail="Unauthorized") + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) return identity async def verify_identity_admin( x_identity_id: str = Header(..., alias="X-Identity-Id"), token: str = Header(..., alias="X-Token"), + request: Request = None, session: AsyncSession = Depends(get_session), ): """Проверяет identity + token и что identity.is_admin; для админских ручек v2.""" @@ -60,12 +66,14 @@ async def verify_identity_admin( raise HTTPException(status_code=401, detail="Unauthorized") if not identity.is_admin: raise HTTPException(status_code=403, detail="Forbidden") + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) return identity async def verify_identity_admin_short( x_identity_id: str = Header(..., alias="X-Identity-Id"), token: str = Header(..., alias="X-Token"), + request: Request = None, ): """Проверка админа с короткой сессией (для broadcast и др.), чтобы не держать соединение с БД.""" async with async_session_maker() as session: @@ -75,12 +83,14 @@ async def verify_identity_admin_short( raise HTTPException(status_code=401, detail="Unauthorized") if not identity.is_admin: raise HTTPException(status_code=403, detail="Forbidden") + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) return identity async def verify_admin_token_short( admin_id: int = Query(..., alias="tg_id"), token: str = Header(..., alias="X-Token"), + request: Request = None, ) -> Admin: """Проверка админа с короткой сессией (для broadcast и др.), чтобы не держать соединение с БД.""" hashed = hash_token(token) @@ -90,4 +100,5 @@ async def verify_admin_token_short( await session.commit() if not admin: raise HTTPException(status_code=401, detail="Unauthorized") + set_api_actor(request, tg_id=admin.tg_id) return admin diff --git a/api/main.py b/api/main.py index 0ac3d54b..f6de29a0 100644 --- a/api/main.py +++ b/api/main.py @@ -1,8 +1,14 @@ +import asyncio +import os from time import perf_counter from fastapi import FastAPI, Request +from starlette.middleware.cors import CORSMiddleware +from starlette.staticfiles import StaticFiles -from config import API_LOGGING, API_VERSION +from audit import ensure_api_context, log_api_access, record_api_access_event_background +from config import API_LOGGING, API_VERSION, API_CORS_ORIGINS +from database import async_session_maker from logger import logger if API_VERSION == 1: @@ -19,22 +25,69 @@ app = FastAPI( openapi_url="/api/openapi.json", ) +app.add_middleware( + CORSMiddleware, + allow_origins=API_CORS_ORIGINS, + allow_credentials=True, + allow_methods=["*"], + allow_headers=["*"], +) + @app.middleware("http") async def api_access_log_middleware(request: Request, call_next): + context = ensure_api_context(request) if not API_LOGGING: - return await call_next(request) + response = await call_next(request) + response.headers["X-Request-Id"] = context.request_id + return response started = perf_counter() - response = await call_next(request) - duration_ms = int((perf_counter() - started) * 1000) - client_ip = request.client.host if request.client else "-" - path_qs = request.url.path - if request.url.query: - path_qs = f"{path_qs}?{request.url.query}" + try: + response = await call_next(request) + except Exception as exc: + duration_ms = int((perf_counter() - started) * 1000) + log_api_access( + request, + status_code=500, + duration_ms=duration_ms, + result="fail", + reason=type(exc).__name__, + ) + asyncio.create_task( + record_api_access_event_background( + async_session_maker, + request, + result="fail", + reason=type(exc).__name__, + status_code=500, + ) + ) + raise - logger.info(f'[API] {client_ip} "{request.method} {path_qs}" {response.status_code} {duration_ms}ms') + duration_ms = int((perf_counter() - started) * 1000) + response.headers["X-Request-Id"] = context.request_id + result = "success" if response.status_code < 400 else "fail" + log_api_access( + request, + status_code=response.status_code, + duration_ms=duration_ms, + result=result, + ) + asyncio.create_task( + record_api_access_event_background( + async_session_maker, + request, + result=result, + reason=None if response.status_code < 400 else str(response.status_code), + status_code=response.status_code, + ) + ) return response app.include_router(api_router) + +_web_uploads_dir = "static/web_uploads" +os.makedirs(_web_uploads_dir, exist_ok=True) +app.mount("/api/web/uploads", StaticFiles(directory=_web_uploads_dir), name="web_uploads") diff --git a/api/v2/router.py b/api/v2/router.py index 3e241076..2134506e 100644 --- a/api/v2/router.py +++ b/api/v2/router.py @@ -17,6 +17,7 @@ from api.v2.routes import ( settings, payment_links, identities, + web, ) router = APIRouter() @@ -27,6 +28,7 @@ router.include_router(users.router, prefix="/api/users", tags=["Users"]) router.include_router(keys.router, prefix="/api/keys", tags=["Keys"]) router.include_router(coupons.router, prefix="/api/coupons", tags=["Coupons"]) router.include_router(servers.router, prefix="/api/servers", tags=["Servers"]) +router.include_router(tariffs.public_router, prefix="/api/tariffs", tags=["Tariffs"]) router.include_router(tariffs.router, prefix="/api/tariffs", tags=["Tariffs"]) router.include_router(gifts.router, prefix="/api/gifts", tags=["Gifts"]) router.include_router(referrals.router, prefix="/api/referrals", tags=["Referrals"]) @@ -37,3 +39,4 @@ router.include_router(misc.router, prefix="/api") router.include_router(modules.router, prefix="/api") router.include_router(management.router, prefix="/api/management", tags=["Management"]) router.include_router(settings.router, prefix="/api/settings", tags=["Settings"]) +router.include_router(web.router, prefix="", tags=["Web"]) \ No newline at end of file diff --git a/api/v2/routes/auth.py b/api/v2/routes/auth.py index cf6f83cd..5d3807de 100644 --- a/api/v2/routes/auth.py +++ b/api/v2/routes/auth.py @@ -1,15 +1,18 @@ -from fastapi import APIRouter, Depends, HTTPException +from fastapi import APIRouter, Depends, HTTPException, Request from sqlalchemy.ext.asyncio import AsyncSession +from audit import set_api_actor from api.depends import get_session, verify_identity_token from api.v2.schemas.identities import ( IdentityResponse, LinkTelegramRequest, + LoginByCodeRequest, LoginRequest, LoginResponse, LoginTelegramRequest, RegisterByEmailRequest, RegisterResponse, + SendLoginCodeRequest, ) from config import API_TOKEN_TTL_DAYS, API_TOKEN from database import identities as idb @@ -23,6 +26,7 @@ TELEGRAM_LOGIN_MAX_AGE = 86400 # 24 часа @router.post("/register", response_model=RegisterResponse) async def register_by_email( body: RegisterByEmailRequest, + request: Request, session: AsyncSession = Depends(get_session), ): ( @@ -39,12 +43,14 @@ async def register_by_email( if existing: raise HTTPException(status_code=409, detail="Идентичность с таким email уже существует") identity, token = await idb.create_identity_with_token(session, email=email, password=body.password) + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) return RegisterResponse(identity_id=identity.id, token=token) @router.post("/login", response_model=LoginResponse) async def login( body: LoginRequest, + request: Request, session: AsyncSession = Depends(get_session), ): """Вход по email и паролю. Возвращает identity_id и новый токен. Срок действия токена: """ + TOKEN_TTL_HINT + "." @@ -55,12 +61,73 @@ async def login( if not result: raise HTTPException(status_code=401, detail="Неверный email или пароль") identity, token = result + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) + return LoginResponse(identity_id=identity.id, token=token) + + +_LOGIN_CODES: dict[str, tuple[str, float]] = {} +_LOGIN_CODE_TTL = 600.0 # 10 min + + +def _clean_login_codes() -> None: + import time + now = time.time() + for k in list(_LOGIN_CODES): + if now - _LOGIN_CODES[k][1] > _LOGIN_CODE_TTL: + del _LOGIN_CODES[k] + + +@router.post("/send-login-code") +async def send_login_code( + body: SendLoginCodeRequest, + request: Request, + session: AsyncSession = Depends(get_session), +): + """Отправить код входа на email. Код хранится на сервере 10 мин (для демо — без реальной отправки письма).""" + _clean_login_codes() + email = body.email.strip().lower() + if not email: + raise HTTPException(status_code=400, detail="Email обязателен") + identity = await idb.get_identity_by_email(session, email) + if not identity: + raise HTTPException(status_code=404, detail="Аккаунт с таким email не найден") + import secrets + import time + code = "".join(secrets.choice("0123456789") for _ in range(6)) + _LOGIN_CODES[email] = (code, time.time()) + return {"ok": True, "message": "Код отправлен на почту"} + + +@router.post("/login-by-code", response_model=LoginResponse) +async def login_by_code( + body: LoginByCodeRequest, + request: Request, + session: AsyncSession = Depends(get_session), +): + """Вход по email и коду из письма.""" + _clean_login_codes() + email = body.email.strip().lower() + if not email or not body.code or not body.code.strip(): + raise HTTPException(status_code=400, detail="Email и код обязательны") + stored = _LOGIN_CODES.get(email) + if not stored: + raise HTTPException(status_code=400, detail="Код не найден или истёк. Запросите новый.") + code_value, _ = stored + if body.code.strip() != code_value: + raise HTTPException(status_code=401, detail="Неверный код") + del _LOGIN_CODES[email] + identity = await idb.get_identity_by_email(session, email) + if not identity: + raise HTTPException(status_code=401, detail="Аккаунт не найден") + token = await idb.issue_token_for_identity(session, identity) + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) return LoginResponse(identity_id=identity.id, token=token) @router.post("/login-telegram", response_model=LoginResponse) async def login_telegram( body: LoginTelegramRequest, + request: Request, session: AsyncSession = Depends(get_session), ): ( @@ -73,22 +140,28 @@ async def login_telegram( raise HTTPException(status_code=401, detail="Неверная подпись или устаревшие данные от Telegram") identity = await idb.get_or_create_identity_for_tg(session, body.id) token = await idb.issue_token_for_identity(session, identity) + set_api_actor(request, identity_id=identity.id, tg_id=identity.tg_id) return LoginResponse(identity_id=identity.id, token=token) @router.post("/link-telegram", response_model=IdentityResponse) async def link_telegram( body: LinkTelegramRequest, + request: Request, session: AsyncSession = Depends(get_session), identity=Depends(verify_identity_token), ): - """Привязывает Telegram (tg_id) к текущей идентичности. Требуется X-Identity-Id и X-Token.""" - result = await idb.attach_telegram(session, identity.id, body.tg_id) + """Привязывает Telegram к текущей идентичности. Требуется подпись от Telegram Login Widget (доказательство владения аккаунтом).""" + payload = body.model_dump(mode="json") + if not verify_telegram_login(payload, API_TOKEN, max_age_seconds=TELEGRAM_LOGIN_MAX_AGE): + raise HTTPException(status_code=401, detail="Неверная подпись или устаревшие данные от Telegram") + result = await idb.attach_telegram(session, identity.id, body.id) if not result: raise HTTPException( status_code=409, detail="Этот Telegram уже привязан к другой идентичности", ) + set_api_actor(request, identity_id=result.id, tg_id=result.tg_id) return IdentityResponse.model_validate(result) diff --git a/api/v2/routes/management.py b/api/v2/routes/management.py index 2e524036..f19f7782 100644 --- a/api/v2/routes/management.py +++ b/api/v2/routes/management.py @@ -9,12 +9,14 @@ import psutil from aiogram import Bot from aiogram.client.default import DefaultBotProperties from aiogram.enums import ParseMode -from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException +from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Query from pydantic import BaseModel from sqlalchemy import distinct, exists, func, select, update from sqlalchemy.ext.asyncio import AsyncSession from api.depends import get_session, verify_identity_admin, verify_identity_admin_short +from api.v2.schemas.audit import AuditEventListResponse, AuditEventResponse +from audit import drain_audit_redis_to_db, list_audit_events from config import API_TOKEN, BOT_SERVICE from database import async_session_maker, save_blocked_user_ids from core.bootstrap import MANAGEMENT_CONFIG @@ -175,6 +177,65 @@ async def get_broadcast_clusters( return {"clusters": clusters} +@router.get("/audit-events", response_model=AuditEventListResponse) +async def get_audit_events_history( + identity=Depends(verify_identity_admin), + session: AsyncSession = Depends(get_session), + identity_id: str | None = Query(None, description="Фильтр по identity_id"), + tg_id: int | None = Query(None, description="Фильтр по Telegram user id"), + channel: str | None = Query(None, description="api или telegram"), + event_type: str | None = Query(None, description="Точный event_type"), + limit: int = Query(100, ge=1, le=500), + offset: int = Query(0, ge=0), +): + """История аудита клиента по identity_id и/или tg_id.""" + if identity_id is None and tg_id is None: + raise HTTPException(status_code=400, detail="Укажите identity_id или tg_id") + + events = await list_audit_events( + session, + identity_id=identity_id, + tg_id=tg_id, + channel=channel, + event_type=event_type, + limit=limit, + offset=offset, + ) + return AuditEventListResponse( + items=[ + AuditEventResponse( + id=getattr(event, "id", None), + event_type=event.event_type, + channel=event.channel, + actor_identity_id=event.actor_identity_id, + actor_tg_id=event.actor_tg_id, + path_or_handler=event.path_or_handler, + entity_type=event.entity_type, + entity_id=event.entity_id, + result=event.result, + reason=event.reason, + metadata=event.metadata_, + request_id=event.request_id, + created_at=event.created_at, + ) + for event in events + ], + limit=limit, + offset=offset, + ) + + +@router.post("/audit-drain") +async def post_audit_drain(identity=Depends(verify_identity_admin_short)): + """Выгружает буфер аудита из Redis в БД. Для вызова по крону (например 0 0 * * * в 00:00).""" + try: + count = await drain_audit_redis_to_db(async_session_maker) + return {"success": True, "drained": count} + except Exception as exc: + logger.warning("audit-drain failed: %s", exc) + raise HTTPException(status_code=500, detail=str(exc)) from exc + + @router.post("/broadcast") async def launch_broadcast( payload: BroadcastLaunchPayload, diff --git a/api/v2/routes/payment_links.py b/api/v2/routes/payment_links.py index 2f7b1c7f..be630383 100644 --- a/api/v2/routes/payment_links.py +++ b/api/v2/routes/payment_links.py @@ -1,4 +1,4 @@ -from fastapi import APIRouter, Depends, HTTPException +from fastapi import APIRouter, Depends, HTTPException, Request from sqlalchemy.ext.asyncio import AsyncSession from api.depends import get_session, verify_identity_token @@ -28,12 +28,16 @@ async def _resolve_tg_id(body: PaymentLinkCreateRequest, session: AsyncSession) @router.post("/", response_model=PaymentLinkCreateResponse) async def create_link( body: PaymentLinkCreateRequest, + http_request: Request, session: AsyncSession = Depends(get_session), identity=Depends(verify_identity_token), ): """Создаёт платёжную ссылку через выбранную кассу (единая точка входа). Принимает identity_id или tg_id.""" - tg_id = await _resolve_tg_id(body, session) - request = PaymentLinkRequest( + try: + tg_id = await _resolve_tg_id(body, session) + except HTTPException: + raise + payment_request = PaymentLinkRequest( tg_id=tg_id, amount=body.amount, currency=body.currency or "RUB", @@ -42,7 +46,7 @@ async def create_link( failure_url=body.failure_url, metadata=body.metadata, ) - result = await create_payment_link(session, request) + result = await create_payment_link(session, payment_request) return PaymentLinkCreateResponse( success=result.success, payment_id=result.payment_id, diff --git a/api/v2/routes/root.py b/api/v2/routes/root.py index d7c2dd7e..43feae24 100644 --- a/api/v2/routes/root.py +++ b/api/v2/routes/root.py @@ -1,5 +1,7 @@ from fastapi import APIRouter +from config import PROJECT_NAME, USERNAME_BOT + router = APIRouter(tags=["Root"]) @@ -11,3 +13,12 @@ async def root(): @router.get("/api/version", include_in_schema=True) async def version(): return {"version": 2, "api": "v2"} + + +@router.get("/api/telegram-widget-bot", include_in_schema=True) +async def telegram_widget_bot(): + """Имя бота и имя проекта для веб-клиента.""" + return { + "bot_username": USERNAME_BOT.replace("@", ""), + "project_name": (PROJECT_NAME or "Solo").strip() if isinstance(PROJECT_NAME, str) else "Solo", + } diff --git a/api/v2/routes/tariffs.py b/api/v2/routes/tariffs.py index d9a15176..b1b69783 100644 --- a/api/v2/routes/tariffs.py +++ b/api/v2/routes/tariffs.py @@ -1,9 +1,76 @@ -from fastapi import APIRouter +from fastapi import APIRouter, Depends, Query +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession +from api.depends import get_session from api.v2.schemas import TariffBase, TariffResponse, TariffUpdate +from api.v2.schemas.tariffs import TariffGroup, TariffPublic from api.v2.base_crud import generate_crud_router from database.models import Tariff + +def _tariff_to_public(t: Tariff) -> TariffPublic: + return TariffPublic( + id=t.id, + name=t.name or "", + group_code=t.group_code or "", + duration_days=t.duration_days or 0, + price_rub=t.price_rub or 0, + traffic_limit=t.traffic_limit, + device_limit=t.device_limit, + subgroup_title=t.subgroup_title, + sort_order=t.sort_order, + vless=bool(getattr(t, "vless", False)), + ) + + +public_router = APIRouter() + + +@public_router.get("/groups", response_model=list[TariffGroup]) +async def get_tariff_groups(session: AsyncSession = Depends(get_session)): + """Публичный список групп тарифов — уникальные значения колонки group_code.""" + q = ( + select(Tariff.group_code) + .where(Tariff.is_active == True, Tariff.group_code.isnot(None), Tariff.group_code != "") + .distinct() + .order_by(Tariff.group_code) + ) + result = await session.execute(q) + values = result.scalars().all() + return [TariffGroup(group_code=v or "") for v in values] + + +@public_router.get("/public", response_model=list[TariffPublic]) +async def get_tariffs_public( + group_code: str | None = Query(None, description="Фильтр по группе тарифов"), + tariff_ids: str | None = Query(None, description="ID тарифов через запятую (приоритет над группой)"), + filter_vless: str | None = Query( + None, + description="vless: только для роутера (vless=True), app: только для приложения (vless=False), иначе все", + ), + session: AsyncSession = Depends(get_session), +): + """Публичный список активных тарифов (без авторизации).""" + q = select(Tariff).where(Tariff.is_active == True).order_by(Tariff.sort_order.asc().nulls_last(), Tariff.price_rub.asc()) + if tariff_ids: + try: + ids = [int(x.strip()) for x in tariff_ids.split(",") if x.strip()] + if ids: + q = q.where(Tariff.id.in_(ids)) + except ValueError: + pass + elif group_code: + q = q.where(Tariff.group_code == group_code) + if filter_vless == "router": + q = q.where(Tariff.vless == True) + elif filter_vless == "app": + q = q.where(Tariff.vless == False) + result = await session.execute(q) + rows = result.scalars().all() + return [_tariff_to_public(t) for t in rows] + + router = generate_crud_router( model=Tariff, schema_response=TariffResponse, diff --git a/api/v2/routes/web.py b/api/v2/routes/web.py new file mode 100644 index 00000000..75b2f0f3 --- /dev/null +++ b/api/v2/routes/web.py @@ -0,0 +1,142 @@ +import uuid +from pathlib import Path + +from fastapi import APIRouter, Depends, File, HTTPException, UploadFile +from pydantic import BaseModel +from sqlalchemy import select, delete +from sqlalchemy.ext.asyncio import AsyncSession + +from api.depends import get_session, verify_identity_admin +from api.v2.schemas import WebPageResponse, WebPageUpdate, WebBlockResponse, WebTheme +from api.v2.schemas.web import WebUploadResponse +from database.models import WebPage, WebBlock, WebTheme as WebThemeModel + +UPLOAD_DIR = Path("static/web_uploads") +ALLOWED_EXTENSIONS = frozenset({".png", ".jpg", ".jpeg", ".gif", ".webp", ".svg", ".mp4", ".webm"}) +MAX_FILE_SIZE = 100 * 1024 * 1024 + +router = APIRouter(tags=["Web"]) + + +class WebPagesListResponse(BaseModel): + slugs: list[str] + + +KNOWN_PAGE_SLUGS = ["landing", "tariffs", "faq", "login", "dashboard"] + + +@router.get("/api/web/pages", response_model=WebPagesListResponse) +async def list_web_pages( + session: AsyncSession = Depends(get_session), +): + """Список slug всех страниц сайта (для экспорта/импорта). Включает известные страницы, даже если запись ещё не создана.""" + result = await session.execute(select(WebPage.slug).order_by(WebPage.slug)) + from_db = {row[0] for row in result.fetchall()} + slugs = sorted(from_db | set(KNOWN_PAGE_SLUGS)) + return WebPagesListResponse(slugs=slugs) + + +async def get_or_create_page(session: AsyncSession, slug: str) -> WebPage: + result = await session.execute(select(WebPage).where(WebPage.slug == slug)) + page = result.scalar_one_or_none() + if page is None: + page = WebPage(slug=slug, title=slug) + session.add(page) + await session.flush() + return page + + +@router.get("/api/web/pages/{slug}", response_model=WebPageResponse) +async def get_web_page( + slug: str, + session: AsyncSession = Depends(get_session), +): + await get_or_create_page(session, slug) + + blocks_result = await session.execute( + select(WebBlock).where(WebBlock.page_slug == slug).order_by(WebBlock.order, WebBlock.id) + ) + blocks = [WebBlockResponse.model_validate(b) for b in blocks_result.scalars().all()] + + theme_result = await session.execute(select(WebThemeModel).where(WebThemeModel.page_slug == slug)) + theme_row = theme_result.scalar_one_or_none() + theme = WebTheme(tokens=theme_row.tokens) if theme_row else None + + return WebPageResponse(slug=slug, blocks=blocks, theme=theme) + + +@router.put("/api/web/pages/{slug}", response_model=WebPageResponse) +async def update_web_page( + slug: str, + body: WebPageUpdate, + session: AsyncSession = Depends(get_session), + identity=Depends(verify_identity_admin), +): + await get_or_create_page(session, slug) + + await session.execute(delete(WebBlock).where(WebBlock.page_slug == slug)) + + for block in body.blocks: + session.add( + WebBlock( + page_slug=slug, + order=block.order, + type=block.type, + data=block.data, + ) + ) + + theme_row = None + if body.theme is not None: + result = await session.execute(select(WebThemeModel).where(WebThemeModel.page_slug == slug)) + theme_row = result.scalar_one_or_none() + if theme_row is None: + theme_row = WebThemeModel(page_slug=slug, tokens=body.theme.tokens) + session.add(theme_row) + else: + theme_row.tokens = body.theme.tokens + + await session.flush() + + blocks_result = await session.execute( + select(WebBlock).where(WebBlock.page_slug == slug).order_by(WebBlock.order, WebBlock.id) + ) + blocks = [WebBlockResponse.model_validate(b) for b in blocks_result.scalars().all()] + + if theme_row is None: + theme_result = await session.execute(select(WebThemeModel).where(WebThemeModel.page_slug == slug)) + theme_row = theme_result.scalar_one_or_none() + theme = WebTheme(tokens=theme_row.tokens) if theme_row else None + + return WebPageResponse(slug=slug, blocks=blocks, theme=theme) + + +@router.post("/api/web/upload", response_model=WebUploadResponse) +async def upload_media( + file: UploadFile = File(...), + identity=Depends(verify_identity_admin), +): + """Upload image or video for landing blocks and return same-origin URL.""" + if not file.filename or "." not in file.filename: + raise HTTPException(400, "Файл должен иметь расширение") + ext = Path(file.filename).suffix.lower() + if ext not in ALLOWED_EXTENSIONS: + raise HTTPException( + 400, + f"Разрешены только: {', '.join(sorted(ALLOWED_EXTENSIONS))}", + ) + UPLOAD_DIR.mkdir(parents=True, exist_ok=True) + size = 0 + for chunk in file.file: + size += len(chunk) + if size > MAX_FILE_SIZE: + raise HTTPException(400, f"Размер файла не более {MAX_FILE_SIZE // (1024*1024)} МБ") + await file.seek(0) + name = f"{uuid.uuid4().hex}{ext}" + path = UPLOAD_DIR / name + with open(path, "wb") as f: + while chunk := await file.read(64 * 1024): + f.write(chunk) + url = f"/api/web/uploads/{name}" + return WebUploadResponse(url=url) + diff --git a/api/v2/schemas/__init__.py b/api/v2/schemas/__init__.py index 671e1658..eb7c015b 100644 --- a/api/v2/schemas/__init__.py +++ b/api/v2/schemas/__init__.py @@ -28,3 +28,4 @@ from api.v1.schemas import ( ) from api.v1.schemas.keys import KeyBase, KeyCreateRequest, KeyUpdate from api.v1.schemas.settings import SettingResponse, SettingUpsert +from api.v2.schemas.web import WebBlockResponse, WebTheme, WebPageResponse, WebPageUpdate \ No newline at end of file diff --git a/api/v2/schemas/audit.py b/api/v2/schemas/audit.py new file mode 100644 index 00000000..13b7e354 --- /dev/null +++ b/api/v2/schemas/audit.py @@ -0,0 +1,28 @@ +from __future__ import annotations + +from datetime import datetime +from typing import Any + +from pydantic import BaseModel + + +class AuditEventResponse(BaseModel): + id: int | None = None + event_type: str + channel: str + actor_identity_id: str | None + actor_tg_id: int | None + path_or_handler: str + entity_type: str | None + entity_id: str | None + result: str + reason: str | None + metadata: dict[str, Any] | None + request_id: str | None + created_at: datetime | None + + +class AuditEventListResponse(BaseModel): + items: list[AuditEventResponse] + limit: int + offset: int diff --git a/api/v2/schemas/identities.py b/api/v2/schemas/identities.py index 57df9259..36a1bdaa 100644 --- a/api/v2/schemas/identities.py +++ b/api/v2/schemas/identities.py @@ -40,6 +40,15 @@ class LoginResponse(BaseModel): token: str +class SendLoginCodeRequest(BaseModel): + email: str = Field(..., min_length=1) + + +class LoginByCodeRequest(BaseModel): + email: str = Field(..., min_length=1) + code: str = Field(..., min_length=1) + + class LoginTelegramRequest(BaseModel): """Данные от Telegram Login Widget (кнопка «Войти через Telegram»).""" @@ -53,7 +62,15 @@ class LoginTelegramRequest(BaseModel): class LinkTelegramRequest(BaseModel): - tg_id: int = Field(...) + """Данные от Telegram Login Widget — обязательны для доказательства владения аккаунтом при привязке.""" + + id: int = Field(..., description="Telegram user id (tg_id)") + first_name: str = Field("") + last_name: str | None = None + username: str | None = None + photo_url: str | None = None + auth_date: int = Field(..., description="Unix timestamp от Telegram") + hash: str = Field(..., description="HMAC подпись для проверки на бэкенде") class IdentityAttachEmail(BaseModel): diff --git a/api/v2/schemas/tariffs.py b/api/v2/schemas/tariffs.py new file mode 100644 index 00000000..c9f8442f --- /dev/null +++ b/api/v2/schemas/tariffs.py @@ -0,0 +1,20 @@ +from pydantic import BaseModel + + +class TariffGroup(BaseModel): + """Группа тарифов (group_code) для выбора в лендинге и др.""" + group_code: str + + +class TariffPublic(BaseModel): + """Публичный список тарифов (без авторизации).""" + id: int + name: str + group_code: str + duration_days: int + price_rub: int + traffic_limit: int | None + device_limit: int | None + subgroup_title: str | None + sort_order: int | None + vless: bool = False \ No newline at end of file diff --git a/api/v2/schemas/web.py b/api/v2/schemas/web.py new file mode 100644 index 00000000..7721ebbe --- /dev/null +++ b/api/v2/schemas/web.py @@ -0,0 +1,36 @@ +from typing import Any + +from pydantic import BaseModel, Field + + +class WebBlockBase(BaseModel): + type: str = Field(..., max_length=64) + order: int + data: dict[str, Any] + + +class WebBlockResponse(WebBlockBase): + id: str + + class Config: + from_attributes = True + + +class WebTheme(BaseModel): + tokens: dict[str, Any] + + +class WebPageResponse(BaseModel): + slug: str + blocks: list[WebBlockResponse] + theme: WebTheme | None = None + + +class WebPageUpdate(BaseModel): + blocks: list[WebBlockBase] + theme: WebTheme | None = None + + +class WebUploadResponse(BaseModel): + url: str + diff --git a/audit.py b/audit.py new file mode 100644 index 00000000..a6ed01e7 --- /dev/null +++ b/audit.py @@ -0,0 +1,822 @@ +from __future__ import annotations + +import json +import uuid + +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone +from types import SimpleNamespace +from typing import Any, Iterable + +from aiogram.types import CallbackQuery, InlineQuery, Message, TelegramObject, User +from fastapi import Request +from sqlalchemy import and_, delete, desc, or_, select +from sqlalchemy.ext.asyncio import AsyncSession + +from database.models import AuditEvent +from logger import logger + +try: + from core.cache_config import ( + AUDIT_REDIS_BUFFER_ENABLED, + AUDIT_REDIS_DRAIN_BATCH, + AUDIT_REDIS_FLUSH_KEY, + AUDIT_REDIS_IDENTITY_PREFIX, + AUDIT_REDIS_USER_PREFIX, + AUDIT_REDIS_USER_TTL_SEC, + ) +except ImportError: + AUDIT_REDIS_BUFFER_ENABLED = False + AUDIT_REDIS_FLUSH_KEY = "audit:flush" + AUDIT_REDIS_USER_PREFIX = "audit:user:tg:" + AUDIT_REDIS_IDENTITY_PREFIX = "audit:user:identity:" + AUDIT_REDIS_USER_TTL_SEC = 25 * 3600 + AUDIT_REDIS_DRAIN_BATCH = 1000 + +_MAX_TEXT_LEN = 160 +_AUDIT_TABLE_READY = False + + +@dataclass +class AuditContext: + request_id: str + channel: str + path_or_handler: str + actor_identity_id: str | None = None + actor_tg_id: int | None = None + + +def new_request_id() -> str: + return uuid.uuid4().hex + + +def _trim(value: Any, limit: int = _MAX_TEXT_LEN) -> str | None: + if value is None: + return None + text = str(value).strip() + if not text: + return None + if len(text) <= limit: + return text + return text[: limit - 3] + "..." + + +def _jsonable(value: Any) -> Any: + if value is None or isinstance(value, (str, int, float, bool)): + return value + if isinstance(value, datetime): + return value.isoformat() + if isinstance(value, dict): + return {str(k): _jsonable(v) for k, v in value.items()} + if isinstance(value, (list, tuple, set)): + return [_jsonable(item) for item in value] + return _trim(value, 500) + + +def _serialize(payload: dict[str, Any]) -> str: + return json.dumps(_jsonable(payload), ensure_ascii=False, sort_keys=True) + + +def _message_text(event: TelegramObject) -> str | None: + if isinstance(event, Message): + return _trim(event.text or event.caption) + if isinstance(event, CallbackQuery): + return _trim(event.data) + if isinstance(event, InlineQuery): + return _trim(event.query) + return None + + +def _event_user(event: TelegramObject) -> User | None: + if hasattr(event, "from_user") and isinstance(event.from_user, User): + return event.from_user + return None + + +def describe_telegram_event(event: TelegramObject) -> str: + if isinstance(event, Message): + return f"message:{_message_text(event) or '-'}" + if isinstance(event, CallbackQuery): + return f"callback:{_message_text(event) or '-'}" + if isinstance(event, InlineQuery): + return f"inline:{_message_text(event) or '-'}" + return type(event).__name__ + + +def ensure_api_context(request: Request) -> AuditContext: + context = getattr(request.state, "audit_context", None) + if isinstance(context, AuditContext): + return context + + path_or_handler = request.url.path + if request.url.query: + path_or_handler = f"{path_or_handler}?{request.url.query}" + + context = AuditContext( + request_id=new_request_id(), + channel="api", + path_or_handler=path_or_handler, + ) + request.state.audit_context = context + request.state.audit_request_id = context.request_id + return context + + +def get_api_context(request: Request | None) -> AuditContext | None: + if request is None: + return None + context = getattr(request.state, "audit_context", None) + if isinstance(context, AuditContext): + return context + return None + + +def set_api_actor( + request: Request, + *, + identity_id: str | None = None, + tg_id: int | None = None, +) -> AuditContext: + context = ensure_api_context(request) + if identity_id is not None: + context.actor_identity_id = identity_id + if tg_id is not None: + context.actor_tg_id = tg_id + return context + + +def ensure_telegram_context( + data: dict[str, Any] | None, + event: TelegramObject, +) -> AuditContext: + if data is not None: + existing = data.get("audit_context") + if isinstance(existing, AuditContext): + return existing + + user = _event_user(event) + context = AuditContext( + request_id=new_request_id(), + channel="telegram", + path_or_handler=describe_telegram_event(event), + actor_tg_id=user.id if user else None, + ) + if data is not None: + data["audit_context"] = context + data["audit_request_id"] = context.request_id + return context + + +def set_telegram_actor( + audit_context: AuditContext | dict[str, Any] | None, + *, + identity_id: str | None = None, + tg_id: int | None = None, +) -> AuditContext | None: + context: AuditContext | None + if isinstance(audit_context, AuditContext): + context = audit_context + elif isinstance(audit_context, dict): + context = audit_context.get("audit_context") + else: + context = None + + if not isinstance(context, AuditContext): + return None + if identity_id is not None: + context.actor_identity_id = identity_id + if tg_id is not None: + context.actor_tg_id = tg_id + return context + + +def get_telegram_context(audit_context: AuditContext | dict[str, Any] | None) -> AuditContext | None: + if isinstance(audit_context, AuditContext): + return audit_context + if isinstance(audit_context, dict): + context = audit_context.get("audit_context") + if isinstance(context, AuditContext): + return context + return None + + +def log_api_access( + request: Request, + *, + status_code: int, + duration_ms: int, + result: str, + reason: str | None = None, +) -> None: + context = ensure_api_context(request) + client_ip = request.client.host if request.client else "-" + logger.debug( + f"[AUDIT_ACCESS] {_serialize({ + 'request_id': context.request_id, + 'channel': 'api', + 'method': request.method, + 'path': context.path_or_handler, + 'status_code': status_code, + 'duration_ms': duration_ms, + 'result': result, + 'reason': reason, + 'client_ip': client_ip, + 'actor_identity_id': context.actor_identity_id, + 'actor_tg_id': context.actor_tg_id, + })}" + ) + + +def log_telegram_access( + event: TelegramObject, + *, + audit_context: AuditContext | None, + result: str, + reason: str | None = None, +) -> None: + context = audit_context or AuditContext( + request_id=new_request_id(), + channel="telegram", + path_or_handler=describe_telegram_event(event), + ) + user = _event_user(event) + logger.debug( + f"[AUDIT_ACCESS] {_serialize({ + 'request_id': context.request_id, + 'channel': 'telegram', + 'path_or_handler': context.path_or_handler, + 'event_type': type(event).__name__, + 'message': _message_text(event), + 'result': result, + 'reason': reason, + 'actor_identity_id': context.actor_identity_id, + 'actor_tg_id': context.actor_tg_id or (user.id if user else None), + 'username': getattr(user, 'username', None) if user else None, + })}" + ) + + +async def record_audit_event( + session: AsyncSession, + *, + event_type: str, + channel: str, + path_or_handler: str, + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + entity_type: str | None = None, + entity_id: str | int | None = None, + result: str = "success", + reason: str | None = None, + metadata: dict[str, Any] | None = None, + request_id: str | None = None, +) -> AuditEvent: + await ensure_audit_table(session) + event = AuditEvent( + event_type=event_type, + channel=channel, + actor_identity_id=actor_identity_id, + actor_tg_id=actor_tg_id, + path_or_handler=_trim(path_or_handler, 255) or channel, + entity_type=_trim(entity_type, 64), + entity_id=_trim(entity_id, 255), + result=_trim(result, 32) or "success", + reason=_trim(reason, 1000), + metadata_=_jsonable(metadata) if metadata else None, + request_id=_trim(request_id, 64), + ) + session.add(event) + await session.flush() + logger.debug( + f"[AUDIT_EVENT] {_serialize({ + 'id': event.id, + 'request_id': event.request_id, + 'channel': event.channel, + 'event_type': event.event_type, + 'actor_identity_id': event.actor_identity_id, + 'actor_tg_id': event.actor_tg_id, + 'path_or_handler': event.path_or_handler, + 'entity_type': event.entity_type, + 'entity_id': event.entity_id, + 'result': event.result, + 'reason': event.reason, + 'metadata': event.metadata_, + })}" + ) + return event + + +async def safe_record_audit_event(session: AsyncSession, **kwargs: Any) -> AuditEvent | None: + try: + return await record_audit_event(session, **kwargs) + except Exception as exc: + logger.warning(f"[Audit] Не удалось записать событие {kwargs.get('event_type')}: {exc}") + return None + + +def _telegram_access_payload( + audit_context: AuditContext | None, + event: TelegramObject, + *, + result: str = "success", + reason: str | None = None, +) -> dict[str, Any]: + """Собирает payload для записи telegram_access (для фоновой задачи или синхронной).""" + ctx = get_telegram_context(audit_context) + user = _event_user(event) + path = describe_telegram_event(event) + if ctx is None: + ctx = AuditContext( + request_id=new_request_id(), + channel="telegram", + path_or_handler=path, + actor_tg_id=user.id if user else None, + ) + return { + "request_id": ctx.request_id, + "path_or_handler": path, + "actor_identity_id": ctx.actor_identity_id, + "actor_tg_id": ctx.actor_tg_id or (user.id if user else None), + "result": result, + "reason": reason, + } + + +async def record_telegram_access_event( + session: AsyncSession, + audit_context: AuditContext | None, + event: TelegramObject, + *, + result: str = "success", + reason: str | None = None, +) -> AuditEvent | None: + """Пишет в БД одно событие «обработчик Telegram вызван» для полного следа пользователя.""" + payload = _telegram_access_payload(audit_context, event, result=result, reason=reason) + return await safe_record_audit_event( + session, + event_type="telegram_access", + channel="telegram", + path_or_handler=payload["path_or_handler"], + actor_identity_id=payload["actor_identity_id"], + actor_tg_id=payload["actor_tg_id"], + result=payload["result"], + reason=payload["reason"], + request_id=payload["request_id"], + ) + + +def _audit_record_for_redis( + *, + event_type: str = "telegram_access", + channel: str = "telegram", + path_or_handler: str, + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + entity_type: str | None = None, + entity_id: str | int | None = None, + result: str = "success", + reason: str | None = None, + request_id: str | None = None, + metadata_: dict | None = None, +) -> dict[str, Any]: + """Формирует запись для буфера Redis (с created_at в ISO).""" + return { + "event_type": event_type, + "channel": channel, + "path_or_handler": _trim(path_or_handler, 255) or channel, + "actor_identity_id": actor_identity_id, + "actor_tg_id": actor_tg_id, + "entity_type": _trim(entity_type, 64) if entity_type else None, + "entity_id": _trim(str(entity_id), 255) if entity_id is not None else None, + "result": _trim(result, 32) or "success", + "reason": _trim(reason, 1000) if reason else None, + "request_id": _trim(request_id, 64) if request_id else None, + "metadata_": _jsonable(metadata_) if metadata_ else None, + "created_at": datetime.now(timezone.utc).isoformat(), + } + + +async def record_audit_event_to_redis( + *, + request_id: str | None = None, + path_or_handler: str = "", + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + result: str = "success", + reason: str | None = None, +) -> None: + """Пишет событие telegram_access в буфер Redis (списки для выгрузки в БД и для чтения по пользователю).""" + from core.redis_cache import cache_expire, cache_rpush + + record = _audit_record_for_redis( + path_or_handler=path_or_handler, + actor_identity_id=actor_identity_id, + actor_tg_id=actor_tg_id, + result=result, + reason=reason, + request_id=request_id, + ) + await cache_rpush(AUDIT_REDIS_FLUSH_KEY, record) + if actor_tg_id is not None: + user_key = f"{AUDIT_REDIS_USER_PREFIX}{actor_tg_id}" + await cache_rpush(user_key, record) + await cache_expire(user_key, AUDIT_REDIS_USER_TTL_SEC) + if actor_identity_id: + identity_key = f"{AUDIT_REDIS_IDENTITY_PREFIX}{actor_identity_id}" + await cache_rpush(identity_key, record) + await cache_expire(identity_key, AUDIT_REDIS_USER_TTL_SEC) + + +async def record_api_access_event_to_redis( + *, + request_id: str | None = None, + path_or_handler: str = "", + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + result: str = "success", + reason: str | None = None, +) -> None: + """Пишет событие api_access в буфер Redis (как telegram_access).""" + from core.redis_cache import cache_expire, cache_rpush + + record = _audit_record_for_redis( + event_type="api_access", + channel="api", + path_or_handler=path_or_handler, + actor_identity_id=actor_identity_id, + actor_tg_id=actor_tg_id, + result=result, + reason=reason, + request_id=request_id, + ) + await cache_rpush(AUDIT_REDIS_FLUSH_KEY, record) + if actor_tg_id is not None: + user_key = f"{AUDIT_REDIS_USER_PREFIX}{actor_tg_id}" + await cache_rpush(user_key, record) + await cache_expire(user_key, AUDIT_REDIS_USER_TTL_SEC) + if actor_identity_id: + identity_key = f"{AUDIT_REDIS_IDENTITY_PREFIX}{actor_identity_id}" + await cache_rpush(identity_key, record) + await cache_expire(identity_key, AUDIT_REDIS_USER_TTL_SEC) + + +async def record_api_access_event_background( + session_factory: Any, + request: Request, + *, + result: str = "success", + reason: str | None = None, + status_code: int = 200, +) -> None: + """Пишет одно событие api_access в Redis-буфер или в БД в фоне (после обработки запроса). + Вызывается из middleware; actor берётся из request.state (set_api_actor в эндпоинтах).""" + context = ensure_api_context(request) + path_or_handler = f"{request.method} {request.url.path}" + if request.url.query: + path_or_handler = f"{path_or_handler}?{request.url.query}" + path_or_handler = _trim(path_or_handler, 255) or "api" + if AUDIT_REDIS_BUFFER_ENABLED: + try: + await record_api_access_event_to_redis( + request_id=context.request_id, + path_or_handler=path_or_handler, + actor_identity_id=context.actor_identity_id, + actor_tg_id=context.actor_tg_id, + result=result, + reason=reason, + ) + except Exception as exc: + logger.warning("[Audit] Запись api_access в Redis-буфер не удалась: %s", exc) + return + try: + async with session_factory() as session: + await ensure_audit_table(session) + await record_audit_event( + session, + event_type="api_access", + channel="api", + path_or_handler=path_or_handler, + actor_identity_id=context.actor_identity_id, + actor_tg_id=context.actor_tg_id, + result=result, + reason=reason, + request_id=_trim(context.request_id, 64), + ) + await session.commit() + except Exception as exc: + logger.warning( + "[Audit] Фоновая запись api_access не удалась: %s", + exc, + extra={"path_or_handler": path_or_handler[:80] if path_or_handler else None}, + ) + + +async def record_telegram_access_event_background( + session_factory: Any, + *, + request_id: str | None, + path_or_handler: str, + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + result: str = "success", + reason: str | None = None, +) -> None: + """Пишет событие telegram_access: в Redis-буфер (если включён) или в БД в отдельной сессии. + При ошибке логирует и не пробрасывает исключение.""" + if AUDIT_REDIS_BUFFER_ENABLED: + try: + await record_audit_event_to_redis( + request_id=request_id, + path_or_handler=path_or_handler, + actor_identity_id=actor_identity_id, + actor_tg_id=actor_tg_id, + result=result, + reason=reason, + ) + except Exception as exc: + logger.warning("[Audit] Запись в Redis-буфер не удалась: %s", exc) + return + try: + async with session_factory() as session: + await ensure_audit_table(session) + await record_audit_event( + session, + event_type="telegram_access", + channel="telegram", + path_or_handler=_trim(path_or_handler, 255) or "telegram", + actor_identity_id=actor_identity_id, + actor_tg_id=actor_tg_id, + result=result, + reason=reason, + request_id=_trim(request_id, 64), + ) + await session.commit() + except Exception as exc: + logger.warning( + "[Audit] Фоновая запись telegram_access не удалась: %s", + exc, + extra={"path_or_handler": path_or_handler[:80] if path_or_handler else None}, + ) + + +async def safe_record_api_event( + session: AsyncSession, + request: Request, + *, + event_type: str, + entity_type: str | None = None, + entity_id: str | int | None = None, + result: str = "success", + reason: str | None = None, + metadata: dict[str, Any] | None = None, + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + path_or_handler: str | None = None, +) -> AuditEvent | None: + context = ensure_api_context(request) + return await safe_record_audit_event( + session, + event_type=event_type, + channel="api", + path_or_handler=path_or_handler or context.path_or_handler, + actor_identity_id=actor_identity_id if actor_identity_id is not None else context.actor_identity_id, + actor_tg_id=actor_tg_id if actor_tg_id is not None else context.actor_tg_id, + entity_type=entity_type, + entity_id=entity_id, + result=result, + reason=reason, + metadata=metadata, + request_id=context.request_id, + ) + + +async def safe_record_telegram_event( + session: AsyncSession, + audit_context: AuditContext | dict[str, Any] | None, + *, + event_type: str, + entity_type: str | None = None, + entity_id: str | int | None = None, + result: str = "success", + reason: str | None = None, + metadata: dict[str, Any] | None = None, + actor_identity_id: str | None = None, + actor_tg_id: int | None = None, + path_or_handler: str | None = None, +) -> AuditEvent | None: + context = get_telegram_context(audit_context) + return await safe_record_audit_event( + session, + event_type=event_type, + channel="telegram", + path_or_handler=path_or_handler or (context.path_or_handler if context else "telegram"), + actor_identity_id=actor_identity_id if actor_identity_id is not None else (context.actor_identity_id if context else None), + actor_tg_id=actor_tg_id if actor_tg_id is not None else (context.actor_tg_id if context else None), + entity_type=entity_type, + entity_id=entity_id, + result=result, + reason=reason, + metadata=metadata, + request_id=context.request_id if context else None, + ) + + +def _redis_record_to_event_like(rec: dict[str, Any]) -> SimpleNamespace: + """Превращает запись из Redis в объект с теми же атрибутами, что и AuditEvent.""" + created = rec.get("created_at") + if isinstance(created, str): + try: + created = datetime.fromisoformat(created.replace("Z", "+00:00")) + except Exception: + created = datetime.now(timezone.utc) + elif created is None: + created = datetime.now(timezone.utc) + return SimpleNamespace( + id=None, + event_type=rec.get("event_type", "telegram_access"), + channel=rec.get("channel", "telegram"), + path_or_handler=rec.get("path_or_handler") or "", + actor_identity_id=rec.get("actor_identity_id"), + actor_tg_id=rec.get("actor_tg_id"), + entity_type=rec.get("entity_type"), + entity_id=rec.get("entity_id"), + result=rec.get("result", "success"), + reason=rec.get("reason"), + metadata_=rec.get("metadata_"), + request_id=rec.get("request_id"), + created_at=created, + ) + + +async def _list_audit_events_from_redis( + tg_id: int | None, + identity_id: str | None, + channel: str | None, + event_types: list[str] | None, + max_events: int = 3000, +) -> list[SimpleNamespace]: + """Читает события пользователя из Redis-буфера (для слияния с БД).""" + if not AUDIT_REDIS_BUFFER_ENABLED: + return [] + from core.redis_cache import cache_lrange + + out: list[SimpleNamespace] = [] + seen: set[tuple[str, str]] = set() + keys_to_read = [] + if tg_id is not None: + keys_to_read.append(f"{AUDIT_REDIS_USER_PREFIX}{tg_id}") + if identity_id: + keys_to_read.append(f"{AUDIT_REDIS_IDENTITY_PREFIX}{identity_id}") + for key in keys_to_read: + raw = await cache_lrange(key, -max_events, -1) + for rec in reversed(raw): + if not isinstance(rec, dict): + continue + created = rec.get("created_at") + rid = rec.get("request_id") or "" + if (created, rid) in seen: + continue + if channel and rec.get("channel") != channel: + continue + if event_types and rec.get("event_type") not in event_types: + continue + seen.add((str(created), rid)) + out.append(_redis_record_to_event_like(rec)) + out.sort(key=lambda e: e.created_at, reverse=True) + return out[:max_events] + + +async def list_audit_events( + session: AsyncSession, + *, + identity_id: str | None = None, + tg_id: int | None = None, + channel: str | None = None, + event_type: str | None = None, + event_types: Iterable[str] | None = None, + limit: int = 100, + offset: int = 0, +) -> list[AuditEvent | SimpleNamespace]: + """Список событий аудита. При включённом Redis-буфере объединяет данные из БД и Redis.""" + await ensure_audit_table(session) + event_types_list = sorted(event_types) if event_types else None + + if not AUDIT_REDIS_BUFFER_ENABLED: + stmt = select(AuditEvent) + actor_filters = [] + if identity_id: + actor_filters.append(AuditEvent.actor_identity_id == identity_id) + actor_filters.append(and_(AuditEvent.entity_type == "identity", AuditEvent.entity_id == identity_id)) + if tg_id is not None: + tg_id_str = str(tg_id) + actor_filters.append(AuditEvent.actor_tg_id == tg_id) + actor_filters.append(and_(AuditEvent.entity_type == "user", AuditEvent.entity_id == tg_id_str)) + actor_filters.append(and_(AuditEvent.entity_type == "telegram_user", AuditEvent.entity_id == tg_id_str)) + if actor_filters: + stmt = stmt.where(or_(*actor_filters)) + if channel: + stmt = stmt.where(AuditEvent.channel == channel) + if event_type: + stmt = stmt.where(AuditEvent.event_type == event_type) + if event_types_list: + stmt = stmt.where(AuditEvent.event_type.in_(event_types_list)) + stmt = stmt.order_by(desc(AuditEvent.created_at), desc(AuditEvent.id)).limit(limit).offset(offset) + result = await session.execute(stmt) + return list(result.scalars().all()) + + redis_events = await _list_audit_events_from_redis( + tg_id, identity_id, channel, event_types_list, max_events=3000 + ) + need = offset + limit + len(redis_events) + stmt = select(AuditEvent) + actor_filters = [] + if identity_id: + actor_filters.append(AuditEvent.actor_identity_id == identity_id) + actor_filters.append(and_(AuditEvent.entity_type == "identity", AuditEvent.entity_id == identity_id)) + if tg_id is not None: + tg_id_str = str(tg_id) + actor_filters.append(AuditEvent.actor_tg_id == tg_id) + actor_filters.append(and_(AuditEvent.entity_type == "user", AuditEvent.entity_id == tg_id_str)) + actor_filters.append(and_(AuditEvent.entity_type == "telegram_user", AuditEvent.entity_id == tg_id_str)) + if actor_filters: + stmt = stmt.where(or_(*actor_filters)) + if channel: + stmt = stmt.where(AuditEvent.channel == channel) + if event_type: + stmt = stmt.where(AuditEvent.event_type == event_type) + if event_types_list: + stmt = stmt.where(AuditEvent.event_type.in_(event_types_list)) + stmt = stmt.order_by(desc(AuditEvent.created_at), desc(AuditEvent.id)).limit(min(5000, need)).offset(0) + result = await session.execute(stmt) + db_events = list(result.scalars().all()) + merged = redis_events + db_events + merged.sort(key=lambda e: (e.created_at, getattr(e, "id", 0)), reverse=True) + return merged[offset : offset + limit] + + +async def ensure_audit_table(session: AsyncSession) -> None: + global _AUDIT_TABLE_READY + if _AUDIT_TABLE_READY: + return + connection = await session.connection() + await connection.run_sync(AuditEvent.__table__.create, checkfirst=True) + _AUDIT_TABLE_READY = True + + +async def delete_old_audit_events( + session: AsyncSession, + *, + older_than_days: int = 90, +) -> int: + """Удаляет события старше N дней. Вызывать по крону/периодике при больших наплывах. + Возвращает количество удалённых строк.""" + await ensure_audit_table(session) + threshold = datetime.now(timezone.utc) - timedelta(days=older_than_days) + stmt = delete(AuditEvent).where(AuditEvent.created_at < threshold) + result = await session.execute(stmt) + return result.rowcount or 0 + + +async def drain_audit_redis_to_db(session_factory: Any) -> int: + """Выгружает буфер аудита из Redis в БД батчами. Вызывать по крону (например в 00:00). + Возвращает количество записанных событий.""" + from core.redis_cache import cache_lpop_batch + + total = 0 + while True: + batch = await cache_lpop_batch(AUDIT_REDIS_FLUSH_KEY, AUDIT_REDIS_DRAIN_BATCH) + if not batch: + break + try: + async with session_factory() as session: + await ensure_audit_table(session) + for rec in batch: + created = rec.get("created_at") + if isinstance(created, str): + try: + created = datetime.fromisoformat(created.replace("Z", "+00:00")) + except Exception: + created = datetime.now(timezone.utc) + elif created is None: + created = datetime.now(timezone.utc) + event = AuditEvent( + event_type=rec.get("event_type", "telegram_access"), + channel=rec.get("channel", "telegram"), + path_or_handler=rec.get("path_or_handler") or "telegram", + actor_identity_id=rec.get("actor_identity_id"), + actor_tg_id=rec.get("actor_tg_id"), + entity_type=rec.get("entity_type"), + entity_id=rec.get("entity_id"), + result=rec.get("result", "success"), + reason=rec.get("reason"), + metadata_=rec.get("metadata_"), + request_id=rec.get("request_id"), + created_at=created, + ) + session.add(event) + await session.commit() + total += len(batch) + except Exception as exc: + logger.warning("[Audit] drain_audit_redis_to_db батч не записан: %s", exc) + break + return total diff --git a/core/cache_config.py b/core/cache_config.py index 9728516f..53e254b6 100644 --- a/core/cache_config.py +++ b/core/cache_config.py @@ -1,3 +1,42 @@ +""" +Сводка по Redis: префиксы ключей и окна жизни (TTL). + +Ключ/префикс │ TTL (сек) │ Назначение +──────────────────────────┼───────────┼──────────────────────────────────────── +throttle_counter │ 1 │ Счётчик троттлинга (окно сброса) +throttle_notice │ 1 │ Показ уведомления «подождите» +concurrency_notice │ 5 │ Сообщение «слишком много запросов» +utm_exists │ 300 │ UTM-код уже обработан (старт) +user_middleware │ 60 │ Снапшот пользователя (debounce*2) +user_snapshot │ 30 │ Снапшот пользователя +user_exists │ 60 │ Пользователь есть в БД +balance │ 25 │ Баланс в профиле +profile_data │ 25 │ Данные профиля +key_count │ 25 │ Количество ключей +ban_status │ 60 │ Статус бана +direct_start_user_exists │ 20 │ Пользователь есть (direct start blocker) +admin_access │ 60 │ Доступ в админку (да/нет) +remna_server │ 300 │ URL сервера Remnawave +remna_profile │ 20/45 │ Профиль Remnawave (45 при ошибке) +runtime_configs │ 86400 │ Рантайм-конфиг (1 сутки) +sub_response │ 20 │ Ответ подписки (subscription) +servers │ 60 │ Список серверов +tariff │ 120 │ Тариф по ID +tariffs_cluster │ 120 │ Тарифы по кластеру +keys_list │ 25 │ Список ключей +key_details │ 45 │ Детали ключа +key_email │ 45 │ email по client_id +payment_pending │ 3600 │ Ожидающий платёж (1 ч) +audit_history │ 300 │ История действий клиента (админка) +audit:flush │ — │ Буфер аудита (список для выгрузки в БД в 00:00) +audit:user:tg:* │ 25 ч │ События по tg_id для чтения до выгрузки +audit:user:identity:* │ 25 ч │ События по identity для чтения до выгрузки +webhook_abuse_fail │ 60 │ Счётчик неудачных вебхуков по IP +webhook_abuse_block │ 300 │ Блокировка IP по злоупотреблению + +При недоступности Redis: повторная попытка подключения через 5 сек +(REDIS_BACKOFF_SEC в redis_cache). +""" UPDATE_STALE_AGE_SEC = 60 CONCURRENCY_MAX_WAIT_SEC = 300 @@ -62,6 +101,16 @@ PROFILE_DATA_CACHE_TTL_SEC = 25 PAYMENT_PENDING_CACHE_TTL_SEC = 3600 + +AUDIT_HISTORY_CACHE_TTL_SEC = 300 + +AUDIT_REDIS_BUFFER_ENABLED = True +AUDIT_REDIS_FLUSH_KEY = "audit:flush" +AUDIT_REDIS_USER_PREFIX = "audit:user:tg:" +AUDIT_REDIS_IDENTITY_PREFIX = "audit:user:identity:" +AUDIT_REDIS_USER_TTL_SEC = 25 * 3600 +AUDIT_REDIS_DRAIN_BATCH = 1000 + ERROR_THROTTLE_WINDOW_SEC = 60 ERROR_THROTTLE_MAX_KEYS = 500 ERROR_THROTTLE_MESSAGE_MAX_LEN = 120 diff --git a/core/redis_cache.py b/core/redis_cache.py index aad1e36c..92e8bf2b 100644 --- a/core/redis_cache.py +++ b/core/redis_cache.py @@ -137,3 +137,68 @@ async def cache_delete_pattern(pattern: str) -> int: except Exception: return deleted return deleted + + +async def cache_rpush(key: str, *values: Any) -> int: + """Добавляет значения в хвост списка. Значения сериализуются в JSON. Возвращает длину списка после или 0 при ошибке.""" + if not values: + return 0 + client = await _get_redis() + if client is None: + return 0 + try: + raw = [json.dumps(v, ensure_ascii=False) for v in values] + return int(await client.rpush(key, *raw)) + except Exception: + return 0 + + +async def cache_expire(key: str, ttl_sec: int) -> bool: + """Устанавливает TTL для ключа. Возвращает True при успехе.""" + client = await _get_redis() + if client is None: + return False + try: + return await client.expire(key, max(1, int(ttl_sec))) + except Exception: + return False + + +async def cache_lrange(key: str, start: int, end: int) -> list[Any]: + """Возвращает срез списка. end=-1 — до конца. Элементы десериализуются из JSON.""" + client = await _get_redis() + if client is None: + return [] + try: + raw_list = await client.lrange(key, start, end) + out = [] + for raw in raw_list: + try: + out.append(json.loads(raw)) + except Exception: + pass + return out + except Exception: + return [] + + +async def cache_lpop_batch(key: str, count: int) -> list[Any]: + """Забирает до count элементов с головы списка (FIFO). Совместимо с Redis < 6.2.""" + if count <= 0: + return [] + client = await _get_redis() + if client is None: + return [] + out = [] + try: + for _ in range(count): + raw = await client.lpop(key) + if raw is None: + break + try: + out.append(json.loads(raw)) + except Exception: + pass + return out + except Exception: + return out diff --git a/database/models.py b/database/models.py index 877d2fa6..ea0ece3a 100644 --- a/database/models.py +++ b/database/models.py @@ -11,6 +11,7 @@ from sqlalchemy import ( DateTime, Float, ForeignKey, + Index, Integer, Numeric, String, @@ -301,6 +302,34 @@ class TrackingSource(DictLikeMixin, Base): created_at = Column(DateTime, default=datetime.utcnow) +class AuditEvent(DictLikeMixin, Base): + """События аудита (флоу пользователя).""" + __tablename__ = "audit_events" + __table_args__ = ( + Index("ix_audit_events_tg_created", "actor_tg_id", "created_at"), + Index("ix_audit_events_identity_created", "actor_identity_id", "created_at"), + ) + + id = Column(Integer, primary_key=True, autoincrement=True) + event_type = Column(String(64), nullable=False, index=True) + channel = Column(String(32), nullable=False, index=True) + actor_identity_id = Column( + String(36), + ForeignKey("identities.id", ondelete="SET NULL", onupdate="CASCADE"), + nullable=True, + index=True, + ) + actor_tg_id = Column(BigInteger, nullable=True, index=True) + path_or_handler = Column(String(255), nullable=False) + entity_type = Column(String(64), nullable=True, index=True) + entity_id = Column(String(255), nullable=True, index=True) + result = Column(String(32), nullable=False, server_default=text("'success'")) + reason = Column(Text, nullable=True) + metadata_ = Column("metadata", JSONB, nullable=True) + request_id = Column(String(64), nullable=True, index=True) + created_at = Column(DateTime, default=datetime.utcnow, index=True) + + class Admin(Base): __tablename__ = "admins" @@ -323,3 +352,27 @@ class Setting(DictLikeMixin, Base): description = Column(Text, nullable=True) created_at = Column(DateTime, default=datetime.utcnow) updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow) + + +class WebPage(DictLikeMixin, Base): + __tablename__ = "web_pages" + + slug = Column(String(64), primary_key=True) + title = Column(String(255), nullable=True) + + +class WebTheme(DictLikeMixin, Base): + __tablename__ = "web_themes" + + page_slug = Column(String(64), ForeignKey("web_pages.slug", ondelete="CASCADE"), primary_key=True) + tokens = Column(JSONB, nullable=False, default=dict) + + +class WebBlock(DictLikeMixin, Base): + __tablename__ = "web_blocks" + + id = Column(String(36), primary_key=True, default=lambda: str(uuid.uuid4())) + page_slug = Column(String(64), ForeignKey("web_pages.slug", ondelete="CASCADE"), index=True, nullable=False) + order = Column(Integer, nullable=False, default=0) + type = Column(String(64), nullable=False) + data = Column(JSONB, nullable=False, default=dict) diff --git a/handlers/admin/users/__init__.py b/handlers/admin/users/__init__.py index 63dcc1e2..0d20543c 100644 --- a/handlers/admin/users/__init__.py +++ b/handlers/admin/users/__init__.py @@ -1,10 +1,11 @@ from aiogram import Router -from . import users_balance, users_bans, users_gifts, users_hwid, users_keys, users_manage, users_tariffs +from . import users_audit, users_balance, users_bans, users_gifts, users_hwid, users_keys, users_manage, users_tariffs router = Router() router.include_router(users_manage.router) +router.include_router(users_audit.router) router.include_router(users_balance.router) router.include_router(users_hwid.router) router.include_router(users_keys.router) diff --git a/handlers/admin/users/keyboard.py b/handlers/admin/users/keyboard.py index a9051d8d..3ee5b5c5 100644 --- a/handlers/admin/users/keyboard.py +++ b/handlers/admin/users/keyboard.py @@ -33,7 +33,12 @@ class AdminUserKeyEditorCallback(CallbackData, prefix="admin_users_key"): edit: bool = False -async def build_user_edit_kb(tg_id: int, key_records: list, is_banned: bool = False) -> InlineKeyboardMarkup: +async def build_user_edit_kb( + tg_id: int, + key_records: list, + is_banned: bool = False, + admin_role: str | None = None, +) -> InlineKeyboardMarkup: builder = InlineKeyboardBuilder() current_time = datetime.now(tz=timezone.utc) @@ -77,6 +82,13 @@ async def build_user_edit_kb(tg_id: int, key_records: list, is_banned: bool = Fa ), ) + builder.row( + InlineKeyboardButton( + text="🕘 История действий", + callback_data=AdminUserEditorCallback(action="users_audit", tg_id=tg_id, data="all|all|0").pack(), + ) + ) + builder.row( InlineKeyboardButton( text="♻️ Восстановить триал", @@ -97,7 +109,7 @@ async def build_user_edit_kb(tg_id: int, key_records: list, is_banned: bool = Fa ), ) - hook_buttons = await run_hooks("admin_user_edit", tg_id=tg_id, is_banned=is_banned) + hook_buttons = await run_hooks("admin_user_edit", tg_id=tg_id, is_banned=is_banned, admin_role=admin_role) builder = insert_hook_buttons(builder, hook_buttons) builder.row(build_editor_btn("🔄 Обновить данные", tg_id, edit=True)) @@ -516,6 +528,88 @@ def build_user_gifts_kb(tg_id: int, gifts: list, page: int = 0) -> InlineKeyboar return builder.as_markup() +def build_user_audit_kb( + tg_id: int, + channel_filter: str = "all", + category_filter: str = "all", + page: int = 0, + has_prev: bool = False, + has_next: bool = False, +) -> InlineKeyboardMarkup: + builder = InlineKeyboardBuilder() + + builder.row( + InlineKeyboardButton( + text="Все", + callback_data=AdminUserEditorCallback(action="users_audit", tg_id=tg_id, data=f"all|{category_filter}|0").pack(), + ), + InlineKeyboardButton( + text="API", + callback_data=AdminUserEditorCallback(action="users_audit", tg_id=tg_id, data=f"api|{category_filter}|0").pack(), + ), + InlineKeyboardButton( + text="Telegram", + callback_data=AdminUserEditorCallback(action="users_audit", tg_id=tg_id, data=f"telegram|{category_filter}|0").pack(), + ), + ) + + category_labels = { + "all": "Все", + "auth": "Auth", + "payments": "Платежи", + "subscriptions": "Подписки", + "marketing": "Маркетинг", + } + category_row = [ + InlineKeyboardButton( + text=category_labels[category_key], + callback_data=AdminUserEditorCallback( + action="users_audit", + tg_id=tg_id, + data=f"{channel_filter}|{category_key}|0", + ).pack(), + ) + for category_key in ("all", "auth", "payments", "subscriptions", "marketing") + ] + builder.row(*category_row[:3]) + builder.row(*category_row[3:]) + + if has_prev or has_next: + nav_buttons = [] + if has_prev: + nav_buttons.append( + InlineKeyboardButton( + text="<", + callback_data=AdminUserEditorCallback( + action="users_audit_page", + tg_id=tg_id, + data=f"{channel_filter}|{category_filter}|{page - 1}", + ).pack(), + ) + ) + nav_buttons.append(InlineKeyboardButton(text=str(page + 1), callback_data="noop")) + if has_next: + nav_buttons.append( + InlineKeyboardButton( + text=">", + callback_data=AdminUserEditorCallback( + action="users_audit_page", + tg_id=tg_id, + data=f"{channel_filter}|{category_filter}|{page + 1}", + ).pack(), + ) + ) + builder.row(*nav_buttons) + + builder.row( + InlineKeyboardButton( + text="Назад", + callback_data=AdminUserEditorCallback(action="users_editor", tg_id=tg_id, edit=True).pack(), + ) + ) + return builder.as_markup() + + def build_gift_delete_confirm_kb(tg_id: int, gift_id: str, page: int = 0) -> InlineKeyboardMarkup: builder = InlineKeyboardBuilder() diff --git a/handlers/admin/users/users_audit.py b/handlers/admin/users/users_audit.py new file mode 100644 index 00000000..78715503 --- /dev/null +++ b/handlers/admin/users/users_audit.py @@ -0,0 +1,490 @@ +import html +from datetime import datetime +from types import SimpleNamespace + +import pytz + +from aiogram import F, Router +from aiogram.types import CallbackQuery, Message +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from audit import list_audit_events +from core.cache_config import AUDIT_HISTORY_CACHE_TTL_SEC +from core.redis_cache import cache_get, cache_key, cache_set +from database.models import User +from filters.admin import IsAdminFilter + +from .keyboard import AdminUserEditorCallback, build_user_audit_kb + + +MOSCOW_TZ = pytz.timezone("Europe/Moscow") +PAGE_SIZE = 10 +router = Router() + + +def _serialize_audit_events(events: list) -> list[dict]: + """Для кэша Redis: список событий в JSON-сериализуемый вид.""" + out = [] + for e in events: + out.append({ + "event_type": e.event_type, + "channel": e.channel, + "path_or_handler": getattr(e, "path_or_handler", None) or "", + "actor_identity_id": getattr(e, "actor_identity_id", None), + "actor_tg_id": getattr(e, "actor_tg_id", None), + "entity_type": getattr(e, "entity_type", None), + "entity_id": getattr(e, "entity_id", None), + "result": getattr(e, "result", "success"), + "reason": getattr(e, "reason", None), + "metadata_": getattr(e, "metadata_", None), + "request_id": getattr(e, "request_id", None), + "created_at": e.created_at.isoformat() if e.created_at else None, + }) + return out + + +def _deserialize_audit_events(cached: list[dict]) -> list: + """Из кэша: список dict → объекты с атрибутами как у AuditEvent.""" + out = [] + for d in cached: + created = d.get("created_at") + if isinstance(created, str): + try: + created = datetime.fromisoformat(created.replace("Z", "+00:00")) + except Exception: + created = None + out.append(SimpleNamespace( + event_type=d.get("event_type", ""), + channel=d.get("channel", "telegram"), + path_or_handler=d.get("path_or_handler") or "", + actor_identity_id=d.get("actor_identity_id"), + actor_tg_id=d.get("actor_tg_id"), + entity_type=d.get("entity_type"), + entity_id=d.get("entity_id"), + result=d.get("result", "success"), + reason=d.get("reason"), + metadata_=d.get("metadata_"), + request_id=d.get("request_id"), + created_at=created, + )) + return out + +EVENT_CATEGORY_MAP = { + "auth": { + "register_success", + "register_failed", + "login_success", + "login_failed", + "login_code_sent", + "login_code_send_failed", + "login_by_code_success", + "login_by_code_failed", + "telegram_login_success", + "telegram_login_failed", + "telegram_link_success", + "telegram_link_failed", + }, + "payments": { + "payment_menu_opened", + "payment_currency_selected", + "balance_screen_opened", + "balance_history_opened", + "payment_link_created", + "payment_link_create_failed", + "coupon_activation_success", + "coupon_activation_failed", + "coupon_renewal_success", + "coupon_renewal_failed", + "coupon_key_selection_opened", + }, + "subscriptions": { + "trial_requested", + "key_purchase_started", + "keys_list_opened", + "key_view_opened", + "key_alias_updated", + "key_hwid_reset", + "key_renew_requested", + "key_renew_blocked", + "key_renew_insufficient_funds", + "key_renew_completed", + "key_renew_failed", + }, + "marketing": { + "start_link_opened", + "coupon_link_opened", + "gift_link_opened", + "referral_link_opened", + "referral_link_applied", + "referral_link_failed", + "referral_screen_opened", + "referral_qr_opened", + "top_referrals_opened", + "utm_link_opened", + "utm_link_failed", + }, +} + +EVENT_TYPE_LABELS = { + "start_entry_opened": "Старт бота", + "start_link_opened": "Переход по ссылке", + "register_success": "Регистрация (успех)", + "register_failed": "Регистрация (ошибка)", + "login_success": "Вход по паролю", + "login_failed": "Вход (ошибка)", + "login_code_sent": "Код входа отправлен", + "login_code_send_failed": "Код входа не отправлен", + "login_by_code_success": "Вход по коду", + "login_by_code_failed": "Вход по коду (ошибка)", + "telegram_login_success": "Вход через Telegram", + "telegram_login_failed": "Вход через Telegram (ошибка)", + "telegram_link_success": "Привязка Telegram", + "telegram_link_failed": "Привязка Telegram (ошибка)", + "payment_menu_opened": "Меню оплаты", + "payment_currency_selected": "Выбор валюты", + "balance_screen_opened": "Экран баланса", + "balance_history_opened": "История пополнений", + "payment_link_created": "Создана платёжная ссылка", + "payment_link_create_failed": "Ошибка создания ссылки", + "coupon_activation_success": "Купон активирован", + "coupon_activation_failed": "Купон не активирован", + "coupon_renewal_success": "Продление по купону", + "coupon_renewal_failed": "Продление по купону (ошибка)", + "coupon_key_selection_opened": "Выбор ключа для купона", + "coupon_link_opened": "Переход по ссылке купона", + "trial_requested": "Запрос триала", + "key_purchase_started": "Начало покупки подписки", + "keys_list_opened": "Список подписок", + "key_view_opened": "Карточка подписки", + "key_alias_updated": "Переименование подписки", + "key_hwid_reset": "Сброс устройств (HWID)", + "key_renew_requested": "Запрос продления", + "key_renew_blocked": "Продление ещё недоступно", + "key_renew_insufficient_funds": "Не хватило баланса на продление", + "key_renew_completed": "Подписка продлена", + "key_renew_failed": "Ошибка продления", + "gift_link_opened": "Переход по подарочной ссылке", + "referral_link_opened": "Переход по реферальной ссылке", + "referral_link_applied": "Реферал применён", + "referral_link_failed": "Реферальная ссылка (ошибка)", + "referral_screen_opened": "Экран рефералов", + "referral_qr_opened": "QR реферальной ссылки", + "top_referrals_opened": "Топ рефералов", + "utm_link_opened": "Переход по UTM", + "utm_link_failed": "UTM не найден", +} + + +def _parse_filter_page(raw_data: str | int | None) -> tuple[str, str, int]: + if isinstance(raw_data, str): + parts = raw_data.split("|") + if len(parts) == 3: + channel_filter, category_filter, page_str = parts + elif len(parts) == 2: + channel_filter, page_str = parts + category_filter = "all" + else: + return "all", "all", 0 + if channel_filter not in {"all", "api", "telegram"}: + channel_filter = "all" + if category_filter not in {"all", "auth", "payments", "subscriptions", "marketing"}: + category_filter = "all" + if page_str.isdigit(): + return channel_filter, category_filter, int(page_str) + return channel_filter, category_filter, 0 + return "all", "all", 0 + + +def _resolve_event_types(category_filter: str) -> list[str] | None: + if category_filter == "all": + return None + category_events = EVENT_CATEGORY_MAP.get(category_filter) + if not category_events: + return None + return sorted(category_events) + + +# Порядок и подписи блоков при показе «все» категории +CATEGORY_BLOCK_ORDER = ("auth", "subscriptions", "payments", "marketing", "other") +CATEGORY_BLOCK_LABELS = { + "auth": "Авторизация", + "payments": "Платежи", + "subscriptions": "Подписки", + "marketing": "Маркетинг", + "other": "Другое", +} + + +def _event_category(event) -> str: + """Определяет категорию события для группировки (при выборке «все»).""" + path = (getattr(event, "path_or_handler", None) or "").lower() + etype = (getattr(event, "event_type", None) or "").lower() + if etype != "telegram_access": + for cat, event_types in EVENT_CATEGORY_MAP.items(): + if etype in event_types: + return cat + return "other" + if "handlers.keys" in path or "key_view" in path or "key_create" in path or "key_renew" in path or "tariffs" in path or "key_tariffs" in path or "addon" in path or "key_addons" in path: + return "subscriptions" + if "handlers.payments" in path or "pay" in path or "balance" in path or "handlers.coupons" in path or "payment" in path: + return "payments" + if "handlers.refferal" in path or "referral" in path or "gift" in path or "utm" in path: + return "marketing" + if "handlers.start" in path or "start_entry" in path or "process_start" in path: + return "marketing" + if "auth" in path or "login" in path or "register" in path: + return "auth" + return "other" + + +def _format_metadata(metadata: dict | None) -> str | None: + if not metadata: + return None + items = [] + for key in sorted(metadata.keys()): + value = metadata[key] + if value is None or value == "": + continue + items.append(f"{key}={value}") + if len(items) >= 3: + break + if not items: + return None + return ", ".join(items) + + +def _event_label(event_type: str) -> str: + return EVENT_TYPE_LABELS.get(event_type, event_type) + + +# Отступ для строки события под статусом (чтобы не слипалось) +_FLOW_INDENT = " " + + +def _humanize_path(path: str) -> str: + """Сокращает типичные callback для админки до читаемого вида.""" + if not path or "callback:" not in path: + return path + # callback:admin_users:users_audit:476217106:telegram | all | 0:0 + if "users_audit:" in path: + rest = path.split("users_audit:", 1)[-1].strip() + parts = [p.strip() for p in rest.split("|")[:2] if p.strip()] + if len(parts) >= 2: + ch, cat = parts[0].split(":")[-1] if ":" in parts[0] else parts[0], parts[1] + return f"История: {ch}, {cat}" + return "История" + if "users_editor:" in path: + return "Карточка пользователя" + if "admin_panel:search_user" in path: + return "Поиск пользователя" + if "admin_panel:admin" in path: + return "Админ-панель" + if "message:" in path: + text = path.split("message:", 1)[-1].strip() + if text and text != "-": + return f"Сообщение: {text[:40]}{'…' if len(text) > 40 else ''}" + return path[:55] + ("…" if len(path) > 55 else "") + + +def _format_event_line( + event, show_request_id: bool = False, inside_block: bool = False, skip_time: bool = False +) -> str: + """Строка события. Если skip_time=True — только описание (время уже в строке статуса).""" + created_at = event.created_at.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%d.%m %H:%M:%S") + if event.event_type == "telegram_access": + raw_path = (event.path_or_handler or "—").strip() + event_name = html.escape(_humanize_path(raw_path)) + if event_name == raw_path and len(raw_path) > 55: + event_name = html.escape(raw_path[:55] + "…") + else: + event_name = html.escape(_event_label(event.event_type)) + req_part = "" + if show_request_id and not inside_block and getattr(event, "request_id", None): + short_id = (event.request_id or "")[:8] + if short_id: + req_part = f" [{short_id}]" + entity = "" + if event.entity_type or event.entity_id: + etype = html.escape(str(event.entity_type or "entity")) + raw_id = str(event.entity_id or "") + eid = html.escape(raw_id[:40] + ("…" if len(raw_id) > 40 else "")) + if event.entity_type == "telegram_user" and raw_id.isdigit(): + entity = f"\n{_FLOW_INDENT}tg_id: {eid}" + else: + entity = f"\n{_FLOW_INDENT}{etype}: {eid}" + + metadata_line = _format_metadata(event.metadata_) + if metadata_line: + metadata_line = f"\n{_FLOW_INDENT}{html.escape(metadata_line)}" + else: + metadata_line = "" + + reason_line = "" + if event.reason: + reason_line = f"\n{_FLOW_INDENT}{html.escape(str(event.reason)[:100])}" + + time_part = "" if skip_time else f"{created_at} " + return f"{time_part}{event_name}{req_part}{entity}{metadata_line}{reason_line}" + + +def _format_event_status(event) -> str: + """Только время и результат (ок/ошибка).""" + created_at = event.created_at.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%d.%m %H:%M:%S") + result_text = "ок" if event.result == "success" else "ошибка" + return f"{created_at} {result_text}" + + +# Разделитель между событиями — сразу видно границу "что где" +_FLOW_SEP = "—" + +def _render_events_as_flow(events: list) -> list[str]: + """Флоу: успешные действия (ок) тянутся через ⤷; если не ок — отдельный блок с │. Между блоками — разделитель.""" + if not events: + return [] + lines = [] + for i, event in enumerate(events): + n = i + 1 + is_ok = event.result == "success" + prefix = "⤷ " if is_ok else "│ " + if not is_ok and i > 0: + lines.append(_FLOW_SEP) + status_line = f"{prefix}{n}. {_format_event_status(event)}" + body = _format_event_line(event, show_request_id=False, inside_block=True, skip_time=True) + quoted_body = f"
{body}
" + lines.append(status_line) + lines.append(quoted_body) + if i < len(events) - 1: + lines.append(_FLOW_SEP) + return lines + + +async def _render_user_audit( + message: Message, + session: AsyncSession, + tg_id: int, + *, + channel_filter: str = "all", + category_filter: str = "all", + page: int = 0, +) -> None: + page = max(0, page) + user_identity_id = await session.scalar(select(User.identity_id).where(User.tg_id == tg_id)) + channel = None if channel_filter == "all" else channel_filter + event_types = _resolve_event_types(category_filter) + + cache_key_str = cache_key( + "audit_history", + tg_id, + user_identity_id or "", + channel_filter, + category_filter, + page, + ) + cached = await cache_get(cache_key_str) + if cached is not None and isinstance(cached, list): + raw = _deserialize_audit_events(cached) + has_prev = page > 0 + has_next = len(raw) > PAGE_SIZE + events = raw[:PAGE_SIZE] + else: + raw = await list_audit_events( + session, + tg_id=tg_id, + identity_id=user_identity_id, + channel=channel, + event_types=event_types, + limit=PAGE_SIZE + 1, + offset=page * PAGE_SIZE, + ) + has_prev = page > 0 + has_next = len(raw) > PAGE_SIZE + events = raw[:PAGE_SIZE] + if raw: + await cache_set( + cache_key_str, + _serialize_audit_events(raw), + AUDIT_HISTORY_CACHE_TTL_SEC, + ) + + full_flow = channel_filter == "all" and category_filter == "all" + lines = [f"🕘 История действий клиента {tg_id}"] + if full_flow: + lines.append("📋 Вся хронология. Цифра — номер действия, черта — граница между действиями.") + lines.append( + f"📎 Канал: {html.escape(channel_filter)} | Категория: {html.escape(category_filter)}" + ) + if user_identity_id: + lines.append(f"🆔 Identity: {html.escape(user_identity_id)}") + + if not events: + lines.append("\nСобытий пока нет.") + else: + lines.append("") + rev = list(reversed(events)) # хронология сверху вниз + if full_flow: + by_cat: dict[str, list] = {} + for e in rev: + c = _event_category(e) + by_cat.setdefault(c, []).append(e) + visible_cats = [c for c in CATEGORY_BLOCK_ORDER if c in by_cat] + for idx, cat in enumerate(visible_cats): + lines.append(f"▸ {CATEGORY_BLOCK_LABELS[cat]}") + lines.extend(_render_events_as_flow(by_cat[cat])) + if idx < len(visible_cats) - 1: + lines.append(_FLOW_SEP) + else: + lines.extend(_render_events_as_flow(rev)) + + await message.edit_text( + text="\n".join(lines).strip(), + reply_markup=build_user_audit_kb( + tg_id=tg_id, + channel_filter=channel_filter, + category_filter=category_filter, + page=page, + has_prev=has_prev, + has_next=has_next, + ), + disable_web_page_preview=True, + ) + + +@router.callback_query( + AdminUserEditorCallback.filter(F.action == "users_audit"), + IsAdminFilter(), +) +async def handle_user_audit( + callback_query: CallbackQuery, + callback_data: AdminUserEditorCallback, + session: AsyncSession, +): + channel_filter, category_filter, page = _parse_filter_page(callback_data.data) + await _render_user_audit( + callback_query.message, + session, + callback_data.tg_id, + channel_filter=channel_filter, + category_filter=category_filter, + page=page, + ) + + +@router.callback_query( + AdminUserEditorCallback.filter(F.action == "users_audit_page"), + IsAdminFilter(), +) +async def handle_user_audit_page( + callback_query: CallbackQuery, + callback_data: AdminUserEditorCallback, + session: AsyncSession, +): + channel_filter, category_filter, page = _parse_filter_page(callback_data.data) + await _render_user_audit( + callback_query.message, + session, + callback_data.tg_id, + channel_filter=channel_filter, + category_filter=category_filter, + page=page, + ) diff --git a/handlers/admin/users/users_balance.py b/handlers/admin/users/users_balance.py index efd0e6bc..21cc5ce7 100644 --- a/handlers/admin/users/users_balance.py +++ b/handlers/admin/users/users_balance.py @@ -280,4 +280,5 @@ async def handle_balance_input(message: Message, state: FSMContext, session: Asy status="success", ) + await state.clear() await message.answer(text=text, reply_markup=build_users_balance_change_kb(tg_id)) diff --git a/handlers/admin/users/users_manage.py b/handlers/admin/users/users_manage.py index a77faefd..7149aa8d 100644 --- a/handlers/admin/users/users_manage.py +++ b/handlers/admin/users/users_manage.py @@ -18,7 +18,7 @@ from database import ( get_key_details, update_trial, ) -from database.models import Key, ManualBan, Payment, Referral, User +from database.models import Admin, Key, ManualBan, Payment, Referral, User from filters.admin import IsAdminFilter from handlers.utils import sanitize_key_name from utils.csv_export import export_referrals_csv @@ -88,7 +88,7 @@ async def handle_key_name_input(message: Message, state: FSMContext, session: As ) return - await process_user_search(message, state, session, key_details["tg_id"]) + await process_user_search(message, state, session, key_details["tg_id"], actor_tg_id=message.from_user.id) @router.message(UserEditorState.waiting_for_user_data, IsAdminFilter()) @@ -97,7 +97,7 @@ async def handle_user_data_input(message: Message, state: FSMContext, session: A if message.forward_from: tg_id = message.forward_from.id - await process_user_search(message, state, session, tg_id) + await process_user_search(message, state, session, tg_id, actor_tg_id=message.from_user.id) return if not message.text: @@ -124,7 +124,7 @@ async def handle_user_data_input(message: Message, state: FSMContext, session: A ) return - await process_user_search(message, state, session, tg_id) + await process_user_search(message, state, session, tg_id, actor_tg_id=message.from_user.id) @router.callback_query( @@ -336,6 +336,7 @@ async def process_user_search( session: AsyncSession, tg_id: int, edit: bool = False, + actor_tg_id: int | None = None, ) -> None: await state.clear() @@ -435,7 +436,12 @@ async def process_user_search( text = text_builder.as_html() - kb = await build_user_edit_kb(tg_id, key_records, is_banned=is_banned) + effective_actor_tg_id = actor_tg_id or (message.from_user.id if message.from_user else None) + admin_role = None + if effective_actor_tg_id is not None: + admin_role = await session.scalar(select(Admin.role).where(Admin.tg_id == effective_actor_tg_id)) + + kb = await build_user_edit_kb(tg_id, key_records, is_banned=is_banned, admin_role=admin_role) if edit: try: @@ -462,4 +468,5 @@ async def handle_users_editor( session=session, tg_id=callback_data.tg_id, edit=callback_data.edit, + actor_tg_id=callback.from_user.id, ) diff --git a/handlers/payments/pay.py b/handlers/payments/pay.py index b6abeddc..33c93196 100644 --- a/handlers/payments/pay.py +++ b/handlers/payments/pay.py @@ -151,7 +151,6 @@ async def _build_pay_menu_for_currency(currency: str) -> InlineKeyboardBuilder: @router.callback_query(F.data.startswith("pay_currency|")) async def handle_pay_currency(callback_query: CallbackQuery, state: FSMContext, session: AsyncSession): currency = callback_query.data.split("|")[1] - if currency == "STARS": return await process_callback_pay_stars(callback_query, state, session) diff --git a/handlers/refferal.py b/handlers/refferal.py index 475ef29c..9e4bffb8 100644 --- a/handlers/refferal.py +++ b/handlers/refferal.py @@ -190,7 +190,6 @@ async def show_referral_qr(callback_query: CallbackQuery): reply_markup=builder.as_markup(), media_path=qr_path, ) - os.remove(qr_path) except Exception as error: diff --git a/handlers/start.py b/handlers/start.py index 6a7c31c1..657223ab 100644 --- a/handlers/start.py +++ b/handlers/start.py @@ -26,7 +26,6 @@ from core.cache_config import START_UTM_EXISTS_TTL_SEC from core.redis_cache import cache_get, cache_key, cache_set from database import ( add_user, - get_coupon_by_code, get_user_snapshot, upsert_source_if_empty, ) @@ -94,7 +93,6 @@ async def start_entry( captcha: bool = True, ): message = event.message if isinstance(event, CallbackQuery) else event - try: await run_hooks("start_entry", message=message, event=event, state=state, session=session, admin=admin) except Exception as e: diff --git a/middlewares/__init__.py b/middlewares/__init__.py index 2962560f..c2a1ee25 100644 --- a/middlewares/__init__.py +++ b/middlewares/__init__.py @@ -60,7 +60,6 @@ def register_middleware( if PROBE_LOGGING: dispatcher.update.outer_middleware(StreamProbeMiddleware("global")) - # Первым делом отвечаем на callback, чтобы не уйти в «query is too old» при очереди dispatcher.update.outer_middleware(EarlyCallbackAnswerMiddleware()) if middleware_enabled("runtime_config_sync"): @@ -81,7 +80,7 @@ def register_middleware( available_middlewares = { "admin": AdminMiddleware(), "maintenance": MaintenanceModeMiddleware(), - "logging": LoggingMiddleware(), + "logging": LoggingMiddleware(sessionmaker) if sessionmaker else LoggingMiddleware(), "throttling": ThrottlingMiddleware(), "user": UserMiddleware(), "answer": CallbackAnswerMiddleware(), diff --git a/middlewares/loggings.py b/middlewares/loggings.py index cfd60b99..f7d42030 100644 --- a/middlewares/loggings.py +++ b/middlewares/loggings.py @@ -5,6 +5,14 @@ from typing import Any, TypedDict from aiogram import BaseMiddleware from aiogram.types import CallbackQuery, InlineQuery, Message, TelegramObject, User +from audit import ( + ensure_telegram_context, + log_telegram_access, + record_telegram_access_event, + record_telegram_access_event_background, + set_telegram_actor, + _telegram_access_payload, +) from logger import logger @@ -25,7 +33,12 @@ def _log_activity_sync(user_info: UserInfo) -> None: class LoggingMiddleware(BaseMiddleware): - """Middleware для логирования действий пользователя. Лог пишется в фоне, не задерживая обработчик.""" + """Middleware для логирования действий пользователя. Лог пишется в фоне, не задерживая обработчик. + Если передан sessionmaker — аудит пишется в отдельной сессии в фоне (не блокирует ответ).""" + + def __init__(self, sessionmaker=None): + super().__init__() + self._sessionmaker = sessionmaker async def __call__( self, @@ -33,12 +46,75 @@ class LoggingMiddleware(BaseMiddleware): event: TelegramObject, data: dict[str, Any], ) -> Any: + audit_context = ensure_telegram_context(data, event) user_info = self._extract_user_info(event) if user_info["user_id"]: asyncio.create_task(asyncio.to_thread(_log_activity_sync, user_info)) - return await handler(event, data) + try: + result = await handler(event, data) + db_user = data.get("user") + if isinstance(db_user, dict): + set_telegram_actor( + audit_context, + identity_id=db_user.get("identity_id"), + tg_id=db_user.get("tg_id"), + ) + if self._sessionmaker is not None: + payload = _telegram_access_payload(audit_context, event, result="success") + asyncio.create_task( + record_telegram_access_event_background(self._sessionmaker, **payload) + ) + else: + session = data.get("session") + if session is not None: + await record_telegram_access_event( + session, + audit_context, + event, + result="success", + ) + asyncio.create_task( + asyncio.to_thread( + log_telegram_access, + event, + audit_context=audit_context, + result="success", + ) + ) + return result + except Exception as exc: + reason = type(exc).__name__ + if self._sessionmaker is not None: + payload = _telegram_access_payload( + audit_context, event, result="fail", reason=reason + ) + asyncio.create_task( + record_telegram_access_event_background( + self._sessionmaker, **payload + ) + ) + else: + session = data.get("session") + if session is not None: + await record_telegram_access_event( + session, + audit_context, + event, + result="fail", + reason=reason, + ) + asyncio.create_task( + asyncio.to_thread( + log_telegram_access, + event, + audit_context=audit_context, + result="fail", + reason=reason, + ) + ) + raise def _extract_user_info(self, event: TelegramObject) -> UserInfo: """Извлекает информацию о пользователе из различных типов событий.""" diff --git a/requirements.txt b/requirements.txt index 675cb6a3..88839962 100644 --- a/requirements.txt +++ b/requirements.txt @@ -48,6 +48,7 @@ pydantic_core==2.23.4 Pygments==2.19.2 python-dateutil==2.9.0.post0 pytz==2025.1 +python-multipart==0.0.22 qrcode==8.2 requests==2.32.4 redis==5.2.0