diff --git a/audit.py b/audit/__init__.py similarity index 58% rename from audit.py rename to audit/__init__.py index a750913d..7e288b32 100644 --- a/audit.py +++ b/audit/__init__.py @@ -4,35 +4,36 @@ import json import uuid from dataclasses import dataclass -from datetime import datetime, timedelta, timezone +from datetime import datetime, 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.audit import ( + create_audit_reset_marker_db, + delete_old_audit_events_db, + ensure_audit_table, + fetch_existing_audit_request_ids_db, + fetch_latest_audit_reset_db, + fetch_audit_events_db, + fetch_audit_events_db_window, + fetch_audit_rows_db, + fetch_successful_payment_rows_db, +) from database.models import AuditEvent from logger import logger -def _get_bot_webhook_path() -> str: - """Путь вебхука бота из конфига (для исключения из шага «успешная оплата»).""" - try: - from config import WEBHOOK_PATH - return ((WEBHOOK_PATH or "").strip().lower()) or "" - except ImportError: - return "" - - -def _is_bot_webhook_path(path: str) -> bool: - """True только если path — именно вебхук бота (точное совпадение сегмента пути), не касса.""" - bot_path = _get_bot_webhook_path() - if not bot_path: - return False - p = (path or "").strip().lower() - path_segment = p.split(" ", 1)[1] if " " in p else p - return path_segment == bot_path or path_segment.rstrip("/") == bot_path.rstrip("/") +from .rules import ( + AUDIT_STEP_LABELS, + DEFAULT_FUNNEL_STEPS, + _funnel_step_counts, + _is_ignored_analytics_event, + _normalize_path_to_step, + _normalize_path_to_steps, +) try: from core.cache_config import ( @@ -51,10 +52,11 @@ except ImportError: AUDIT_REDIS_USER_TTL_SEC = 25 * 3600 AUDIT_REDIS_DRAIN_BATCH = 1000 + _MAX_TEXT_LEN = 160 -_AUDIT_TABLE_READY = False - - +_AUDIT_REDIS_PROCESSING_KEY = f"{AUDIT_REDIS_FLUSH_KEY}:processing" +_AUDIT_REDIS_DRAIN_LOCK_KEY = f"{AUDIT_REDIS_FLUSH_KEY}:drain_lock" +_AUDIT_REDIS_DRAIN_LOCK_TTL_SEC = 15 * 60 @dataclass class AuditContext: request_id: str @@ -102,6 +104,36 @@ def _serialize(payload: dict[str, Any]) -> str: return json.dumps(_jsonable(payload), ensure_ascii=False, sort_keys=True) +async def get_audit_db_reset_at(session: AsyncSession) -> datetime | None: + reset_at = await fetch_latest_audit_reset_db(session, source="db") + if reset_at is None: + return None + if reset_at.tzinfo is None: + return reset_at.replace(tzinfo=timezone.utc) + return reset_at.astimezone(timezone.utc) + + +async def set_audit_db_reset_at(session: AsyncSession, at: datetime | None = None) -> datetime: + created = at.astimezone(timezone.utc).replace(tzinfo=None) if at is not None and at.tzinfo else (at or datetime.utcnow()) + await create_audit_reset_marker_db(session, source="db", created_at=created) + await session.commit() + return created.replace(tzinfo=timezone.utc) + + +async def clear_audit_redis_buffers() -> int: + from core.redis_cache import cache_delete_pattern + + patterns = ( + AUDIT_REDIS_FLUSH_KEY, + _AUDIT_REDIS_PROCESSING_KEY, + _AUDIT_REDIS_DRAIN_LOCK_KEY, + ) + deleted = 0 + for pattern in patterns: + deleted += await cache_delete_pattern(pattern) + return deleted + + def _message_text(event: TelegramObject) -> str | None: if isinstance(event, Message): return _trim(event.text or event.caption) @@ -346,7 +378,6 @@ def _telegram_access_payload( 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) @@ -375,7 +406,6 @@ async def record_telegram_access_event( 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, @@ -404,7 +434,7 @@ def _audit_record_for_redis( request_id: str | None = None, metadata_: dict | None = None, ) -> dict[str, Any]: - """Формирует запись для буфера Redis (с created_at в ISO).""" + effective_request_id = _trim(request_id, 64) if request_id else new_request_id() return { "event_type": event_type, "channel": channel, @@ -415,12 +445,31 @@ def _audit_record_for_redis( "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, + "request_id": effective_request_id, "metadata_": _jsonable(metadata_) if metadata_ else None, "created_at": datetime.now(timezone.utc).isoformat(), } +async def _push_audit_record_to_redis(record: dict[str, Any]) -> None: + from core.redis_cache import cache_expire, cache_rpush + + n = await cache_rpush(AUDIT_REDIS_FLUSH_KEY, record) + if n == 0: + raise RuntimeError("Redis unavailable (cache_rpush returned 0)") + + actor_tg_id = record.get("actor_tg_id") + actor_identity_id = record.get("actor_identity_id") + 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_audit_event_to_redis( *, request_id: str | None = None, @@ -430,9 +479,6 @@ async def record_audit_event_to_redis( 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, @@ -441,17 +487,7 @@ async def record_audit_event_to_redis( reason=reason, request_id=request_id, ) - n = await cache_rpush(AUDIT_REDIS_FLUSH_KEY, record) - if n == 0: - raise RuntimeError("Redis unavailable (cache_rpush returned 0)") - 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) + await _push_audit_record_to_redis(record) async def record_api_access_event_to_redis( @@ -463,9 +499,6 @@ async def record_api_access_event_to_redis( 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", @@ -476,15 +509,7 @@ async def record_api_access_event_to_redis( 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) + await _push_audit_record_to_redis(record) async def record_api_access_event_background( @@ -495,8 +520,6 @@ async def record_api_access_event_background( 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: @@ -548,8 +571,6 @@ async def record_telegram_access_event_background( result: str = "success", reason: str | None = None, ) -> None: - """Пишет событие telegram_access: в Redis-буфер (если включён), иначе в БД. - При сбое записи в Redis — fallback в БД, чтобы события не терялись.""" if AUDIT_REDIS_BUFFER_ENABLED: try: await record_audit_event_to_redis( @@ -562,9 +583,7 @@ async def record_telegram_access_event_background( ) return except Exception as exc: - logger.warning( - f"[Audit] Запись в Redis-буфер не удалась, пишем в БД: {exc}" - ) + logger.warning(f"[Audit] Запись в Redis-буфер не удалась, пишем в БД: {exc}") try: async with session_factory() as session: await ensure_audit_table(session) @@ -651,7 +670,6 @@ async def safe_record_telegram_event( def _redis_record_to_event_like(rec: dict[str, Any]) -> SimpleNamespace: - """Превращает запись из Redis в объект с теми же атрибутами, что и AuditEvent.""" created = rec.get("created_at") if isinstance(created, str): try: @@ -684,7 +702,6 @@ async def _list_audit_events_from_redis( 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 @@ -716,7 +733,6 @@ async def _list_audit_events_from_redis( async def list_audit_events_from_redis_buffer(max_events: int = 5000) -> list[SimpleNamespace]: - """Читает последние события из глобального буфера Redis (audit:flush). Не удаляет записи.""" if not AUDIT_REDIS_BUFFER_ENABLED: return [] from core.redis_cache import cache_lrange @@ -734,26 +750,26 @@ async def list_audit_events_from_redis_buffer(max_events: int = 5000) -> list[Si def _aggregate_audit_rows( rows: list[tuple[Any, Any, Any, Any]], ) -> tuple[dict[str, dict[str, Any]], list[dict[str, Any]], set[tuple[str, int | str]]]: - """Собирает by_step, by_path_list, all_actors из строк (path, result, actor_tg_id, actor_identity_id).""" by_step: dict[str, dict[str, Any]] = {} - total_events = 0 all_actors: set[tuple[str, int | str]] = set() - for path, res, tg_id, identity_id in rows: - total_events += 1 - step = _normalize_path_to_step(path or "") - actor = (identity_id or "", tg_id or 0) - if identity_id or tg_id: + for row in rows: + path, res, tg_id, identity_id = row[0], row[1], row[2], row[3] + if _is_ignored_analytics_event(path or ""): + continue + actor = _audit_actor_key(identity_id, tg_id) + if actor is not None: all_actors.add(actor) - if step not in by_step: - by_step[step] = {"total": 0, "success": 0, "fail": 0, "actors": set()} - by_step[step]["total"] += 1 - if res == "success": - by_step[step]["success"] += 1 - else: - by_step[step]["fail"] += 1 - if identity_id or tg_id: - by_step[step]["actors"].add(actor) + for step in _normalize_path_to_steps(path or ""): + if step not in by_step: + by_step[step] = {"total": 0, "success": 0, "fail": 0, "actors": set()} + by_step[step]["total"] += 1 + if res == "success": + by_step[step]["success"] += 1 + else: + by_step[step]["fail"] += 1 + if actor is not None: + by_step[step]["actors"].add(actor) by_path_list = [] for step, data in sorted(by_step.items(), key=lambda x: -x[1]["total"]): @@ -761,18 +777,63 @@ def _aggregate_audit_rows( fail = data["fail"] unique = len(data["actors"]) fail_rate = round(100.0 * fail / total, 1) if total else 0 - by_path_list.append({ - "step": step, - "label": AUDIT_STEP_LABELS.get(step, step), - "total": total, - "success": data["success"], - "fail": fail, - "unique_users": unique, - "fail_rate_pct": fail_rate, - }) + by_path_list.append( + { + "step": step, + "label": AUDIT_STEP_LABELS.get(step, step), + "total": total, + "success": data["success"], + "fail": fail, + "unique_users": unique, + "fail_rate_pct": fail_rate, + } + ) return by_step, by_path_list, all_actors +def _audit_actor_key(identity_id: Any, tg_id: Any) -> tuple[str, str | int] | None: + if tg_id not in (None, 0, "0", ""): + try: + return ("tg", int(tg_id)) + except (TypeError, ValueError): + return ("tg", str(tg_id)) + if identity_id not in (None, ""): + return ("identity", str(identity_id)) + return None + + +def _event_like_dedupe_key(event: AuditEvent | SimpleNamespace) -> tuple[Any, ...]: + request_id = getattr(event, "request_id", None) + event_type = getattr(event, "event_type", None) + path = getattr(event, "path_or_handler", None) + result = getattr(event, "result", None) + if request_id: + return ("request", request_id, event_type, path, result) + return ( + "raw", + event_type, + getattr(event, "channel", None), + path, + getattr(event, "actor_tg_id", None), + getattr(event, "actor_identity_id", None), + result, + getattr(event, "reason", None), + str(getattr(event, "created_at", None)), + ) + + +def _dedupe_event_like(events: list[AuditEvent | SimpleNamespace]) -> list[AuditEvent | SimpleNamespace]: + seen: set[tuple[Any, ...]] = set() + unique: list[AuditEvent | SimpleNamespace] = [] + for event in events: + key = _event_like_dedupe_key(event) + if key in seen: + continue + seen.add(key) + unique.append(event) + return unique + + async def list_audit_events( session: AsyncSession, *, @@ -784,262 +845,44 @@ async def list_audit_events( 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()) + return await fetch_audit_events_db( + session, + identity_id=identity_id, + tg_id=tg_id, + channel=channel, + event_type=event_type, + event_types=event_types_list, + limit=limit, + offset=offset, + ) 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 + db_events = await fetch_audit_events_db_window( + session, + identity_id=identity_id, + tg_id=tg_id, + channel=channel, + event_type=event_type, + event_types=event_types_list, + limit=min(5000, need), + ) + merged = _dedupe_event_like(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 - - -AUDIT_STEP_LABELS: dict[str, str] = { - "start": "Старт", - "profile": "Профиль", - "view_keys": "Мои ключи", - "key_create": "Подписка оформлена", - "pay_start": "Начало оплаты (ссылка создана)", - "pay": "Успешная оплата (пополнение)", - "key_view": "Ключ (карточка)", - "connect": "Подписка подключена", - "renew": "Продление", - "addons": "Аддоны", - "referral": "Рефералы", - "coupons": "Купоны", - "register": "Регистрация (API)", - "login": "Вход (API)", - "api_other": "API прочее", - "admin": "Админ-панель", - "other": "Прочее", -} - - -DEFAULT_FUNNEL_STEPS = ("start", "profile", "view_keys", "key_create", "pay_start", "pay", "key_view", "connect") - - -_CALLBACK_EXACT: dict[str, str] = { - "profile": "profile", - "view_keys": "view_keys", - "create_key": "key_create", - "buy": "key_create", - "pay": "pay_start", - "balance": "pay_start", -} - - -_CALLBACK_PREFIX: list[tuple[str, str]] = [ - ("view_key|", "key_view"), - ("view_keys|", "view_keys"), - ("connect_device|", "connect"), - ("connect_router|", "connect"), - ("connect_tv|", "connect"), - ("connect_pc|", "connect"), - ("connect_ios|", "connect"), - ("connect_android|", "connect"), - ("show_qr|", "connect"), - ("continue_tv|", "connect"), - ("pay_currency", "pay_start"), - ("balance_history", "pay_start"), - ("cfg_user_confirm|", "key_create"), - ("choose_payment_provider|", "pay_start"), - ("cfg_renew", "renew"), - ("key_addons", "addons"), - ("extend_key", "coupons"), -] - - -_CALLBACK_CONTAINS: list[tuple[str, str]] = [ - ("connect_", "connect"), - ("renew", "renew"), - ("addon", "addons"), - ("referral", "referral"), - ("invite", "referral"), - ("coupon", "coupons"), - ("users_audit", "admin"), - ("users_editor", "admin"), - ("search_user", "admin"), - ("admin_panel", "admin"), -] - -_MESSAGE_START: set[str] = {"/start", "start"} - -_HANDLER_CONTAINS: list[tuple[str, str] | tuple[str, str, str]] = [ - ("process_start", "start"), - ("start_entry", "start"), - ("show_start_menu", "start"), - ("process_callback_view_profile", "profile"), - ("process_callback_or_message_view_keys", "view_keys"), - ("key_view", "key_view", "key_create"), - ("key_create", "key_create"), - ("handle_key_creation", "key_create"), - ("confirm_create", "key_create"), - ("complete_key_renewal", "renew"), - ("handle_connect_device", "connect"), - ("process_connect_", "connect"), - ("process_callback_connect", "connect"), - ("process_continue_tv", "connect"), - ("show_qr_code", "connect"), - ("pay", "pay_start"), - ("balance", "pay_start"), - ("renew", "renew"), - ("addon", "addons"), - ("referral", "referral"), - ("refferal", "referral"), - ("coupon", "coupons"), - ("admin_panel", "admin"), - ("users_audit", "admin"), - ("users_editor", "admin"), - ("search_user", "admin"), - ("auth/register", "register"), - ("auth/login", "login"), - ("auth/send-login", "login"), - ("auth/login-by-code", "login"), - ("auth/login-telegram", "login"), -] - -def _funnel_step_counts(path: str, result: str, step: str) -> bool: - """Решает, считать ли событие достижением шага воронки. - pay_start: любое успешное «начало оплаты» (меню, валюта, создание ссылки). - pay: вебхук с «webhook» в path, кроме пути вебхука бота из конфига (WEBHOOK_PATH). - key_create: только факт создания ключа (cfg_user_confirm или /keys/create). - connect: любое успешное подключение.""" - if result != "success": - return False - p = (path or "").lower() - if step == "pay_start": - return True - if step == "pay": - if "webhook" not in p: - return False - if _is_bot_webhook_path(path or ""): - return False - return True - if step == "key_create": - return "cfg_user_confirm" in p or "/keys/create" in p - return True - - -def _normalize_path_to_step(path: str) -> str: - """Сводит path_or_handler к шагу по правилам из маппингов. - Вебхук с «webhook» в path → pay, кроме пути вебхука бота из конфига (WEBHOOK_PATH).""" - if not path: - return "other" - p = path.lower().strip() - if "webhook" in p: - if _is_bot_webhook_path(path): - pass - else: - return "pay" - - if p.startswith("callback:"): - callback_data = p.split(":", 1)[-1] - data = callback_data.split("|")[0] - step = _CALLBACK_EXACT.get(data) - if step: - return step - for prefix, step in _CALLBACK_PREFIX: - if callback_data.startswith(prefix): - return step - for item in _CALLBACK_CONTAINS: - substr, step = item[0], item[1] - if substr in callback_data: - return step - return "other" - - if p.startswith("message:"): - text = (p.split(":", 1)[-1] or "").strip() - return "start" if any(s in text for s in _MESSAGE_START) else "other" - - if p.startswith("post ") or p.startswith("get "): - if "auth/register" in p: - return "register" - if "auth/" in p and "login" in p: - return "login" - if "payment-links" in p or ("payment" in p and "webhook" not in p): - return "pay_start" - return "api_other" - - for item in _HANDLER_CONTAINS: - if len(item) == 3: - substr, step, exclude = item[0], item[1], item[2] - if substr in p and exclude not in p: - return step - else: - substr, step = item[0], item[1] - if substr in p: - return step - return "other" + return await delete_old_audit_events_db(session, older_than_days=older_than_days) async def get_audit_stats( @@ -1049,32 +892,20 @@ async def get_audit_stats( date_to: datetime, max_events: int = 100_000, ) -> dict[str, Any]: - """Агрегаты по аудиту за период (только БД): по шагам объём, успехи/ошибки, уникальные пользователи. - Удобно для «какие пути хорошо отрабатывают, какие нет». Данные из Redis в расчёт не берутся.""" - await ensure_audit_table(session) d_from = _naive_utc(date_from) d_to = _naive_utc(date_to) - stmt = ( - select( - AuditEvent.path_or_handler, - AuditEvent.result, - AuditEvent.actor_tg_id, - AuditEvent.actor_identity_id, - ) - .where( - AuditEvent.created_at >= d_from, - AuditEvent.created_at < d_to, - ) - .limit(max_events) - ) - result = await session.execute(stmt) - rows = result.all() + rows = await fetch_audit_rows_db(session, date_from=d_from, date_to=d_to, limit=max_events) + rows.extend(await fetch_successful_payment_rows_db(session, date_from=d_from, date_to=d_to, limit=max_events)) by_step, by_path_list, all_actors = _aggregate_audit_rows(rows) + raw_total_events = sum(1 for row in rows if not _is_ignored_analytics_event(row[0] or "")) + analytics_total_events = sum(row["total"] for row in by_path_list) return { "summary": { "date_from": date_from.isoformat(), "date_to": date_to.isoformat(), - "total_events": sum(d["total"] for d in by_step.values()), + "total_events": raw_total_events, + "raw_total_events": raw_total_events, + "analytics_total_events": analytics_total_events, "unique_users": len(all_actors), }, "by_path": by_path_list, @@ -1082,16 +913,48 @@ async def get_audit_stats( async def get_audit_stats_from_redis(max_events: int = 5000) -> dict[str, Any] | None: - """Агрегаты по аудиту из буфера Redis (без БД). Возвращает None, если буфер выключен.""" if not AUDIT_REDIS_BUFFER_ENABLED: return None events = await list_audit_events_from_redis_buffer(max_events=max_events) - rows = [(e.path_or_handler, e.result, e.actor_tg_id, e.actor_identity_id) for e in events] + filtered_events = events + rows = [(e.path_or_handler, e.result, e.actor_tg_id, e.actor_identity_id) for e in filtered_events] by_step, by_path_list, all_actors = _aggregate_audit_rows(rows) + raw_total_events = sum(1 for row in rows if not _is_ignored_analytics_event(row[0] or "")) + analytics_total_events = sum(row["total"] for row in by_path_list) return { "summary": { "source": "redis", - "total_events": len(events), + "total_events": raw_total_events, + "raw_total_events": raw_total_events, + "analytics_total_events": analytics_total_events, + "unique_users": len(all_actors), + }, + "by_path": by_path_list, + } + + +async def get_audit_stats_from_redis_since( + *, + date_from: datetime | None = None, + max_events: int = 5000, +) -> dict[str, Any] | None: + if not AUDIT_REDIS_BUFFER_ENABLED: + return None + events = await list_audit_events_from_redis_buffer(max_events=max_events) + filtered_events = events + if date_from is not None: + threshold = date_from.astimezone(timezone.utc) if date_from.tzinfo else date_from.replace(tzinfo=timezone.utc) + filtered_events = [e for e in events if getattr(e, "created_at", None) and e.created_at >= threshold] + rows = [(e.path_or_handler, e.result, e.actor_tg_id, e.actor_identity_id) for e in filtered_events] + by_step, by_path_list, all_actors = _aggregate_audit_rows(rows) + raw_total_events = sum(1 for row in rows if not _is_ignored_analytics_event(row[0] or "")) + analytics_total_events = sum(row["total"] for row in by_path_list) + return { + "summary": { + "source": "redis", + "total_events": raw_total_events, + "raw_total_events": raw_total_events, + "analytics_total_events": analytics_total_events, "unique_users": len(all_actors), }, "by_path": by_path_list, @@ -1106,27 +969,11 @@ async def get_audit_funnel( steps_ordered: tuple[str, ...] | None = None, max_events: int = 50_000, ) -> list[dict[str, Any]]: - """Воронка: сколько уникальных пользователей достигли каждого шага. - Шаг pay = только вебхук кассы (пополнение). Шаг key_create = только факт создания ключа. Остальное — по result success.""" - await ensure_audit_table(session) steps = steps_ordered or DEFAULT_FUNNEL_STEPS d_from = _naive_utc(date_from) d_to = _naive_utc(date_to) - stmt = ( - select( - AuditEvent.path_or_handler, - AuditEvent.result, - AuditEvent.actor_tg_id, - AuditEvent.actor_identity_id, - ) - .where( - AuditEvent.created_at >= d_from, - AuditEvent.created_at < d_to, - ) - .limit(max_events) - ) - result = await session.execute(stmt) - rows = result.all() + rows = await fetch_audit_rows_db(session, date_from=d_from, date_to=d_to, limit=max_events) + rows.extend(await fetch_successful_payment_rows_db(session, date_from=d_from, date_to=d_to, limit=max_events)) return _funnel_from_rows(rows, steps) @@ -1134,47 +981,43 @@ def _funnel_from_rows( rows: list[tuple[Any, ...]], steps_ordered: tuple[str, ...] | None = None, ) -> list[dict[str, Any]]: - """Воронка по строкам: (path, result, actor_tg_id, actor_identity_id). Учитываются только завершающие события.""" steps = steps_ordered or DEFAULT_FUNNEL_STEPS - actor_steps: dict[tuple[str, int], set[str]] = {} + step_actors: dict[str, set[tuple[str, str | int]]] = {step: set() for step in steps} for row in rows: if len(row) >= 4: path, result, tg_id, identity_id = row[0], row[1], row[2], row[3] else: path, tg_id, identity_id = row[0], row[1], row[2] result = "success" - step = _normalize_path_to_step(path or "") - if not _funnel_step_counts(path or "", str(result or ""), step): + if _is_ignored_analytics_event(path or ""): continue - key = (str(identity_id or ""), int(tg_id or 0)) - if key == ("", 0): + key = _audit_actor_key(identity_id, tg_id) + if key is None: continue - if key not in actor_steps: - actor_steps[key] = set() - actor_steps[key].add(step) - - step_index = {s: i for i, s in enumerate(steps)} - reached: list[int] = [0] * len(steps) - for _actor, reached_steps in actor_steps.items(): - max_idx = -1 - for st in reached_steps: - if st in step_index and step_index[st] > max_idx: - max_idx = step_index[st] - for i in range(max_idx + 1): - reached[i] += 1 + for step in _normalize_path_to_steps(path or ""): + if not _funnel_step_counts(path or "", str(result or ""), step): + continue + if step in step_actors: + step_actors[step].add(key) funnel_list = [] - prev_count = None - for i, step in enumerate(steps): - count = reached[i] - conversion = round(100.0 * count / prev_count, 1) if prev_count and prev_count > 0 else 100.0 - funnel_list.append({ - "step": step, - "label": AUDIT_STEP_LABELS.get(step, step), - "count": count, - "conversion_from_prev_pct": conversion if prev_count else None, - }) - prev_count = count + prev_actors: set[tuple[str, str | int]] | None = None + for step in steps: + actors = step_actors.get(step, set()) + count = len(actors) + conversion = None + if prev_actors is not None and prev_actors: + overlap = len(prev_actors & actors) + conversion = round(100.0 * overlap / len(prev_actors), 1) + funnel_list.append( + { + "step": step, + "label": AUDIT_STEP_LABELS.get(step, step), + "count": count, + "conversion_from_prev_pct": conversion, + } + ) + prev_actors = actors return funnel_list @@ -1182,7 +1025,6 @@ async def get_audit_funnel_from_redis( max_events: int = 5000, steps_ordered: tuple[str, ...] | None = None, ) -> list[dict[str, Any]] | None: - """Воронка по событиям из буфера Redis. None, если буфер выключен.""" if not AUDIT_REDIS_BUFFER_ENABLED: return None events = await list_audit_events_from_redis_buffer(max_events=max_events) @@ -1190,46 +1032,91 @@ async def get_audit_funnel_from_redis( return _funnel_from_rows(rows, steps_ordered) -async def drain_audit_redis_to_db(session_factory: Any) -> int: - """Выгружает буфер аудита из Redis в БД батчами. Вызывать по крону (например в 00:00). - Возвращает количество записанных событий.""" - from core.redis_cache import cache_lpop_batch +async def get_audit_funnel_from_redis_since( + *, + date_from: datetime | None = None, + max_events: int = 5000, + steps_ordered: tuple[str, ...] | None = None, +) -> list[dict[str, Any]] | None: + if not AUDIT_REDIS_BUFFER_ENABLED: + return None + events = await list_audit_events_from_redis_buffer(max_events=max_events) + filtered_events = events + if date_from is not None: + threshold = date_from.astimezone(timezone.utc) if date_from.tzinfo else date_from.replace(tzinfo=timezone.utc) + filtered_events = [e for e in events if getattr(e, "created_at", None) and e.created_at >= threshold] + rows = [(e.path_or_handler, e.result, e.actor_tg_id, e.actor_identity_id) for e in filtered_events] + return _funnel_from_rows(rows, steps_ordered) + +async def drain_audit_redis_to_db(session_factory: Any) -> int: + from core.redis_cache import cache_delete, cache_lmove_batch, cache_lpop_batch, cache_lrange, cache_setnx + + if not await cache_setnx(_AUDIT_REDIS_DRAIN_LOCK_KEY, 1, _AUDIT_REDIS_DRAIN_LOCK_TTL_SEC): + logger.info("[Audit] drain_audit_redis_to_db пропущен: уже выполняется другой drain") + return 0 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, + try: + while True: + raw_batch = await cache_lrange(_AUDIT_REDIS_PROCESSING_KEY, 0, AUDIT_REDIS_DRAIN_BATCH - 1) + if not raw_batch: + raw_batch = await cache_lmove_batch( + AUDIT_REDIS_FLUSH_KEY, + _AUDIT_REDIS_PROCESSING_KEY, + AUDIT_REDIS_DRAIN_BATCH, + ) + if not raw_batch: + break + batch = [rec for rec in raw_batch if isinstance(rec, dict)] + if not batch: + logger.warning("[Audit] drain_audit_redis_to_db: отброшен пустой/битый батч (%s элементов)", len(raw_batch)) + await cache_lpop_batch(_AUDIT_REDIS_PROCESSING_KEY, len(raw_batch)) + continue + try: + async with session_factory() as session: + await ensure_audit_table(session) + existing_request_ids = await fetch_existing_audit_request_ids_db( + session, + [str(rec.get("request_id")) for rec in batch if rec.get("request_id")], ) - 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 + inserted_count = 0 + seen_batch_request_ids: set[str] = set() + for rec in batch: + request_id = rec.get("request_id") + if request_id and (request_id in existing_request_ids or request_id in seen_batch_request_ids): + continue + if request_id: + seen_batch_request_ids.add(request_id) + 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=request_id, + created_at=created, + ) + session.add(event) + inserted_count += 1 + await session.commit() + await cache_lpop_batch(_AUDIT_REDIS_PROCESSING_KEY, len(raw_batch)) + total += inserted_count + except Exception as exc: + logger.warning("[Audit] drain_audit_redis_to_db батч не записан, останется в Redis processing: %s", exc) + break + return total + finally: + await cache_delete(_AUDIT_REDIS_DRAIN_LOCK_KEY) diff --git a/audit/rules.py b/audit/rules.py new file mode 100644 index 00000000..ae7209af --- /dev/null +++ b/audit/rules.py @@ -0,0 +1,433 @@ +from __future__ import annotations + +from typing import Iterable + + +def _get_bot_webhook_path() -> str: + """Путь вебхука бота из конфига (для исключения из шага «успешная оплата»).""" + try: + from config import WEBHOOK_PATH + + return ((WEBHOOK_PATH or "").strip().lower()) or "" + except ImportError: + return "" + + +def _is_bot_webhook_path(path: str) -> bool: + """True только если path — именно вебхук бота (точное совпадение сегмента пути), не касса.""" + bot_path = _get_bot_webhook_path() + if not bot_path: + return False + p = (path or "").strip().lower() + path_segment = p.split(" ", 1)[1] if " " in p else p + return path_segment == bot_path or path_segment.rstrip("/") == bot_path.rstrip("/") + + +AUDIT_STEP_LABELS: dict[str, str] = { + "start": "Старт", + "start_coupon": "Старт: купон", + "start_gift": "Старт: подарок", + "start_referral": "Старт: рефералка", + "start_utm": "Старт: UTM", + "profile": "Профиль", + "about": "О VPN", + "instructions": "Инструкции", + "balance": "Баланс / история оплат", + "view_keys": "Мои ключи", + "buy_entry": "Оформление: вход", + "tariff_config": "Оформление: выбор тарифа/конфига", + "key_create": "Подписка оформлена (ключ создан)", + "pay_start": "Оплата: вход / создание ссылки", + "pay": "Успешная оплата (пополнение)", + "key_view": "Ключ (карточка)", + "connect": "Подключение: экран / инструкции / QR", + "key_manage": "Управление подпиской", + "renew": "Продление", + "addons": "Аддоны", + "referral": "Рефералы", + "coupons": "Купоны", + "register": "Регистрация (API)", + "login": "Вход (API)", + "api_other": "API прочее", + "admin": "Админ-панель", + "other": "Прочее", +} + + +DEFAULT_FUNNEL_STEPS = ( + "start", + "profile", + "view_keys", + "buy_entry", + "tariff_config", + "key_create", + "pay_start", + "pay", + "key_view", + "connect", +) + + +_CALLBACK_EXACT: dict[str, str] = { + "start": "start", + "profile": "profile", + "about_vpn": "about", + "instructions": "instructions", + "view_keys": "view_keys", + "create_key": "buy_entry", + "buy": "buy_entry", + "pay": "pay_start", + "balance": "balance", + "balance_history": "balance", + "activate_coupon": "coupons", + "cancel_coupon_activation": "coupons", + "exit_coupon_input": "coupons", + "invite": "referral", + "top_referrals": "referral", + "check_subscription": "start", + "back_to_tariff_group_list": "tariff_config", + "back_to_subgroup_tariffs": "tariff_config", + "cancel_and_back_to_view_keys": "view_keys", + "fastflow_coupon": "coupons", + "fastflow_coupon_back": "coupons", + "fastflow_back": "pay_start", + "pay_kassai": "pay_start", + "pay_kassai_cards": "pay_start", + "pay_kassai_sbp": "pay_start", + "pay_heleket_crypto": "pay_start", + "pay_freekassa": "pay_start", + "pay_robokassa": "pay_start", +} + + +_CALLBACK_PREFIX: list[tuple[str, str]] = [ + ("view_key|", "key_view"), + ("view_keys|", "view_keys"), + ("show_referral_qr|", "referral"), + ("tariff_subgroup_user|", "tariff_config"), + ("select_tariff_plan|", "tariff_config"), + ("cfg_user_devices|", "tariff_config"), + ("cfg_user_traffic|", "tariff_config"), + ("rename_key|", "key_manage"), + ("reset_hwid|", "key_manage"), + ("freeze_subscription|", "key_manage"), + ("freeze_subscription_confirm|", "key_manage"), + ("unfreeze_subscription|", "key_manage"), + ("unfreeze_subscription_confirm|", "key_manage"), + ("delete_key|", "key_manage"), + ("confirm_delete|", "key_manage"), + ("update_subscription|", "key_manage"), + ("change_location|", "key_manage"), + ("select_country|", "key_manage"), + ("connect_device|", "connect"), + ("connect_router|", "connect"), + ("connect_tv|", "connect"), + ("connect_pc|", "connect"), + ("connect_ios|", "connect"), + ("connect_android|", "connect"), + ("show_qr|", "connect"), + ("continue_tv|", "connect"), + ("pay_currency|", "pay_start"), + ("choose_payment_currency|", "pay_start"), + ("cfg_user_confirm|", "key_create"), + ("choose_payment_provider|", "pay_start"), + ("kassai_method|", "pay_start"), + ("kassai_cards_amount|", "pay_start"), + ("kassai_sbp_amount|", "pay_start"), + ("kassai_custom_amount|", "pay_start"), + ("heleket_method|", "pay_start"), + ("heleket_crypto_amount|", "pay_start"), + ("heleket_custom_amount|", "pay_start"), + ("freekassa_amount|", "pay_start"), + ("robokassa_", "pay_start"), + ("cfg_renew", "renew"), + ("key_addons", "addons"), + ("extend_key", "coupons"), +] + + +_CALLBACK_CONTAINS: list[tuple[str, str]] = [ + ("connect_", "connect"), + ("renew", "renew"), + ("addon", "addons"), + ("referral", "referral"), + ("invite", "referral"), + ("coupon", "coupons"), + ("users_audit", "admin"), + ("users_editor", "admin"), + ("search_user", "admin"), + ("admin_panel", "admin"), +] + + +_MESSAGE_STEP_BY_COMMAND: dict[str, str] = { + "/start": "start", + "start": "start", + "/buy": "buy_entry", + "buy": "buy_entry", + "/subs": "view_keys", + "subs": "view_keys", + "/profile": "profile", + "profile": "profile", + "/invite": "referral", + "invite": "referral", + "/instructions": "instructions", + "instructions": "instructions", + "/activate_coupon": "coupons", + "activate_coupon": "coupons", +} + +_START_PAYLOAD_MAX_LEN = 256 +_START_PAYLOAD_MAX_PARTS = 20 + + +_IGNORED_CALLBACK_EXACT: set[str] = { + " ", + "back", + "cancel", + "back_to_pay", + "back_to_currency", + "back_to_tariff_group_list", + "back_to_subgroup_tariffs", + "cancel_and_back_to_view_keys", + "cancel_coupon_activation", + "exit_coupon_input", + "fastflow_back", + "fastflow_coupon_back", + "cfg_back_menu", + "cfg_cancel_input", + "cancel_broadcast", +} + + +_IGNORED_CALLBACK_PREFIX: tuple[str, ...] = ( + "back:", + "back_to_", + "cancel_", + "gifts_page|", +) + +_IGNORED_CALLBACK_EXACT_RULES: dict[str, str] = {key: "ignore" for key in _IGNORED_CALLBACK_EXACT} +_IGNORED_CALLBACK_PREFIX_RULES: tuple[tuple[str, str], ...] = tuple( + (prefix, "ignore") for prefix in _IGNORED_CALLBACK_PREFIX +) + + +_HANDLER_CONTAINS: list[tuple[str, str] | tuple[str, str, str]] = [ + ("process_start", "start"), + ("start_entry", "start"), + ("show_start_menu", "start"), + ("process_callback_view_profile", "profile"), + ("handle_about_vpn", "about"), + ("send_instructions", "instructions"), + ("process_callback_or_message_view_keys", "view_keys"), + ("key_view", "key_view", "key_create"), + ("handle_user_config_confirm", "key_create"), + ("finalize_config_and_purchase", "key_create"), + ("proceed_purchase_with_values", "tariff_config"), + ("key_create", "buy_entry"), + ("handle_key_creation", "buy_entry"), + ("complete_key_renewal", "renew"), + ("handle_connect_device", "connect"), + ("process_connect_", "connect"), + ("process_callback_connect", "connect"), + ("process_continue_tv", "connect"), + ("show_qr_code", "connect"), + ("balance_history", "balance"), + ("balance_handler", "balance"), + ("pay", "pay_start"), + ("choose_payment_provider", "pay_start"), + ("kassai_", "pay_start"), + ("heleket_", "pay_start"), + ("freekassa_", "pay_start"), + ("robokassa_", "pay_start"), + ("rename_key", "key_manage"), + ("reset_hwid", "key_manage"), + ("freeze_subscription", "key_manage"), + ("unfreeze_subscription", "key_manage"), + ("delete_key", "key_manage"), + ("change_location", "key_manage"), + ("select_country", "key_manage"), + ("renew", "renew"), + ("addon", "addons"), + ("referral", "referral"), + ("refferal", "referral"), + ("coupon", "coupons"), + ("admin_panel", "admin"), + ("users_audit", "admin"), + ("users_editor", "admin"), + ("search_user", "admin"), + ("auth/register", "register"), + ("auth/login", "login"), + ("auth/send-login", "login"), + ("auth/login-by-code", "login"), + ("auth/login-telegram", "login"), +] + + +_API_CONTAINS: list[tuple[str, str]] = [ + ("auth/register", "register"), + ("auth/login", "login"), + ("auth/send-login", "login"), + ("auth/login-by-code", "login"), + ("auth/login-telegram", "login"), + ("/keys/create", "key_create"), +] + + +def _match_step_rules( + value: str, + *, + exact: dict[str, str] | None = None, + prefixes: Iterable[tuple[str, str]] | None = None, + contains: Iterable[tuple[str, str] | tuple[str, str, str]] | None = None, +) -> str | None: + """Универсальный matcher шага по exact/prefix/contains правилам.""" + normalized = (value or "").lower().strip() + if not normalized: + return None + if exact: + step = exact.get(normalized) + if step: + return step + if prefixes: + for prefix, step in prefixes: + if normalized.startswith(prefix): + return step + if contains: + for item in contains: + substr, step = item[0], item[1] + exclude = item[2] if len(item) > 2 else None + if substr in normalized and (exclude is None or exclude not in normalized): + return step + return None + + +def _is_ignored_analytics_event(path: str) -> bool: + """Исключает чисто навигационные события из аналитики шагов и воронки.""" + p = (path or "").lower().strip() + if not p.startswith("callback:"): + return False + callback_data = p.split(":", 1)[-1] + return _match_step_rules( + callback_data, + exact=_IGNORED_CALLBACK_EXACT_RULES, + prefixes=_IGNORED_CALLBACK_PREFIX_RULES, + ) is not None + + +def _message_command_step(path: str) -> str | None: + """Определяет шаг только по точной команде/первому токену сообщения.""" + steps = _message_command_steps(path) + return steps[0] if steps else None + + +def _message_command_steps(path: str) -> list[str]: + """Определяет один или несколько шагов из текстового сообщения.""" + p = (path or "").strip() + if not p.lower().startswith("message:"): + return [] + text = (p.split(":", 1)[-1] or "").strip().lower() + if not text: + return [] + parts = text.split(None, 1) + token = parts[0] + if token.startswith("/"): + token = token.split("@", 1)[0] + if token in ("/start", "start"): + steps = ["start"] + payload = parts[1].strip() if len(parts) > 1 else "" + if payload: + if len(payload) > _START_PAYLOAD_MAX_LEN: + payload = payload[:_START_PAYLOAD_MAX_LEN] + payload_parts = payload.split("-") + if len(payload_parts) > _START_PAYLOAD_MAX_PARTS: + payload_parts = payload_parts[:_START_PAYLOAD_MAX_PARTS] + for part in payload_parts: + part = part.strip() + if not part: + continue + if "coupons" in part: + steps.append("start_coupon") + continue + if "gift" in part: + steps.append("start_gift") + continue + if "referral" in part: + steps.append("start_referral") + continue + if "utm" in part: + steps.append("start_utm") + continue + normalized_payload = "-".join(payload_parts).strip().lower() + if normalized_payload in _MESSAGE_STEP_BY_COMMAND: + steps.append(_MESSAGE_STEP_BY_COMMAND[normalized_payload]) + return list(dict.fromkeys(steps)) + return steps + step = _MESSAGE_STEP_BY_COMMAND.get(token) + return [step] if step else [] + + +def _callback_step(path: str) -> str | None: + callback_data = path.split(":", 1)[-1] + callback_key = callback_data.split("|")[0] + return _match_step_rules( + callback_key, + exact=_CALLBACK_EXACT, + ) or _match_step_rules( + callback_data, + prefixes=_CALLBACK_PREFIX, + contains=_CALLBACK_CONTAINS, + ) + + +def _api_step(path: str) -> str: + step = _match_step_rules(path, contains=_API_CONTAINS) + if step: + return step + if "payment-links" in path or ("payment" in path and "webhook" not in path): + return "pay_start" + return "api_other" + + +def _handler_step(path: str) -> str | None: + return _match_step_rules(path, contains=_HANDLER_CONTAINS) + + +def _funnel_step_counts(path: str, result: str, step: str) -> bool: + """Решает, считать ли событие достижением шага воронки.""" + if result != "success": + return False + p = (path or "").lower() + if step == "pay_start": + return True + if step == "pay": + return p.startswith("payment_success:") + if step == "key_create": + return "cfg_user_confirm" in p or "/keys/create" in p + return True + + +def _normalize_path_to_step(path: str) -> str: + """Сводит path_or_handler к шагу по правилам из маппингов.""" + steps = _normalize_path_to_steps(path) + return steps[0] if steps else "other" + + +def _normalize_path_to_steps(path: str) -> list[str]: + """Сводит path_or_handler к одному или нескольким шагам по правилам аудита.""" + if not path: + return ["other"] + p = path.lower().strip() + if p.startswith("payment_success:"): + return ["pay"] + if p.startswith("callback:"): + return [_callback_step(p) or "other"] + message_steps = _message_command_steps(path) + if message_steps: + return message_steps + if p.startswith("message:"): + return ["other"] + if p.startswith("post ") or p.startswith("get "): + return [_api_step(p)] + return [_handler_step(p) or "other"] diff --git a/cli_launcher.py b/cli_launcher.py index ee677e36..94d87ba6 100755 --- a/cli_launcher.py +++ b/cli_launcher.py @@ -303,7 +303,7 @@ def install_rsync_if_needed(): os.system("sudo apt update && sudo apt install -y rsync") -def clean_project_dir_safe(update_buttons=False, update_img=False): +def clean_project_dir_safe(update_buttons=False, update_img=False, update_redis_cache=False): console.print("[yellow]Очистка проекта перед обновлением...[/yellow]") preserved_paths = set() @@ -328,6 +328,9 @@ def clean_project_dir_safe(update_buttons=False, update_img=False): for name in dirs + files: preserved_paths.add(os.path.join(root, name)) + if not update_redis_cache: + preserved_paths.add(os.path.join(PROJECT_DIR, "core", "redis_cache.py")) + for root, dirs, files in os.walk(PROJECT_DIR, topdown=False): for file in files: path = os.path.join(root, file) @@ -474,6 +477,7 @@ def update_from_beta(): update_buttons = safe_confirm("[yellow]Обновлять файл buttons.py?[/yellow]", default=False) update_img = safe_confirm("[yellow]Обновлять папку img?[/yellow]", default=False) + update_redis_cache = safe_confirm("[yellow]Обновлять файл core/redis_cache.py?[/yellow]", default=False) backup_project() install_git_if_needed() @@ -488,13 +492,19 @@ def update_from_beta(): return subprocess.run(["sudo", "rm", "-rf", os.path.join(PROJECT_DIR, "venv")]) - clean_project_dir_safe(update_buttons=update_buttons, update_img=update_img) + clean_project_dir_safe( + update_buttons=update_buttons, + update_img=update_img, + update_redis_cache=update_redis_cache, + ) exclude_options = "" if not update_img: exclude_options += "--exclude=img " if not update_buttons: exclude_options += "--exclude=handlers/buttons.py " + if not update_redis_cache: + exclude_options += "--exclude=core/redis_cache.py " exclude_options += "--exclude=modules " rsync_cmd = ["rsync", "-a"] + [x for x in exclude_options.split() if x] + [f"{TEMP_DIR}/", f"{PROJECT_DIR}/"] @@ -520,7 +530,7 @@ def update_from_beta(): console.print("[green]Обновление с ветки dev завершено.[/green]") -def _do_update_to_tag(tag_name: str, update_buttons: bool, update_img: bool) -> None: +def _do_update_to_tag(tag_name: str, update_buttons: bool, update_img: bool, update_redis_cache: bool) -> None: """Общая логика обновления до указанного тега (релиз или произвольный тег).""" subprocess.run(["rm", "-rf", TEMP_DIR]) subprocess.run( @@ -530,13 +540,19 @@ def _do_update_to_tag(tag_name: str, update_buttons: bool, update_img: bool) -> console.print("[red]Начинается перезапись файлов бота![/red]") subprocess.run(["sudo", "rm", "-rf", os.path.join(PROJECT_DIR, "venv")]) - clean_project_dir_safe(update_buttons=update_buttons, update_img=update_img) + clean_project_dir_safe( + update_buttons=update_buttons, + update_img=update_img, + update_redis_cache=update_redis_cache, + ) exclude_options = "" if not update_img: exclude_options += "--exclude=img " if not update_buttons: exclude_options += "--exclude=handlers/buttons.py " + if not update_redis_cache: + exclude_options += "--exclude=core/redis_cache.py " exclude_options += "--exclude=modules " rsync_cmd = ["rsync", "-a"] + exclude_options.split() + [f"{TEMP_DIR}/", f"{PROJECT_DIR}/"] @@ -567,12 +583,13 @@ def update_from_release(): return console.print("[red]ВНИМАНИЕ! Папка бота будет полностью перезаписана![/red]") - console.print("[red] Исключения: папка img и файл handlers/buttons.py[/red]") + console.print("[red] Исключения: папка img, файл handlers/buttons.py и файл core/redis_cache.py[/red]") if not safe_confirm("[red]Вы точно хотите продолжить?[/red]"): return update_buttons = safe_confirm("[yellow]Обновлять файл buttons.py?[/yellow]", default=False) update_img = safe_confirm("[yellow]Обновлять папку img?[/yellow]", default=False) + update_redis_cache = safe_confirm("[yellow]Обновлять файл core/redis_cache.py?[/yellow]", default=False) backup_project() install_git_if_needed() @@ -618,7 +635,7 @@ def update_from_release(): return console.print(f"[cyan]Клонируем {tag_name} во временную папку...[/cyan]") - _do_update_to_tag(tag_name, update_buttons, update_img) + _do_update_to_tag(tag_name, update_buttons, update_img, update_redis_cache) except Exception as e: console.print(f"[red]❌ Ошибка при обновлении: {e}[/red]") diff --git a/core/app.cpython-312-x86_64-linux-gnu.so b/core/app.cpython-312-x86_64-linux-gnu.so index 16c390e2..46277b43 100644 Binary files a/core/app.cpython-312-x86_64-linux-gnu.so and b/core/app.cpython-312-x86_64-linux-gnu.so differ diff --git a/core/cache_config.py b/core/cache_config.py index 53e254b6..ba76e512 100644 --- a/core/cache_config.py +++ b/core/cache_config.py @@ -29,6 +29,8 @@ key_email │ 45 │ email по client_id payment_pending │ 3600 │ Ожидающий платёж (1 ч) audit_history │ 300 │ История действий клиента (админка) audit:flush │ — │ Буфер аудита (список для выгрузки в БД в 00:00) +audit:flush:processing │ — │ Батч аудита, перенесённый в processing до commit в БД +audit:flush:drain_lock │ 900 │ Лок nightly/manual drain аудита audit:user:tg:* │ 25 ч │ События по tg_id для чтения до выгрузки audit:user:identity:* │ 25 ч │ События по identity для чтения до выгрузки webhook_abuse_fail │ 60 │ Счётчик неудачных вебхуков по IP diff --git a/core/redis_cache.py b/core/redis_cache.py index 7401a8f6..a53bcae6 100644 --- a/core/redis_cache.py +++ b/core/redis_cache.py @@ -216,3 +216,42 @@ async def cache_lpop_batch(key: str, count: int) -> list[Any]: return out except Exception: return out + + +async def cache_lmove_batch(source: str, destination: str, count: int) -> list[Any]: + """Атомарно переносит до count элементов из головы source в хвост destination.""" + if count <= 0: + return [] + client = await _get_redis() + if client is None: + return [] + try: + raw_list = await client.eval( + """ + local moved = {} + local count = tonumber(ARGV[1]) + for i = 1, count do + local item = redis.call('LPOP', KEYS[1]) + if not item then + break + end + redis.call('RPUSH', KEYS[2], item) + table.insert(moved, item) + end + return moved + """, + 2, + source, + destination, + int(count), + ) + out = [] + for raw in raw_list or []: + try: + out.append(json.loads(raw)) + except Exception: + pass + return out + except Exception as exc: + logger.warning(f"[Redis] lmove_batch({source}->{destination}) не удался: {exc}") + return [] diff --git a/database/__init__.py b/database/__init__.py index 0ee6db64..664ca325 100644 --- a/database/__init__.py +++ b/database/__init__.py @@ -1,3 +1,4 @@ +from .audit import * from .bans import * from .coupons import * from .db import async_session_maker diff --git a/database/audit.py b/database/audit.py new file mode 100644 index 00000000..ef30be54 --- /dev/null +++ b/database/audit.py @@ -0,0 +1,200 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from typing import Iterable + +from sqlalchemy import DateTime as SQLADateTime, and_, cast, delete, desc, func, or_, select +from sqlalchemy.ext.asyncio import AsyncSession + +from database.models import AuditEvent, Payment + +try: + from core.constants import PAYMENT_SYSTEMS_EXCLUDED +except ImportError: + PAYMENT_SYSTEMS_EXCLUDED = ("referral", "coupon", "cashback", "admin") + + +_AUDIT_TABLE_READY = False + + +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_db( + session: AsyncSession, + *, + older_than_days: int = 90, +) -> int: + 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 fetch_audit_rows_db( + session: AsyncSession, + *, + date_from: datetime, + date_to: datetime, + limit: int, +) -> list[tuple[str, str, int | None, str | None]]: + await ensure_audit_table(session) + stmt = ( + select( + AuditEvent.path_or_handler, + AuditEvent.result, + AuditEvent.actor_tg_id, + AuditEvent.actor_identity_id, + ) + .where( + AuditEvent.event_type != "audit_reset", + AuditEvent.created_at >= date_from, + AuditEvent.created_at < date_to, + ) + .order_by(AuditEvent.created_at.desc()) + .limit(limit) + ) + result = await session.execute(stmt) + return list(result.all()) + + +async def fetch_successful_payment_rows_db( + session: AsyncSession, + *, + date_from: datetime, + date_to: datetime, + limit: int, +) -> list[tuple[str, str, int | None, None]]: + success_at_expr = func.coalesce( + cast(Payment.metadata_["status_changed_at"].astext, SQLADateTime), + Payment.created_at, + ) + stmt = ( + select(Payment.payment_system, Payment.payment_id, Payment.tg_id) + .where( + Payment.status == "success", + Payment.payment_system.notin_(PAYMENT_SYSTEMS_EXCLUDED), + success_at_expr >= date_from, + success_at_expr < date_to, + ) + .order_by(desc(success_at_expr)) + .limit(limit) + ) + result = await session.execute(stmt) + return [ + (f"payment_success:{payment_system or '-'}:{payment_id or '-'}", "success", tg_id, None) + for payment_system, payment_id, tg_id in result.all() + ] + + +async def fetch_latest_audit_reset_db( + session: AsyncSession, + *, + source: str = "db", +) -> datetime | None: + await ensure_audit_table(session) + stmt = select(func.max(AuditEvent.created_at)).where( + AuditEvent.event_type == "audit_reset", + AuditEvent.channel == "system", + AuditEvent.path_or_handler == f"audit_reset:{source}", + ) + result = await session.execute(stmt) + return result.scalar_one_or_none() + + +async def create_audit_reset_marker_db( + session: AsyncSession, + *, + source: str = "db", + created_at: datetime | None = None, +) -> datetime: + await ensure_audit_table(session) + created = created_at or datetime.utcnow() + event = AuditEvent( + event_type="audit_reset", + channel="system", + path_or_handler=f"audit_reset:{source}", + result="success", + created_at=created, + ) + session.add(event) + await session.flush() + return created + + +async def fetch_existing_audit_request_ids_db( + session: AsyncSession, + request_ids: Iterable[str], +) -> set[str]: + await ensure_audit_table(session) + request_ids_list = sorted({rid for rid in request_ids if rid}) + if not request_ids_list: + return set() + stmt = select(AuditEvent.request_id).where(AuditEvent.request_id.in_(request_ids_list)) + result = await session.execute(stmt) + return {rid for rid in result.scalars().all() if rid} + + +async def fetch_audit_events_db( + 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]: + await ensure_audit_table(session) + 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) + event_types_list = sorted(event_types) if event_types else None + 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()) + + +async def fetch_audit_events_db_window( + 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 = 5000, +) -> list[AuditEvent]: + return await fetch_audit_events_db( + session, + identity_id=identity_id, + tg_id=tg_id, + channel=channel, + event_type=event_type, + event_types=event_types, + limit=limit, + offset=0, + ) diff --git a/database/payments.py b/database/payments.py index 19dbdbb6..8590039f 100644 --- a/database/payments.py +++ b/database/payments.py @@ -165,9 +165,12 @@ async def update_payment_status( payment.status = new_status if payment_id is not None: payment.payment_id = payment_id + base = payment.metadata_ or {} + if new_status == "success" and "status_changed_at" not in base: + base["status_changed_at"] = datetime.utcnow().replace(tzinfo=None).isoformat() if metadata_patch: - base = payment.metadata_ or {} base.update(metadata_patch) + if base: payment.metadata_ = base await session.commit() diff --git a/handlers/admin/stats/keyboard.py b/handlers/admin/stats/keyboard.py index 03e550ea..90e402fb 100644 --- a/handlers/admin/stats/keyboard.py +++ b/handlers/admin/stats/keyboard.py @@ -4,10 +4,33 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder from ..panel.keyboard import AdminPanelCallback, build_admin_back_btn -def build_audit_refresh_kb() -> InlineKeyboardMarkup: - """Клавиатура под сообщением аудита: кнопка «Обновить».""" +def build_audit_refresh_kb(source: str = "db") -> InlineKeyboardMarkup: + """Клавиатура под сообщением аудита: выбор источника данных.""" builder = InlineKeyboardBuilder() - builder.button(text="🔄 Обновить", callback_data=AdminPanelCallback(action="audit_refresh").pack()) + redis_text = "• Redis raw" if source == "redis" else "Redis raw" + db_text = "• БД вчера" if source == "db" else "БД вчера" + reset_text = "Сбросить Redis" if source == "redis" else "Сбросить БД" + builder.button(text=redis_text, callback_data=AdminPanelCallback(action="audit_refresh_redis").pack()) + builder.button(text=db_text, callback_data=AdminPanelCallback(action="audit_refresh_db").pack()) + builder.button(text=reset_text, callback_data=AdminPanelCallback(action=f"audit_reset_ask_{source}").pack()) + builder.adjust(2, 1) + return builder.as_markup() + + +def build_audit_source_kb() -> InlineKeyboardMarkup: + """Клавиатура выбора источника аудита при первом открытии.""" + builder = InlineKeyboardBuilder() + builder.button(text="Redis raw", callback_data=AdminPanelCallback(action="audit_refresh_redis").pack()) + builder.button(text="БД вчера", callback_data=AdminPanelCallback(action="audit_refresh_db").pack()) + builder.adjust(2) + return builder.as_markup() + + +def build_audit_reset_confirm_kb(source: str) -> InlineKeyboardMarkup: + builder = InlineKeyboardBuilder() + builder.button(text="Да, сбросить", callback_data=AdminPanelCallback(action=f"audit_reset_do_{source}").pack()) + builder.button(text="Отмена", callback_data=AdminPanelCallback(action=f"audit_refresh_{source}").pack()) + builder.adjust(1) return builder.as_markup() diff --git a/handlers/admin/stats/stats_handler.py b/handlers/admin/stats/stats_handler.py index 88a2dfdc..1695a5e7 100644 --- a/handlers/admin/stats/stats_handler.py +++ b/handlers/admin/stats/stats_handler.py @@ -1,5 +1,5 @@ from collections import Counter -from datetime import datetime, timedelta +from datetime import date, datetime, timedelta from html import escape import pytz @@ -9,7 +9,16 @@ from aiogram.exceptions import TelegramBadRequest from aiogram.types import CallbackQuery, Message from sqlalchemy.ext.asyncio import AsyncSession -from audit import get_audit_funnel, get_audit_funnel_from_redis, get_audit_stats, get_audit_stats_from_redis +from audit import ( + AUDIT_STEP_LABELS, + clear_audit_redis_buffers, + get_audit_funnel, + get_audit_funnel_from_redis, + get_audit_db_reset_at, + get_audit_stats, + get_audit_stats_from_redis, + set_audit_db_reset_at, +) from bot import bot from config import ADMIN_ID from database import ( @@ -40,11 +49,63 @@ from utils.csv_export import ( ) from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb -from .keyboard import build_audit_refresh_kb, build_stats_kb +from .keyboard import build_audit_refresh_kb, build_audit_reset_confirm_kb, build_audit_source_kb, build_stats_kb router = Router() +KEY_AUDIT_STEPS_ORDER = ( + "start", + "start_coupon", + "start_gift", + "start_referral", + "start_utm", + "profile", + "about", + "instructions", + "balance", + "view_keys", + "buy_entry", + "tariff_config", + "key_create", + "pay_start", + "pay", + "key_view", + "connect", + "key_manage", + "renew", + "addons", + "referral", + "coupons", +) + + +def _previous_moscow_day_window(now: datetime | None = None) -> tuple[datetime, datetime, date]: + moscow_tz = pytz.timezone("Europe/Moscow") + current = now or datetime.now(moscow_tz) + yesterday_date = current.date() - timedelta(days=1) + start = moscow_tz.localize(datetime.combine(yesterday_date, datetime.min.time())) + end = start + timedelta(days=1) + return start.astimezone(pytz.UTC), end.astimezone(pytz.UTC), yesterday_date + + +def _audit_success_event_counts(by_path: list[dict]) -> dict[str, int]: + """Успешные события по шагам для бизнес-метрик в отчёте.""" + return {row["step"]: int(row.get("success", 0) or 0) for row in by_path} + + +def _append_key_audit_steps(lines: list[str], by_path: list[dict]) -> None: + """Добавляет фиксированный список ключевых шагов клиента, включая нулевые значения.""" + totals_by_step = {row["step"]: row for row in by_path} + lines.append("Ключевые шаги клиента:") + for step in KEY_AUDIT_STEPS_ORDER: + row = totals_by_step.get(step) + total = row["total"] if row else 0 + success = row["success"] if row else 0 + fail = row["fail"] if row else 0 + label = AUDIT_STEP_LABELS.get(step, step) + lines.append(f" • {label}: {total} (ок: {success}, ошибок: {fail})") + @router.callback_query(AdminPanelCallback.filter(F.action == "stats"), IsAdminFilter()) async def handle_stats(callback_query: CallbackQuery, session: AsyncSession): @@ -277,7 +338,11 @@ async def handle_stats_audit(callback_query: CallbackQuery, session: AsyncSessio lines = [ f"📊 Аудит за {yesterday_date.strftime('%d.%m.%Y')} (МСК)", "", - f"📎 Событий: {summary['total_events']} │ Уникальных пользователей: {summary['unique_users']}", + ( + f"📎 Сырых событий: {summary.get('raw_total_events', summary['total_events'])} │ " + f"Аналитических шагов: {summary.get('analytics_total_events', summary['total_events'])} │ " + f"Уникальных пользователей: {summary['unique_users']}" + ), "", "По шагам (топ по объёму):", ] @@ -287,18 +352,19 @@ async def handle_stats_audit(callback_query: CallbackQuery, session: AsyncSessio f"{fail_mark} {row['label']}: {row['total']} (ок: {row['success']}, ошибок: {row['fail']}, {row['fail_rate_pct']}% ошибок)" ) lines.append("") - by_step_totals = {row["step"]: row["total"] for row in by_path} - funnel_counts = {s["step"]: s["count"] for s in funnel} - pay_start = by_step_totals.get("pay_start", 0) - pay_ok = max(by_step_totals.get("pay", 0), funnel_counts.get("pay", 0)) - key_created = by_step_totals.get("key_create", 0) - connect_ok = max(by_step_totals.get("connect", 0), funnel_counts.get("connect", 0)) - pct_pay = round(100.0 * pay_ok / pay_start, 1) if pay_start else 0 - pct_connect = round(100.0 * connect_ok / key_created, 1) if key_created else 0 - lines.append("Оплата: начало {0}, успешная {1}, % успешных от созданных: {2}%".format(pay_start, pay_ok, pct_pay)) - lines.append("Подписка: оформлена {0}, подключена {1}, % успешных подключений: {2}%".format(key_created, connect_ok, pct_connect)) + _append_key_audit_steps(lines, by_path) lines.append("") - lines.append("Воронка (уник. пользователей по шагам):") + success_by_step = _audit_success_event_counts(by_path) + pay_start = success_by_step.get("pay_start", 0) + pay_ok = success_by_step.get("pay", 0) + key_created = success_by_step.get("key_create", 0) + connect_opened = success_by_step.get("connect", 0) + pct_pay = round(100.0 * pay_ok / pay_start, 1) if pay_start else 0 + pct_connect = round(100.0 * connect_opened / key_created, 1) if key_created else 0 + lines.append("Оплата: начало {0}, успешная {1}, % успешных от созданных: {2}%".format(pay_start, pay_ok, pct_pay)) + lines.append("Подписка: оформлена {0}, открыто подключение {1}, % от оформленных: {2}%".format(key_created, connect_opened, pct_connect)) + lines.append("") + lines.append("Воронка (уник. пользователей по точным шагам):") for step in funnel: conv = f" → {step['conversion_from_prev_pct']}%" if step["conversion_from_prev_pct"] is not None else "" lines.append(f" • {step['label']}: {step['count']} польз.{conv}") @@ -318,52 +384,70 @@ async def handle_stats_audit(callback_query: CallbackQuery, session: AsyncSessio ) -async def _build_audit_report(session: AsyncSession) -> tuple[str | None, str | None]: - """Собирает текст отчёта аудита. Возвращает (text, error): при успехе error=None; при отключённом Redis text=None, error=None.""" +async def _build_audit_report(session: AsyncSession, source: str = "db") -> tuple[str | None, str | None]: + """Собирает текст отчёта аудита из выбранного источника.""" try: - stats = await get_audit_stats_from_redis(max_events=5000) - funnel = await get_audit_funnel_from_redis(max_events=5000) - if stats is None or funnel is None: - return (None, None) - summary = stats["summary"] - by_path = stats["by_path"] - if summary["total_events"] == 0: - moscow_tz = pytz.timezone("Europe/Moscow") - now = datetime.now(moscow_tz) - end_utc = now.astimezone(pytz.UTC) - start_utc = end_utc - timedelta(hours=24) - stats = await get_audit_stats(session, date_from=start_utc, date_to=end_utc) - funnel = await get_audit_funnel(session, date_from=start_utc, date_to=end_utc) + moscow_tz = pytz.timezone("Europe/Moscow") + now = datetime.now(moscow_tz) + reset_note = "" + if source == "redis": + stats = await get_audit_stats_from_redis(max_events=5000) + funnel = await get_audit_funnel_from_redis(max_events=5000) + if stats is None or funnel is None: + return (None, "Буфер Redis для аудита выключен.") summary = stats["summary"] by_path = stats["by_path"] - header = "📊 Аудит из БД (буфер Redis пуст; последние 24 ч)" + header = "📊 Аудит (Redis raw, без добора из БД)" else: - header = "📊 Аудит из Redis (буфер, последние события)" + start_utc, end_utc, report_date = _previous_moscow_day_window(now) + reset_at = await get_audit_db_reset_at(session) + effective_start_utc = start_utc + if reset_at is not None: + effective_start_utc = max(start_utc, reset_at.astimezone(pytz.UTC)) + funnel = await get_audit_funnel(session, date_from=effective_start_utc, date_to=end_utc) + stats = await get_audit_stats(session, date_from=effective_start_utc, date_to=end_utc) + if effective_start_utc != start_utc: + reset_note = ( + "🧹 Сброс БД: " + f"{reset_at.astimezone(moscow_tz).strftime('%d.%m.%Y %H:%M')}" + ) + summary = stats["summary"] + by_path = stats["by_path"] + header = f"📊 Аудит (БД за {report_date.strftime('%d.%m.%Y')} МСК)" lines = [ header, "", - f"📎 Событий: {summary['total_events']} │ Уникальных пользователей: {summary['unique_users']}", + ( + f"📎 Сырых событий: {summary.get('raw_total_events', summary['total_events'])} │ " + f"Аналитических шагов: {summary.get('analytics_total_events', summary['total_events'])} │ " + f"Уникальных пользователей: {summary['unique_users']}" + ), "", - "По шагам (топ по объёму):", ] + if reset_note: + lines.extend([reset_note, ""]) + lines.extend([ + "По шагам (топ по объёму):", + ]) for row in by_path[:8]: fail_mark = "⚠️" if row["fail_rate_pct"] > 10 else "✅" lines.append( f"{fail_mark} {row['label']}: {row['total']} (ок: {row['success']}, ошибок: {row['fail']}, {row['fail_rate_pct']}% ошибок)" ) lines.append("") - by_step_totals = {row["step"]: row["total"] for row in by_path} - funnel_counts = {s["step"]: s["count"] for s in funnel} - pay_start = by_step_totals.get("pay_start", 0) - pay_ok = max(by_step_totals.get("pay", 0), funnel_counts.get("pay", 0)) - key_created = by_step_totals.get("key_create", 0) - connect_ok = max(by_step_totals.get("connect", 0), funnel_counts.get("connect", 0)) - pct_pay = round(100.0 * pay_ok / pay_start, 1) if pay_start else 0 - pct_connect = round(100.0 * connect_ok / key_created, 1) if key_created else 0 - lines.append("Оплата: начало {0}, успешная {1}, % успешных от созданных: {2}%".format(pay_start, pay_ok, pct_pay)) - lines.append("Подписка: оформлена {0}, подключена {1}, % успешных подключений: {2}%".format(key_created, connect_ok, pct_connect)) + _append_key_audit_steps(lines, by_path) lines.append("") - lines.append("Воронка (уник. пользователей по шагам):") + success_by_step = _audit_success_event_counts(by_path) + pay_start = success_by_step.get("pay_start", 0) + pay_ok = success_by_step.get("pay", 0) + key_created = success_by_step.get("key_create", 0) + connect_opened = success_by_step.get("connect", 0) + pct_pay = round(100.0 * pay_ok / pay_start, 1) if pay_start else 0 + pct_connect = round(100.0 * connect_opened / key_created, 1) if key_created else 0 + lines.append("Оплата: начало {0}, успешная {1}, % успешных от созданных: {2}%".format(pay_start, pay_ok, pct_pay)) + lines.append("Подписка: оформлена {0}, открыто подключение {1}, % от оформленных: {2}%".format(key_created, connect_opened, pct_connect)) + lines.append("") + lines.append("Воронка (уник. пользователей по точным шагам):") for step in funnel: conv = f" → {step['conversion_from_prev_pct']}%" if step["conversion_from_prev_pct"] is not None else "" lines.append(f" • {step['label']}: {step['count']} польз.{conv}") @@ -375,40 +459,104 @@ async def _build_audit_report(session: AsyncSession) -> tuple[str | None, str | @router.message(F.text.in_(["Аудит", "аудит"]), IsAdminFilter()) async def handle_audit_command(message: Message, session: AsyncSession): - """По команде «Аудит» — отправить статистику аудита из буфера Redis; при пустом буфере — из БД за 24 ч.""" - text, err = await _build_audit_report(session) - if err: - await message.answer(f"❗ Ошибка: {escape(err)}") - return - if text is None: - await message.answer( - "📊 Буфер аудита в Redis выключен (AUDIT_REDIS_BUFFER_ENABLED=False). " - "Используйте кнопку «Аудит (воронки)» в Статистике для отчёта из БД.", - ) - return - await message.answer(text, reply_markup=build_audit_refresh_kb()) + """По команде «Аудит» — предложить выбрать источник данных для отчёта.""" + help_text = ( + "📊 Аудит\n\n" + "Выберите источник данных:\n" + "• Redis raw — сырые последние события из буфера Redis, без добора успешных оплат из БД.\n" + "• БД вчера — агрегированный отчёт из базы за прошлые сутки: с 00:00 до 00:00 по Москве.\n\n" + "Сброс Redis очищает кэш аудита. Сброс БД применяется отдельно только при явном выборе в режиме БД." + ) + await message.answer(help_text, reply_markup=build_audit_source_kb()) @router.callback_query(AdminPanelCallback.filter(F.action == "audit_refresh"), IsAdminFilter()) async def handle_audit_refresh(callback_query: CallbackQuery, session: AsyncSession): - """Обновить отчёт аудита по нажатию кнопки «Обновить» под сообщением.""" + """Совместимость со старыми сообщениями аудита: открыть DB-отчёт за прошлые сутки.""" await callback_query.answer() - text, err = await _build_audit_report(session) + source = "db" + text, err = await _build_audit_report(session, source=source) if err: await callback_query.message.edit_text(f"❗ Ошибка: {escape(err)}") return - if text is None: - await callback_query.message.edit_text( - "📊 Буфер аудита в Redis выключен (AUDIT_REDIS_BUFFER_ENABLED=False). " - "Используйте кнопку «Аудит (воронки)» в Статистике для отчёта из БД.", - ) - return try: - await callback_query.message.edit_text(text, reply_markup=build_audit_refresh_kb()) + await callback_query.message.edit_text(text, reply_markup=build_audit_refresh_kb(source)) except TelegramBadRequest: pass +@router.callback_query(AdminPanelCallback.filter(F.action == "audit_refresh_redis"), IsAdminFilter()) +async def handle_audit_refresh_redis(callback_query: CallbackQuery, session: AsyncSession): + """Показать сырые данные аудита из Redis.""" + await callback_query.answer() + source = "redis" + text, err = await _build_audit_report(session, source=source) + if err: + await callback_query.message.edit_text(f"❗ Ошибка: {escape(err)}", reply_markup=build_audit_refresh_kb(source)) + return + try: + await callback_query.message.edit_text(text, reply_markup=build_audit_refresh_kb(source)) + except TelegramBadRequest: + pass + + +@router.callback_query(AdminPanelCallback.filter(F.action == "audit_refresh_db"), IsAdminFilter()) +async def handle_audit_refresh_db(callback_query: CallbackQuery, session: AsyncSession): + """Показать аудит из БД за прошлые московские сутки.""" + await callback_query.answer() + source = "db" + text, err = await _build_audit_report(session, source=source) + if err: + await callback_query.message.edit_text(f"❗ Ошибка: {escape(err)}", reply_markup=build_audit_refresh_kb(source)) + return + try: + await callback_query.message.edit_text(text, reply_markup=build_audit_refresh_kb(source)) + except TelegramBadRequest: + pass + + +@router.callback_query(AdminPanelCallback.filter(F.action.in_(["audit_reset_ask_redis", "audit_reset_ask_db"])), IsAdminFilter()) +async def handle_audit_reset_ask(callback_query: CallbackQuery): + await callback_query.answer() + source = "redis" if callback_query.data and "redis" in callback_query.data else "db" + source_label = "Redis raw" if source == "redis" else "БД вчера" + text = ( + f"🧹 Сброс аудита\n\n" + f"Источник: {source_label}\n" + "Подтвердите действие." + ) + if source == "redis": + text += "\n\nБудет очищен только Redis-буфер отчёта аудита. История по пользователям останется." + else: + text += "\n\nВ БД будет записан отдельный сброс, который применится к DB-отчёту." + try: + await callback_query.message.edit_text(text, reply_markup=build_audit_reset_confirm_kb(source)) + except TelegramBadRequest: + pass + + +@router.callback_query(AdminPanelCallback.filter(F.action.in_(["audit_reset_do_redis", "audit_reset_do_db"])), IsAdminFilter()) +async def handle_audit_reset_do(callback_query: CallbackQuery, session: AsyncSession): + source = "redis" if callback_query.data and "redis" in callback_query.data else "db" + try: + if source == "redis": + await clear_audit_redis_buffers() + else: + await set_audit_db_reset_at(session) + await callback_query.answer("Аудит сброшен") + text, err = await _build_audit_report(session, source=source) + if err: + await callback_query.message.edit_text(f"❗ Ошибка: {escape(err)}", reply_markup=build_audit_refresh_kb(source)) + return + try: + await callback_query.message.edit_text(text, reply_markup=build_audit_refresh_kb(source)) + except TelegramBadRequest: + pass + except Exception as e: + logger.exception("Ошибка при сбросе аудита (%s): %s", source, e) + await callback_query.answer("Не удалось сбросить аудит", show_alert=True) + + @router.callback_query(AdminPanelCallback.filter(F.action == "stats_export_users_csv"), IsAdminFilter()) async def handle_export_users_csv(callback_query: CallbackQuery, session: AsyncSession): kb = build_admin_back_kb("stats")