diff --git a/api/v2/routes/management.py b/api/v2/routes/management.py
index f19f7782..c84f636f 100644
--- a/api/v2/routes/management.py
+++ b/api/v2/routes/management.py
@@ -3,6 +3,7 @@ import os
import re
import subprocess
import sys
+from datetime import datetime, timezone, timedelta
from typing import Literal
import psutil
@@ -15,8 +16,12 @@ 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 api.v2.schemas.audit import (
+ AuditEventListResponse,
+ AuditEventResponse,
+ AuditStatsResponse,
+)
+from audit import drain_audit_redis_to_db, get_audit_funnel, get_audit_stats, 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
@@ -177,6 +182,60 @@ async def get_broadcast_clusters(
return {"clusters": clusters}
+def _parse_date_range(
+ date: str | None = None,
+ date_from: str | None = None,
+ date_to: str | None = None,
+) -> tuple[datetime, datetime]:
+ """Возвращает (date_from, date_to) в UTC. Либо date=YYYY-MM-DD (один день), либо date_from + date_to."""
+ tz = timezone.utc
+ if date:
+ try:
+ d = datetime.strptime(date, "%Y-%m-%d").date()
+ start = datetime(d.year, d.month, d.day, 0, 0, 0, tzinfo=tz)
+ end = start + timedelta(days=1)
+ return start, end
+ except ValueError:
+ raise HTTPException(status_code=400, detail="date должен быть YYYY-MM-DD")
+ if date_from and date_to:
+ try:
+ start = datetime.fromisoformat(date_from.replace("Z", "+00:00"))
+ end = datetime.fromisoformat(date_to.replace("Z", "+00:00"))
+ if start.tzinfo is None:
+ start = start.replace(tzinfo=tz)
+ if end.tzinfo is None:
+ end = end.replace(tzinfo=tz)
+ if start >= end:
+ raise HTTPException(status_code=400, detail="date_from должен быть раньше date_to")
+ return start, end
+ except ValueError as e:
+ raise HTTPException(status_code=400, detail=f"Неверный формат дат: {e}")
+ # по умолчанию — вчера
+ end = datetime.now(tz).replace(hour=0, minute=0, second=0, microsecond=0)
+ start = end - timedelta(days=1)
+ return start, end
+
+
+@router.get("/audit-stats", response_model=AuditStatsResponse)
+async def get_audit_stats_endpoint(
+ identity=Depends(verify_identity_admin),
+ session: AsyncSession = Depends(get_session),
+ date: str | None = Query(None, description="Один день: YYYY-MM-DD"),
+ date_from: str | None = Query(None, description="Начало периода (ISO)"),
+ date_to: str | None = Query(None, description="Конец периода (ISO)"),
+):
+ """Статистика аудита за период: какие пути отрабатывают хорошо/плохо, воронка старт→оплата.
+ Данные только из БД (события из Redis учитываются после drain)."""
+ start, end = _parse_date_range(date=date, date_from=date_from, date_to=date_to)
+ stats = await get_audit_stats(session, date_from=start, date_to=end)
+ funnel = await get_audit_funnel(session, date_from=start, date_to=end)
+ return AuditStatsResponse(
+ summary=stats["summary"],
+ by_path=stats["by_path"],
+ funnel=funnel,
+ )
+
+
@router.get("/audit-events", response_model=AuditEventListResponse)
async def get_audit_events_history(
identity=Depends(verify_identity_admin),
diff --git a/api/v2/schemas/audit.py b/api/v2/schemas/audit.py
index 13b7e354..6cd10c24 100644
--- a/api/v2/schemas/audit.py
+++ b/api/v2/schemas/audit.py
@@ -26,3 +26,33 @@ class AuditEventListResponse(BaseModel):
items: list[AuditEventResponse]
limit: int
offset: int
+
+
+class AuditPathStat(BaseModel):
+ step: str
+ label: str
+ total: int
+ success: int
+ fail: int
+ unique_users: int
+ fail_rate_pct: float
+
+
+class AuditFunnelStep(BaseModel):
+ step: str
+ label: str
+ count: int
+ conversion_from_prev_pct: float | None
+
+
+class AuditStatsSummary(BaseModel):
+ date_from: str
+ date_to: str
+ total_events: int
+ unique_users: int
+
+
+class AuditStatsResponse(BaseModel):
+ summary: AuditStatsSummary
+ by_path: list[AuditPathStat]
+ funnel: list[AuditFunnelStep]
diff --git a/audit.py b/audit.py
index a6ed01e7..a750913d 100644
--- a/audit.py
+++ b/audit.py
@@ -16,6 +16,24 @@ from sqlalchemy.ext.asyncio import AsyncSession
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("/")
+
try:
from core.cache_config import (
AUDIT_REDIS_BUFFER_ENABLED,
@@ -50,6 +68,13 @@ def new_request_id() -> str:
return uuid.uuid4().hex
+def _naive_utc(dt: datetime) -> datetime:
+ """Приводит datetime к naive UTC для запросов к колонкам DateTime (без timezone)."""
+ if dt.tzinfo is not None:
+ dt = dt.astimezone(timezone.utc)
+ return dt.replace(tzinfo=None)
+
+
def _trim(value: Any, limit: int = _MAX_TEXT_LEN) -> str | None:
if value is None:
return None
@@ -416,7 +441,9 @@ async def record_audit_event_to_redis(
reason=reason,
request_id=request_id,
)
- await cache_rpush(AUDIT_REDIS_FLUSH_KEY, record)
+ 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)
@@ -521,8 +548,8 @@ async def record_telegram_access_event_background(
result: str = "success",
reason: str | None = None,
) -> None:
- """Пишет событие telegram_access: в Redis-буфер (если включён) или в БД в отдельной сессии.
- При ошибке логирует и не пробрасывает исключение."""
+ """Пишет событие telegram_access: в Redis-буфер (если включён), иначе в БД.
+ При сбое записи в Redis — fallback в БД, чтобы события не терялись."""
if AUDIT_REDIS_BUFFER_ENABLED:
try:
await record_audit_event_to_redis(
@@ -533,9 +560,11 @@ async def record_telegram_access_event_background(
result=result,
reason=reason,
)
+ return
except Exception as exc:
- logger.warning("[Audit] Запись в Redis-буфер не удалась: %s", exc)
- return
+ logger.warning(
+ f"[Audit] Запись в Redis-буфер не удалась, пишем в БД: {exc}"
+ )
try:
async with session_factory() as session:
await ensure_audit_table(session)
@@ -686,6 +715,64 @@ async def _list_audit_events_from_redis(
return out[:max_events]
+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
+
+ raw = await cache_lrange(AUDIT_REDIS_FLUSH_KEY, -max_events, -1)
+ out = []
+ for rec in raw:
+ if not isinstance(rec, dict):
+ continue
+ out.append(_redis_record_to_event_like(rec))
+ out.sort(key=lambda e: e.created_at)
+ return out
+
+
+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:
+ 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)
+
+ by_path_list = []
+ for step, data in sorted(by_step.items(), key=lambda x: -x[1]["total"]):
+ total = data["total"]
+ 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,
+ })
+ return by_step, by_path_list, all_actors
+
+
async def list_audit_events(
session: AsyncSession,
*,
@@ -777,6 +864,332 @@ async def delete_old_audit_events(
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"
+
+
+async def get_audit_stats(
+ session: AsyncSession,
+ *,
+ date_from: datetime,
+ 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()
+ by_step, by_path_list, all_actors = _aggregate_audit_rows(rows)
+ return {
+ "summary": {
+ "date_from": date_from.isoformat(),
+ "date_to": date_to.isoformat(),
+ "total_events": sum(d["total"] for d in by_step.values()),
+ "unique_users": len(all_actors),
+ },
+ "by_path": by_path_list,
+ }
+
+
+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]
+ by_step, by_path_list, all_actors = _aggregate_audit_rows(rows)
+ return {
+ "summary": {
+ "source": "redis",
+ "total_events": len(events),
+ "unique_users": len(all_actors),
+ },
+ "by_path": by_path_list,
+ }
+
+
+async def get_audit_funnel(
+ session: AsyncSession,
+ *,
+ date_from: datetime,
+ date_to: datetime,
+ 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()
+ return _funnel_from_rows(rows, steps)
+
+
+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]] = {}
+ 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):
+ continue
+ key = (str(identity_id or ""), int(tg_id or 0))
+ if key == ("", 0):
+ 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
+
+ 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
+ return funnel_list
+
+
+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)
+ rows = [(e.path_or_handler, e.result, e.actor_tg_id, e.actor_identity_id) for e in events]
+ return _funnel_from_rows(rows, steps_ordered)
+
+
async def drain_audit_redis_to_db(session_factory: Any) -> int:
"""Выгружает буфер аудита из Redis в БД батчами. Вызывать по крону (например в 00:00).
Возвращает количество записанных событий."""
diff --git a/bot.py b/bot.py
index 05bece4a..da0ad4ae 100644
--- a/bot.py
+++ b/bot.py
@@ -15,15 +15,11 @@ from utils.modules_loader import load_modules_from_folder, modules_hub
apply_button_icons_patch()
bot = Bot(token=API_TOKEN, default=DefaultBotProperties(parse_mode=ParseMode.HTML))
-storage = MemoryStorage()
-try:
- RedisStorage = import_module("aiogram.fsm.storage.redis").RedisStorage
- redis_from_url = import_module("redis.asyncio").from_url
- redis = redis_from_url(REDIS_URL, encoding="utf-8", decode_responses=True)
- storage = RedisStorage(redis=redis)
-except Exception:
- storage = MemoryStorage()
+RedisStorage = import_module("aiogram.fsm.storage.redis").RedisStorage
+redis_from_url = import_module("redis.asyncio").from_url
+redis = redis_from_url(REDIS_URL, encoding="utf-8", decode_responses=True)
+storage = RedisStorage(redis=redis)
dp = Dispatcher(bot=bot, storage=storage)
diff --git a/core/redis_cache.py b/core/redis_cache.py
index 92e8bf2b..7401a8f6 100644
--- a/core/redis_cache.py
+++ b/core/redis_cache.py
@@ -4,6 +4,7 @@ from importlib import import_module
from typing import Any
from config import REDIS_URL
+from logger import logger
_REDIS_CLIENT = None
_REDIS_UNAVAILABLE_UNTIL = 0.0
@@ -17,6 +18,12 @@ def _now() -> float:
async def _get_redis() -> Any | None:
global _REDIS_CLIENT, _REDIS_UNAVAILABLE_UNTIL
+ try:
+ from bot import redis as bot_redis
+ if bot_redis is not None:
+ return bot_redis
+ except ImportError:
+ pass
if _REDIS_CLIENT is not None:
return _REDIS_CLIENT
if _REDIS_UNAVAILABLE_UNTIL > _now():
@@ -33,7 +40,11 @@ async def _get_redis() -> Any | None:
await client.ping()
_REDIS_CLIENT = client
return _REDIS_CLIENT
- except Exception:
+ except Exception as exc:
+ url_display = REDIS_URL.split("@")[-1] if "@" in REDIS_URL else REDIS_URL
+ logger.warning(
+ f"[Redis] Подключение не удалось ({url_display}): {exc}. Повтор через {_REDIS_BACKOFF_SEC} с."
+ )
_REDIS_UNAVAILABLE_UNTIL = _now() + _REDIS_BACKOFF_SEC
_REDIS_CLIENT = None
return None
@@ -141,6 +152,7 @@ async def cache_delete_pattern(pattern: str) -> int:
async def cache_rpush(key: str, *values: Any) -> int:
"""Добавляет значения в хвост списка. Значения сериализуются в JSON. Возвращает длину списка после или 0 при ошибке."""
+ global _REDIS_CLIENT
if not values:
return 0
client = await _get_redis()
@@ -149,7 +161,9 @@ async def cache_rpush(key: str, *values: Any) -> int:
try:
raw = [json.dumps(v, ensure_ascii=False) for v in values]
return int(await client.rpush(key, *raw))
- except Exception:
+ except Exception as exc:
+ logger.warning(f"[Redis] rpush({key}) не удался: {exc}")
+ _REDIS_CLIENT = None
return 0
diff --git a/handlers/admin/stats/keyboard.py b/handlers/admin/stats/keyboard.py
index 06480040..03e550ea 100644
--- a/handlers/admin/stats/keyboard.py
+++ b/handlers/admin/stats/keyboard.py
@@ -4,6 +4,13 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder
from ..panel.keyboard import AdminPanelCallback, build_admin_back_btn
+def build_audit_refresh_kb() -> InlineKeyboardMarkup:
+ """Клавиатура под сообщением аудита: кнопка «Обновить»."""
+ builder = InlineKeyboardBuilder()
+ builder.button(text="🔄 Обновить", callback_data=AdminPanelCallback(action="audit_refresh").pack())
+ return builder.as_markup()
+
+
def build_stats_kb() -> InlineKeyboardMarkup:
builder = InlineKeyboardBuilder()
builder.button(text="🔄 Обновить", callback_data=AdminPanelCallback(action="stats").pack())
diff --git a/handlers/admin/stats/stats_handler.py b/handlers/admin/stats/stats_handler.py
index 2bdc712a..88a2dfdc 100644
--- a/handlers/admin/stats/stats_handler.py
+++ b/handlers/admin/stats/stats_handler.py
@@ -1,5 +1,6 @@
from collections import Counter
from datetime import datetime, timedelta
+from html import escape
import pytz
@@ -8,6 +9,7 @@ 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 bot import bot
from config import ADMIN_ID
from database import (
@@ -38,7 +40,7 @@ from utils.csv_export import (
)
from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb
-from .keyboard import build_stats_kb
+from .keyboard import build_audit_refresh_kb, build_stats_kb
router = Router()
@@ -253,6 +255,160 @@ async def handle_stats(callback_query: CallbackQuery, session: AsyncSession):
await callback_query.answer("Произошла ошибка при получении статистики", show_alert=True)
+@router.callback_query(AdminPanelCallback.filter(F.action == "stats_audit"), IsAdminFilter())
+async def handle_stats_audit(callback_query: CallbackQuery, session: AsyncSession):
+ """Статистика аудита за вчера (МСК): объём по шагам, % ошибок, воронка."""
+ kb = build_admin_back_kb("stats")
+ try:
+ moscow_tz = pytz.timezone("Europe/Moscow")
+ now = datetime.now(moscow_tz)
+ yesterday_date = (now.date() - timedelta(days=1))
+ start = moscow_tz.localize(datetime.combine(yesterday_date, datetime.min.time()))
+ end = start + timedelta(days=1)
+ start_utc = start.astimezone(pytz.UTC)
+ end_utc = end.astimezone(pytz.UTC)
+
+ 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)
+
+ summary = stats["summary"]
+ by_path = stats["by_path"]
+
+ lines = [
+ f"📊 Аудит за {yesterday_date.strftime('%d.%m.%Y')} (МСК)",
+ "",
+ f"📎 Событий: {summary['total_events']} │ Уникальных пользователей: {summary['unique_users']}",
+ "",
+ "По шагам (топ по объёму):",
+ ]
+ 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))
+ 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}")
+
+ text = "\n".join(lines)
+ await callback_query.message.edit_text(
+ text,
+ reply_markup=kb,
+ )
+ await callback_query.answer()
+ except Exception as e:
+ logger.exception("Ошибка при получении статистики аудита: %s", e)
+ await callback_query.answer("Ошибка при загрузке статистики аудита", show_alert=True)
+ await callback_query.message.edit_text(
+ f"❗ Ошибка: {e}",
+ reply_markup=kb,
+ )
+
+
+async def _build_audit_report(session: AsyncSession) -> tuple[str | None, str | None]:
+ """Собирает текст отчёта аудита. Возвращает (text, error): при успехе error=None; при отключённом Redis text=None, error=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)
+ summary = stats["summary"]
+ by_path = stats["by_path"]
+ header = "📊 Аудит из БД (буфер Redis пуст; последние 24 ч)"
+ else:
+ header = "📊 Аудит из Redis (буфер, последние события)"
+ lines = [
+ header,
+ "",
+ f"📎 Событий: {summary['total_events']} │ Уникальных пользователей: {summary['unique_users']}",
+ "",
+ "По шагам (топ по объёму):",
+ ]
+ 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))
+ 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}")
+ return ("\n".join(lines), None)
+ except Exception as e:
+ logger.exception("Ошибка при получении аудита: %s", e)
+ return (None, str(e))
+
+
+@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())
+
+
+@router.callback_query(AdminPanelCallback.filter(F.action == "audit_refresh"), IsAdminFilter())
+async def handle_audit_refresh(callback_query: CallbackQuery, session: AsyncSession):
+ """Обновить отчёт аудита по нажатию кнопки «Обновить» под сообщением."""
+ await callback_query.answer()
+ text, err = await _build_audit_report(session)
+ 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())
+ except TelegramBadRequest:
+ pass
+
+
@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")
diff --git a/middlewares/loggings.py b/middlewares/loggings.py
index f7d42030..00988fae 100644
--- a/middlewares/loggings.py
+++ b/middlewares/loggings.py
@@ -8,11 +8,17 @@ from aiogram.types import CallbackQuery, InlineQuery, Message, TelegramObject, U
from audit import (
ensure_telegram_context,
log_telegram_access,
+ record_audit_event_to_redis,
record_telegram_access_event,
record_telegram_access_event_background,
set_telegram_actor,
_telegram_access_payload,
)
+try:
+ from core.cache_config import AUDIT_REDIS_BUFFER_ENABLED
+except ImportError:
+ AUDIT_REDIS_BUFFER_ENABLED = False
+
from logger import logger
@@ -33,8 +39,9 @@ def _log_activity_sync(user_info: UserInfo) -> None:
class LoggingMiddleware(BaseMiddleware):
- """Middleware для логирования действий пользователя. Лог пишется в фоне, не задерживая обработчик.
- Если передан sessionmaker — аудит пишется в отдельной сессии в фоне (не блокирует ответ)."""
+ """Middleware для логирования действий пользователя. Лог пишется в фоне.
+ Аудит: при включённом Redis-буфере пишем только в Redis (в БД — раз в сутки через drain);
+ при выключенном буфере или сбое Redis — пишем в БД."""
def __init__(self, sessionmaker=None):
super().__init__()
@@ -61,20 +68,32 @@ class LoggingMiddleware(BaseMiddleware):
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",
+ payload = _telegram_access_payload(audit_context, event, result="success")
+ if AUDIT_REDIS_BUFFER_ENABLED:
+ if self._sessionmaker is not None:
+ asyncio.create_task(
+ record_telegram_access_event_background(
+ self._sessionmaker, **payload
+ )
)
+ else:
+ asyncio.create_task(record_audit_event_to_redis(**payload))
+ else:
+ if self._sessionmaker is not None:
+ 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,
@@ -86,25 +105,35 @@ class LoggingMiddleware(BaseMiddleware):
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
+ payload = _telegram_access_payload(
+ audit_context, event, result="fail", reason=reason
+ )
+ if AUDIT_REDIS_BUFFER_ENABLED:
+ if self._sessionmaker is not None:
+ asyncio.create_task(
+ record_telegram_access_event_background(
+ self._sessionmaker, **payload
+ )
)
- )
+ else:
+ asyncio.create_task(record_audit_event_to_redis(**payload))
else:
- session = data.get("session")
- if session is not None:
- await record_telegram_access_event(
- session,
- audit_context,
- event,
- result="fail",
- reason=reason,
+ if self._sessionmaker is not None:
+ 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,