Глобально уникальные ID (docs/changes/016): - у устройств, групп, задач, резервных копий и записей журнала ID вида <префикс>_<uuid7> (dev_, grp_, job_, bkp_, evt_): типы не пересекаются, внутри типа ID не повторяются и сортируются по времени; - миграция БД v1 заменяет числовые ID с пересчётом ссылок (копия файла БД перед миграцией, сверка числа строк, одна транзакция); - резервная копия = пара файлов с одним bkp_ ID, файлы в S3 получают метаданные backup-id/device-id; синхронизация метаданных с бакетом; - журнал событий в БД (events): создание/изменение/удаление, задачи, бэкапы, смена online/offline, вход в UI; API чтения GET /api/v1/events. Журнал в UI, ротация и очистка (docs/changes/017): - страница «Журнал»: фильтры, подгрузка «Показать ещё», окно записи; - настройки ротации (срок и максимум записей) хранятся в БД (схема v2), ротация при старте, раз в час и после сохранения настроек; - очистка журнала только через окно с паролем пользователя, блокировка после 5 неверных попыток, остаётся запись о факте очистки; через API очистки нет. Тесты: 20 из 20. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
144 lines
6.8 KiB
Python
144 lines
6.8 KiB
Python
"""Журнал событий: каждая запись имеет свой ID (evt_…) и ссылается на ID сущности.
|
|
|
|
Актор и текущая задача передаются через ContextVar: их выставляют зависимости авторизации (UI/API),
|
|
фоновый опрос и запуск задач (asyncio наследует контекст при create_task).
|
|
"""
|
|
import json
|
|
import logging
|
|
from contextvars import ContextVar
|
|
from datetime import date, datetime, time as dtime, timedelta, timezone
|
|
|
|
from sqlalchemy import func, or_
|
|
|
|
from app.db import session_scope
|
|
from app.models import Event, now
|
|
|
|
log = logging.getLogger("ros_control.events")
|
|
|
|
_actor: ContextVar[str] = ContextVar("actor", default="system")
|
|
_job: ContextVar[str | None] = ContextVar("job_id", default=None)
|
|
|
|
|
|
def set_actor(actor: str) -> None:
|
|
"""ui:<пользователь> | api | poller | system."""
|
|
_actor.set(actor)
|
|
|
|
|
|
def set_job(job_id: str | None) -> None:
|
|
_job.set(job_id)
|
|
|
|
|
|
def current_job() -> str | None:
|
|
return _job.get()
|
|
|
|
|
|
def record(type_: str, entity_type: str, entity_id: str | None, message: str = "", *, device_id: str | None = None,
|
|
job_id: str | None = None, data: dict | None = None, s=None) -> str | None:
|
|
"""Записывает событие. Внутри транзакции вызывающего (s=сессия) — атомарно с изменением; без сессии —
|
|
отдельной короткой транзакцией, сбой записи журнала не ломает основную операцию."""
|
|
ev = Event(type=type_, entity_type=entity_type, entity_id=entity_id, device_id=device_id,
|
|
job_id=job_id or _job.get(), actor=_actor.get(), message=message,
|
|
data=json.dumps(data, ensure_ascii=False) if data is not None else None)
|
|
if s is not None:
|
|
s.add(ev)
|
|
s.flush()
|
|
return ev.id
|
|
try:
|
|
with session_scope() as sess:
|
|
sess.add(ev)
|
|
sess.flush()
|
|
return ev.id
|
|
except Exception: # noqa: BLE001
|
|
log.exception("не удалось записать событие %s", type_)
|
|
return None
|
|
|
|
|
|
def _filtered(q, *, entity_id=None, type_=None, device_id=None, job_id=None, actor=None,
|
|
date_from: date | None = None, date_to: date | None = None, text: str | None = None):
|
|
if entity_id:
|
|
q = q.filter(Event.entity_id == entity_id)
|
|
if type_:
|
|
q = q.filter(Event.type == type_) if "." in type_ else q.filter(Event.type.like(type_ + ".%"))
|
|
if device_id:
|
|
q = q.filter(Event.device_id == device_id)
|
|
if job_id:
|
|
q = q.filter(Event.job_id == job_id)
|
|
if actor: # «ui» находит ui:<пользователь>
|
|
q = q.filter(or_(Event.actor == actor, Event.actor.like(actor + ":%")))
|
|
if date_from:
|
|
q = q.filter(Event.ts >= datetime.combine(date_from, dtime.min, tzinfo=timezone.utc))
|
|
if date_to: # включительно, по UTC
|
|
q = q.filter(Event.ts < datetime.combine(date_to + timedelta(days=1), dtime.min, tzinfo=timezone.utc))
|
|
if text:
|
|
like = f"%{text.strip()}%"
|
|
q = q.filter(or_(Event.message.ilike(like), Event.id.ilike(like), Event.entity_id.ilike(like),
|
|
Event.device_id.ilike(like), Event.job_id.ilike(like)))
|
|
return q
|
|
|
|
|
|
def list_events(*, before: str | None = None, limit: int = 100, **filters) -> list[Event]:
|
|
"""Новые сверху. before — курсор: ID последней записи предыдущей страницы (ID сортируется по времени)."""
|
|
with session_scope() as s:
|
|
q = _filtered(s.query(Event), **filters)
|
|
if before:
|
|
q = q.filter(Event.id < before)
|
|
return list(q.order_by(Event.id.desc()).limit(max(1, min(limit, 500))))
|
|
|
|
|
|
def count(**filters) -> int:
|
|
with session_scope() as s:
|
|
return _filtered(s.query(func.count(Event.id)), **filters).scalar() or 0
|
|
|
|
|
|
def stats() -> dict:
|
|
"""Сводка по журналу: число записей и время самой старой."""
|
|
with session_scope() as s:
|
|
n, oldest = s.query(func.count(Event.id), func.min(Event.ts)).one()
|
|
return {"count": n or 0, "oldest": oldest}
|
|
|
|
|
|
def types() -> list[str]:
|
|
with session_scope() as s:
|
|
return sorted(t for (t,) in s.query(Event.type).distinct())
|
|
|
|
|
|
def get_event(event_id: str) -> Event:
|
|
with session_scope() as s:
|
|
ev = s.get(Event, event_id)
|
|
if ev is None:
|
|
raise LookupError(f"Запись журнала {event_id} не найдена")
|
|
return ev
|
|
|
|
|
|
def rotate() -> dict:
|
|
"""Ротация по настройкам (services.settings): удаляет записи старше срока и самые старые сверх лимита.
|
|
Если что-то удалено — одна запись journal.rotated. Возвращает {'by_age': n, 'by_count': m}."""
|
|
from app.services import settings # локально: settings импортирует events
|
|
|
|
cfg = settings.journal()
|
|
by_age = by_count = 0
|
|
with session_scope() as s:
|
|
if cfg["retention_days"] > 0:
|
|
cutoff = now() - timedelta(days=cfg["retention_days"])
|
|
by_age = s.query(Event).filter(Event.ts < cutoff).delete(synchronize_session=False)
|
|
if cfg["max_rows"] > 0:
|
|
extra = (s.query(func.count(Event.id)).scalar() or 0) - cfg["max_rows"]
|
|
if extra > 0:
|
|
extra += 1 # место под запись journal.rotated: после ротации в журнале ровно max_rows записей
|
|
# граница по ID: ID сортируется по времени, удаляем самые старые
|
|
bound = s.query(Event.id).order_by(Event.id).offset(extra - 1).limit(1).scalar()
|
|
by_count = s.query(Event).filter(Event.id <= bound).delete(synchronize_session=False)
|
|
if by_age or by_count:
|
|
record("journal.rotated", "journal", None, f"Ротация журнала: удалено записей {by_age + by_count}",
|
|
data={"by_age": by_age, "by_count": by_count, **cfg}, s=s)
|
|
return {"by_age": by_age, "by_count": by_count}
|
|
|
|
|
|
def clear(user: str) -> int:
|
|
"""Удаляет все записи и оставляет одну — о факте очистки (в одной транзакции). Возвращает число удалённых."""
|
|
with session_scope() as s:
|
|
n = s.query(Event).delete(synchronize_session=False)
|
|
record("journal.cleared", "journal", None, f"Журнал очищен пользователем {user}: удалено записей {n}",
|
|
data={"deleted": n, "user": user}, s=s)
|
|
return n
|