diff --git a/.env.example b/.env.example index 8c2471a..6aba488 100644 --- a/.env.example +++ b/.env.example @@ -27,3 +27,7 @@ MAX_CONCURRENCY=10 ROS_CONNECT_TIMEOUT=4 # секунд на установление соединения POLL_INTERVAL=30 # период фонового опроса, секунд (0 — выключить) UPDATE_CHECK_INTERVAL=1800 # как часто проверять обновления ROS в опросе, секунд + +# Ротация журнала событий: значения по умолчанию (действующие настройки меняются в UI: Журнал → Настройки) +EVENTS_RETENTION_DAYS=90 # хранить записи, дней (0 — без ограничения) +EVENTS_MAX_ROWS=100000 # максимум записей (0 — без ограничения) diff --git a/README.md b/README.md index aa52b14..93b0c62 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,15 @@ Admin Dashboard ⇄ Control Server ⇄ RouterOS REST (на каждом устр - Канал обновлений; обновление ROS (устройство перезагружается после скачивания); обновление FW (перезагрузка сразу после появления в журнале устройства записи «Firmware upgraded successfully…»). - Долгие и групповые операции выполняются задачами (карточка «Задачи», состояние — в UI и через API). +## Идентификаторы и журнал событий + +У каждой сущности глобально уникальный ID вида `<префикс>_`: `dev_…` устройство, `grp_…` группа, `job_…` задача, `bkp_…` резервная копия, `evt_…` запись журнала. Префикс исключает пересечение типов, UUIDv7 не повторяется и сортируется по времени создания; ID не переиспользуются после удаления. Числовых ID нет: `GET /api/v1/devices/dev_…`, ID чужого типа или неверного формата — 404. +Резервная копия — пара файлов (`.backup` + `.rsc`) с одним `bkp_` ID; файлы в бакете имеют метаданные `backup-id`/`device-id` (Yandex возвращает их как `Backup-Id`/`Device-Id`), в списке файлов API есть `backup_id`; копия, у которой пропали файлы, помечается удалённой, но остаётся в истории. +**Журнал событий** пишется в БД: создание, изменение и удаление устройств и групп, перенос устройств, жизненный цикл задач и бэкапов, смена online/offline, вход в UI, ротация и очистка журнала. Каждая запись содержит `evt_` ID, ID сущности, актора (`ui:<пользователь>`, `api`, `poller`, `system`) и данные. +В UI — страница «Журнал» (фильтры по типу, актору, устройству, периоду и поиску по сообщению/ID; подгрузка «Показать ещё 100»; клик по строке открывает запись с полными ID и данными). +- **Ротация**: «Журнал → Настройки» — срок хранения (дней) и максимум записей, `0` — без ограничения. Настройки хранятся в БД (значения по умолчанию — `EVENTS_RETENTION_DAYS=90`, `EVENTS_MAX_ROWS=100000` в `.env`); ротация выполняется при старте, раз в час и сразу после сохранения настроек и оставляет запись `journal.rotated`. +- **Очистка**: «Журнал → Очистить журнал» открывает окно с вводом пароля пользователя. Неверный пароль журнал не трогает (`journal.clear_denied`), 5 неверных попыток за 10 минут блокируют очистку на 10 минут. После успешной очистки остаётся одна запись `journal.cleared` (кто и сколько удалил). Очистки через API нет — только через это окно; API даёт чтение журнала и настроек. + ## Интерфейс Экран строится сверху вниз: приложение (навигация) → страница (заголовок, счётчики, главное действие) → вкладки групп → полоса инструментов таблицы → данные; «Задачи» — отдельная карточка. Пока ничего не выбрано, полоса показывает фильтры; при выборе строк — действия над выбранными. @@ -53,7 +62,7 @@ cp .env.example .env # заполнить SECRET_KEY, ADMIN_PASSWORD, API_TOK docker compose up -d --build # UI: http://localhost:8000, OpenAPI: /docs ``` -БД (SQLite) хранится в Docker-томе `ros_data` (`docker volume inspect ros_control_ros_data`), переживает пересборку контейнера; схема обновляется автоматически при старте. Остановка: `docker compose down` (том сохраняется; `down -v` удалит БД). +БД (SQLite) хранится в Docker-томе `ros_data` (`docker volume inspect ros_control_ros_data`), переживает пересборку контейнера; схема обновляется автоматически при старте (перед миграцией ID создаётся копия `ros_control.db.bak-<метка>` рядом с БД). Остановка: `docker compose down` (том сохраняется; `down -v` удалит БД). Ключ Fernet: `python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())"`. @@ -72,6 +81,7 @@ docker compose up -d --build # UI: http://localhost:8000, OpenAPI: /docs | `ROS_TIMEOUT` | `30` | секунд на ответ устройства (долгие операции) | | `POLL_INTERVAL` | `30` | период фонового опроса, с; `0` — выключить | | `UPDATE_CHECK_INTERVAL` | `1800` | как часто опрос проверяет обновления ROS, с | +| `EVENTS_RETENTION_DAYS` / `EVENTS_MAX_ROWS` | `90` / `100000` | ротация журнала по умолчанию (действующие значения меняются в UI); `0` — без ограничения | ### Разработка @@ -79,7 +89,7 @@ docker compose up -d --build # UI: http://localhost:8000, OpenAPI: /docs python3 -m venv venv && ./venv/bin/pip install -r requirements.txt set -a; . ./.env; set +a ./venv/bin/uvicorn app.main:app --reload -./venv/bin/python -m pytest # 12 тестов, фоновый опрос в тестах выключен +./venv/bin/python -m pytest # 20 тестов, фоновый опрос в тестах выключен ``` ## API v1 @@ -90,7 +100,7 @@ set -a; . ./.env; set +a curl -s -H "Authorization: Bearer $API_TOKEN" http://localhost:8000/api/v1/devices | python3 -m json.tool ``` -Устройство: `name`, `host`, `port`, `username`, `password` (только при записи), `verify_tls`, `use_tls` (`false` — HTTP), `group_id` (`null` — без группы), `note`. +Все `id` — строки (`dev_…`, `grp_…`, `job_…`); `group_id` — `grp_…`. Устройство: `name`, `host`, `port`, `username`, `password` (только при записи), `verify_tls`, `use_tls` (`false` — HTTP), `group_id` (`null` — без группы), `note`. | Метод | Путь | Назначение | |---|---|---| @@ -108,6 +118,8 @@ curl -s -H "Authorization: Bearer $API_TOKEN" http://localhost:8000/api/v1/devic | GET/DELETE | `/api/v1/backups`, `/backups/download?key=` | бэкапы в бакете (фильтры `device_id`, `group`, `kind`, `date_from`, `date_to`, `q`) | | POST | `/api/v1/backups/delete` | групповое удаление файлов: `{"keys": [...]}` → `{"deleted": N, "failed": M}` | | GET | `/api/v1/jobs`, `/jobs/{id}` | состояние задач | +| GET | `/api/v1/events`, `/events/{id}` | журнал событий (фильтры `entity_id`, `type`, `device_id`, `job_id`, `actor`, `date_from`, `date_to`, `q`; постранично `before=`, `limit`) | +| GET/PUT | `/api/v1/events/settings` | настройки ротации журнала `{"retention_days", "max_rows"}` и сводка (очистки журнала через API нет) | Внешние справочники: [RouterOS REST API](https://help.mikrotik.com/docs/spaces/ROS/pages/47579162/REST+API#RESTAPI-HTTPMethods), [Yandex Object Storage S3 API](https://yandex.cloud/ru/docs/storage/s3/api-ref/). @@ -138,3 +150,5 @@ curl -s -H "Authorization: Bearer $API_TOKEN" http://localhost:8000/api/v1/devic - `013-bulk-delete-backups` — выбор файлов чекбоксами и групповое удаление бэкапов. - `014-chr-support` — поддержка CHR (нет `/system/routerboard`), строгое сравнение версий ROS, причина недоступности в таблице. - `015-immutable-device-name` — имя устройства задаётся только при создании. +- `016-unique-ids-and-event-log` — глобально уникальные ID (префикс + UUIDv7) для всех сущностей, журнал событий в БД, миграция числовых ID. +- `017-events-ui-rotation` — страница «Журнал» в UI, настройки ротации, очистка журнала с подтверждением паролем. diff --git a/app/api/v1.py b/app/api/v1.py index d62c621..4af97c2 100644 --- a/app/api/v1.py +++ b/app/api/v1.py @@ -1,14 +1,15 @@ +import json from datetime import date, datetime from typing import Literal from fastapi import APIRouter, Depends from fastapi.responses import RedirectResponse -from pydantic import BaseModel, Field +from pydantic import BaseModel, Field, field_validator -from app import s3 +from app import ids, s3 from app.ros.operations import CHANNELS from app.security import require_api_token -from app.services import backups, devices, groups, jobs, ops +from app.services import backups, devices, events, groups, jobs, ops, settings router = APIRouter(prefix="/api/v1", dependencies=[Depends(require_api_token)], tags=["api"]) @@ -21,7 +22,7 @@ class DeviceIn(BaseModel): password: str verify_tls: bool = False use_tls: bool = True # False = HTTP (сервис www, порт 80): пароль идёт открытым текстом - group_id: int | None = None + group_id: str | None = None # grp_… note: str | None = None @@ -33,20 +34,20 @@ class DevicePatch(BaseModel): password: str | None = None verify_tls: bool | None = None use_tls: bool | None = None - group_id: int | None = None # явный null — открепить от группы + group_id: str | None = None # grp_…; явный null — открепить от группы note: str | None = None # пустая строка или null — очистить примечание class DeviceOut(BaseModel): """Пароль устройства в ответы не попадает.""" - id: int + id: str # dev_… name: str host: str port: int username: str verify_tls: bool use_tls: bool - group_id: int | None + group_id: str | None note: str | None online: bool | None status: dict @@ -60,8 +61,8 @@ class DeviceOut(BaseModel): class JobOut(BaseModel): - id: int - device_id: int | None + id: str # job_… + device_id: str | None device_name: str type: str status: str @@ -81,25 +82,25 @@ class GroupIn(BaseModel): class GroupOut(BaseModel): - id: int + id: str # grp_… name: str device_count: int = 0 class BatchIn(BaseModel): """Устройства, над которыми выполняется операция: список id и/или целая группа.""" - device_ids: list[int] = Field(default_factory=list) - group_id: int | None = None + device_ids: list[str] = Field(default_factory=list) # dev_… + group_id: str | None = None # grp_… - def resolve(self) -> list[int]: - ids = list(self.device_ids) + def resolve(self) -> list[str]: + found = list(self.device_ids) if self.group_id is not None: groups.get_group(self.group_id) # 404, если группы нет - ids += [d.id for d in devices.list_devices() if d.group_id == self.group_id] - ids = list(dict.fromkeys(ids)) - if not ids: + found += [d.id for d in devices.list_devices() if d.group_id == self.group_id] + found = list(dict.fromkeys(found)) + if not found: raise ValueError("Не найдено ни одного устройства: укажите device_ids или непустую group_id") - return ids + return found class BatchChannelIn(BatchIn): @@ -110,7 +111,7 @@ class BatchChannelIn(BatchIn): @router.get("/devices", response_model=list[DeviceOut]) async def list_devices( - group: str = "", # "" — все, "none" — без группы, иначе id группы + group: str = "", # "" — все, "none" — без группы, иначе ID группы (grp_…) q: str = "", status: Literal["", "online", "offline"] = "", updates: Literal["", "ros", "fw", "any", "none"] = "", @@ -134,13 +135,13 @@ async def create_group(body: GroupIn): @router.patch("/groups/{group_id}", response_model=GroupOut) -async def rename_group(group_id: int, body: GroupIn): +async def rename_group(group_id: str, body: GroupIn): groups.rename_group(group_id, body.name) return next(g for g in groups.list_groups() if g["id"] == group_id) @router.delete("/groups/{group_id}", status_code=204) -async def delete_group(group_id: int): +async def delete_group(group_id: str): groups.delete_group(group_id) @@ -150,22 +151,22 @@ async def create_device(body: DeviceIn): @router.get("/devices/{device_id}", response_model=DeviceOut) -async def get_device(device_id: int): +async def get_device(device_id: str): return DeviceOut.of(devices.get_device(device_id)) @router.patch("/devices/{device_id}", response_model=DeviceOut) -async def patch_device(device_id: int, body: DevicePatch): +async def patch_device(device_id: str, body: DevicePatch): return DeviceOut.of(devices.update_device(device_id, **body.model_dump(exclude_unset=True))) @router.delete("/devices/{device_id}", status_code=204) -async def delete_device(device_id: int): +async def delete_device(device_id: str): devices.delete_device(device_id) @router.post("/devices/{device_id}/refresh", response_model=DeviceOut) -async def refresh_device(device_id: int): +async def refresh_device(device_id: str): devices.get_device(device_id) await ops.refresh_status(device_id) return DeviceOut.of(devices.get_device(device_id)) @@ -180,24 +181,24 @@ async def refresh_all(): # --- операции (долгие — через jobs) --- @router.post("/devices/{device_id}/backups", status_code=202) -async def create_backup(device_id: int): +async def create_backup(device_id: str): return {"job_ids": jobs.start_jobs("backup", [device_id])} @router.put("/devices/{device_id}/update/channel", response_model=DeviceOut) -async def put_channel(device_id: int, body: ChannelIn): +async def put_channel(device_id: str, body: ChannelIn): devices.get_device(device_id) await ops.set_channel(device_id, body.channel) return DeviceOut.of(devices.get_device(device_id)) @router.post("/devices/{device_id}/update/install", status_code=202) -async def install_update(device_id: int): +async def install_update(device_id: str): return {"job_ids": jobs.start_jobs("ros_update", [device_id])} @router.post("/devices/{device_id}/firmware/upgrade", status_code=202) -async def upgrade_firmware(device_id: int): +async def upgrade_firmware(device_id: str): return {"job_ids": jobs.start_jobs("fw_update", [device_id])} @@ -209,20 +210,20 @@ async def batch(action: Literal["backup", "ros_update", "fw_update"], body: Batc @router.put("/batch/channel", response_model=list[DeviceOut]) async def batch_channel(body: BatchChannelIn): - ids = body.resolve() - for i in ids: + targets = body.resolve() + for i in targets: devices.get_device(i) - for i in ids: + for i in targets: await ops.set_channel(i, body.channel) - return [DeviceOut.of(devices.get_device(i)) for i in ids] + return [DeviceOut.of(devices.get_device(i)) for i in targets] # --- резервные копии в S3 --- @router.get("/backups") async def list_backups( - device_id: int | None = None, - group: str = "", # "" — все, "none" — без группы, иначе id группы + device_id: str | None = None, + group: str = "", # "" — все, "none" — без группы, иначе ID группы (grp_…) kind: Literal["", "backup", "rsc"] = "", date_from: date | None = None, date_to: date | None = None, @@ -252,9 +253,65 @@ async def delete_backups(body: KeysIn): @router.delete("/backups", status_code=204) async def delete_backup(key: str): - if not s3.key_allowed(key): - raise ValueError("Недопустимый ключ") - await s3.delete_object(key) + deleted, failed = await backups.delete_many([key]) + if failed: + raise RuntimeError("Не удалось удалить файл из бакета") + + +# --- журнал событий --- + +class EventOut(BaseModel): + id: str # evt_… + ts: datetime + type: str + entity_type: str + entity_id: str | None + device_id: str | None + job_id: str | None + actor: str + message: str + data: dict | None = None + + model_config = {"from_attributes": True} + + @field_validator("data", mode="before") + @classmethod + def _parse(cls, v): + return json.loads(v) if isinstance(v, str) else v + + +class JournalSettingsIn(BaseModel): + retention_days: int # 0 — без ограничения по сроку + max_rows: int # 0 — без ограничения по числу записей + + +@router.get("/events/settings") +async def get_journal_settings(): + """Настройки ротации журнала и сводка (число записей, самая старая запись).""" + return {**settings.journal(), **events.stats()} + + +@router.put("/events/settings") +async def put_journal_settings(body: JournalSettingsIn): + """Сохраняет настройки ротации и сразу применяет ротацию. Очистки журнала через API нет — она выполняется + только в UI через окно с вводом пароля пользователя.""" + return {**settings.save_journal(body.retention_days, body.max_rows), **events.stats()} + + +@router.get("/events", response_model=list[EventOut]) +async def list_events(entity_id: str | None = None, type: str | None = None, device_id: str | None = None, + job_id: str | None = None, actor: str | None = None, date_from: date | None = None, + date_to: date | None = None, q: str | None = None, before: str | None = None, limit: int = 100): + """Журнал событий, новые сверху. type — точный тип (device.created) или группа (device, job, backup); + actor — ui, api, poller, system или anonymous. Постранично: before=.""" + return events.list_events(entity_id=entity_id, type_=type, device_id=device_id, job_id=job_id, actor=actor, + date_from=date_from, date_to=date_to, text=q, before=before, limit=limit) + + +@router.get("/events/{event_id}", response_model=EventOut) +async def get_event(event_id: str): + ids.check(event_id, "evt") + return events.get_event(event_id) # --- задачи --- @@ -265,5 +322,5 @@ async def list_jobs(): @router.get("/jobs/{job_id}", response_model=JobOut) -async def get_job(job_id: int): +async def get_job(job_id: str): return jobs.get_job(job_id) diff --git a/app/config.py b/app/config.py index 4ac9d8b..c4a3e3b 100644 --- a/app/config.py +++ b/app/config.py @@ -27,6 +27,10 @@ class Settings(BaseSettings): poll_interval: int = 30 # период фонового опроса устройств, с; 0 — отключить update_check_interval: int = 1800 # как часто в опросе проверять обновления ROS (тяжёлый запрос), с + # значения по умолчанию для ротации журнала; действующие настройки хранятся в БД и меняются из UI + events_retention_days: int = 90 # 0 — без ограничения по сроку + events_max_rows: int = 100000 # 0 — без ограничения по числу записей + @lru_cache def get_settings() -> Settings: diff --git a/app/db.py b/app/db.py index cdca38d..9fe6f4d 100644 --- a/app/db.py +++ b/app/db.py @@ -25,9 +25,14 @@ def init_db(url: str | None = None) -> None: _SessionLocal = sessionmaker(_engine, expire_on_commit=False) from app import models # noqa: F401 (регистрация моделей) - Base.metadata.create_all(_engine) + Base.metadata.create_all(_engine) # у новой БД — все таблицы, у старой — только недостающие (events) _migrate() + from app import migrations + + db_file = Path(url.removeprefix("sqlite:///")) if url.startswith("sqlite:///") and ":memory:" not in url else None + migrations.run(_engine, db_file) + def _migrate() -> None: """create_all не добавляет колонки в существующие таблицы — докидываем вручную.""" @@ -36,7 +41,7 @@ def _migrate() -> None: if "use_tls" not in cols: conn.exec_driver_sql("ALTER TABLE devices ADD COLUMN use_tls BOOLEAN NOT NULL DEFAULT 1") if "group_id" not in cols: - conn.exec_driver_sql("ALTER TABLE devices ADD COLUMN group_id INTEGER") + conn.exec_driver_sql("ALTER TABLE devices ADD COLUMN group_id VARCHAR(40)") if "note" not in cols: conn.exec_driver_sql("ALTER TABLE devices ADD COLUMN note TEXT") diff --git a/app/ids.py b/app/ids.py new file mode 100644 index 0000000..b9f22a3 --- /dev/null +++ b/app/ids.py @@ -0,0 +1,74 @@ +"""Глобально уникальные идентификаторы сущностей: <префикс>_. + +Пример: dev_0192f3a1-7c4e-7b1d-9a3f-5e2c1d0b8a47. +- Между типами — префикс: ID разных сущностей не пересекаются по построению. +- Внутри типа — UUIDv7 (RFC 9562): 48 бит времени (мс) + 12-битный счётчик + 62 случайных бита. + Внутри процесса ID строго возрастают (в том числе при откате часов), поэтому сортировка по ID + совпадает с порядком создания; между процессами уникальность держится на 62 случайных битах и PRIMARY KEY. +- ID не переиспользуются: значение не зависит от того, какие записи уже удалены. +""" +import re +import secrets +import threading +import time +import uuid +from datetime import datetime, timezone + +PREFIXES = { + "dev": "устройство", + "grp": "группа", + "job": "задача", + "bkp": "резервная копия", + "evt": "запись журнала", +} +_TAIL = r"[0-9a-f]{8}-[0-9a-f]{4}-7[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}" +_lock = threading.Lock() +_last_ms = 0 +_seq = 0 +_hist: dict[int, int] = {} # счётчики для «исторических» меток времени (миграция) + + +def _build(ms: int, seq: int) -> uuid.UUID: + rand_b = secrets.randbits(62) + return uuid.UUID(int=(ms << 80) | (0x7 << 76) | (seq << 64) | (0b10 << 62) | rand_b) + + +def new_id(prefix: str, at: datetime | None = None) -> str: + """Новый ID. at — момент создания записи (для миграции старых данных: сохраняет порядок).""" + global _last_ms, _seq + if prefix not in PREFIXES: + raise ValueError(f"Неизвестный префикс ID: {prefix}") + with _lock: + if at is not None: + ms = int(at.replace(tzinfo=at.tzinfo or timezone.utc).timestamp() * 1000) + _hist[ms] = _hist.get(ms, -1) + 1 + return f"{prefix}_{_build(ms, min(_hist[ms], 0xFFF))}" + now_ms = time.time_ns() // 1_000_000 + if now_ms > _last_ms: + _last_ms, _seq = now_ms, secrets.randbits(11) # случайный старт: ≥2048 ID на миллисекунду + else: # тот же миллисекунда или часы ушли назад — время не уменьшаем + _seq += 1 + if _seq > 0xFFF: + _last_ms, _seq = _last_ms + 1, 0 + return f"{prefix}_{_build(_last_ms, _seq)}" + + +def is_id(value: object, prefix: str | None = None) -> bool: + """Строка похожа на ID (нужного типа, если prefix задан).""" + if not isinstance(value, str): + return False + pat = rf"({'|'.join(PREFIXES)})_{_TAIL}" if prefix is None else rf"{prefix}_{_TAIL}" + return re.fullmatch(pat, value) is not None + + +def check(value: str, prefix: str) -> str: + """Проверка ID из пути/запроса: чужой тип или мусор — LookupError (в API это 404 с пояснением).""" + if not is_id(value, prefix): + hint = f"ожидается {PREFIXES[prefix]} ({prefix}_…)" + raise LookupError(f"Неверный идентификатор «{str(value)[:60]}»: {hint}") + return value + + +def short_id(value: str | None) -> str: + """Короткий вид для таблиц: последние 8 символов (случайная часть); полный ID — в подсказке.""" + return (value or "")[-8:] diff --git a/app/main.py b/app/main.py index 062b676..36b218b 100644 --- a/app/main.py +++ b/app/main.py @@ -11,7 +11,7 @@ from starlette.middleware.sessions import SessionMiddleware from app.api import v1 from app.config import get_settings from app.db import init_db -from app.services import jobs, poller +from app.services import backups, jobs, poller, rotation from app.ui import routes as ui @@ -20,11 +20,14 @@ async def lifespan(app: FastAPI): init_db() jobs.fail_stale_jobs() task = asyncio.create_task(poller.run()) if get_settings().poll_interval > 0 else None + reconcile = asyncio.create_task(backups.reconcile()) # связать файлы бакета с метаданными (best-effort) + rotator = asyncio.create_task(rotation.run()) # ротация журнала по настройкам yield - if task: - task.cancel() - with suppress(asyncio.CancelledError): - await task + for t in (task, reconcile, rotator): + if t: + t.cancel() + with suppress(asyncio.CancelledError): + await t def create_app() -> FastAPI: diff --git a/app/migrations.py b/app/migrations.py new file mode 100644 index 0000000..c9be774 --- /dev/null +++ b/app/migrations.py @@ -0,0 +1,139 @@ +"""Миграции схемы SQLite. Версия хранится в PRAGMA user_version. + +v1 — глобально уникальные ID: числовые ID таблиц (devices, device_groups, jobs, backups) заменяются на +<префикс>_ (см. app/ids.py) с пересчётом ссылок; порядок записей сохраняется (время в UUIDv7 берётся из +created_at/requested_at). Перед изменениями делается копия файла БД (<файл>.bak-<метка>); всё выполняется +одной транзакцией на «сыром» соединении (DDL в pysqlite иначе не транзакционен) и при любом расхождении в числе строк откатывается. +""" +import json +import logging +import sqlite3 +from datetime import datetime, timezone +from pathlib import Path + +from sqlalchemy import Engine +from sqlalchemy.dialects import sqlite as sqlite_dialect +from sqlalchemy.schema import CreateTable + +from app.ids import new_id + +log = logging.getLogger("ros_control.migrations") +SCHEMA_VERSION = 2 # v2: таблица app_settings (создаётся create_all), данные не меняются +TABLES = ("device_groups", "devices", "backups", "jobs") + + +def _dt(value) -> datetime | None: + try: + return datetime.fromisoformat(str(value)) if value else None + except ValueError: + return None + + +def _int(value) -> int | None: + try: + return int(value) + except (TypeError, ValueError): + return None + + +def _is_legacy(con: sqlite3.Connection) -> bool: + cols = {r[1]: (r[2] or "").upper() for r in con.execute("PRAGMA table_info(devices)")} + return bool(cols) and cols.get("id", "").startswith("INT") + + +def run(engine: Engine, db_path: Path | None) -> None: + """Приводит БД к текущей версии схемы. Идемпотентна.""" + if db_path is None: # не файловая БД (:memory:) — старых данных быть не может + with engine.begin() as c: + c.exec_driver_sql(f"PRAGMA user_version = {SCHEMA_VERSION}") + return + engine.dispose() # закрыть пул: дальше работаем отдельным соединением + con = sqlite3.connect(db_path, isolation_level=None) + try: + version = con.execute("PRAGMA user_version").fetchone()[0] + if version >= SCHEMA_VERSION: + return + if version < 1 and _is_legacy(con): + _to_v1(con, db_path) + con.execute(f"PRAGMA user_version = {SCHEMA_VERSION}") + finally: + con.close() + + +def _to_v1(con: sqlite3.Connection, db_path: Path) -> None: + from app.db import Base # noqa: F401 + from app import models + + stamp = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + bak = db_path.with_name(f"{db_path.name}.bak-{stamp}") + dst = sqlite3.connect(bak) + con.backup(dst) + dst.close() + log.warning("Миграция ID: копия БД сохранена в %s", bak) + + con.row_factory = sqlite3.Row + con.execute("PRAGMA foreign_keys=OFF") + con.execute("BEGIN IMMEDIATE") + try: + old = {t: con.execute(f"SELECT * FROM {t} ORDER BY id").fetchall() for t in TABLES} + for t in TABLES: + con.execute(f"ALTER TABLE {t} RENAME TO {t}_old") + for tbl in (models.Group.__table__, models.Device.__table__, models.Backup.__table__, models.Job.__table__): + con.execute(str(CreateTable(tbl).compile(dialect=sqlite_dialect.dialect()))) + + grp: dict[int, str] = {} + for r in old["device_groups"]: + grp[r["id"]] = new_id("grp") + con.execute("INSERT INTO device_groups (id, name) VALUES (?, ?)", (grp[r["id"]], r["name"])) + + dev: dict[int, str] = {} + names: dict[int, str] = {} + for r in old["devices"]: + dev[r["id"]] = new_id("dev", _dt(r["created_at"])) + names[r["id"]] = r["name"] + con.execute( + "INSERT INTO devices (id, name, host, port, username, password_enc, verify_tls, use_tls, group_id, note," + " created_at, online, status_json, status_at, last_error, last_backup_requested_at)" + " VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", + (dev[r["id"]], r["name"], r["host"], r["port"], r["username"], r["password_enc"], r["verify_tls"], + r["use_tls"], grp.get(_int(r["group_id"])), r["note"], r["created_at"], r["online"], r["status_json"], + r["status_at"], r["last_error"], r["last_backup_requested_at"])) + + def dev_id(old_id): # ссылка на уже удалённое устройство сохраняет идентичность (один новый ID на старый) + return None if old_id is None else dev.setdefault(old_id, new_id("dev")) + + for r in old["backups"]: + key = r["key_binary"] or r["key_rsc"] or "" + parts = key.split("/") + name = names.get(_int(r["device_id"])) or (parts[1] if len(parts) >= 3 else "") + con.execute( + "INSERT INTO backups (id, device_id, device_name, job_id, deleted_at, requested_at, status, key_binary, key_rsc, error)" + " VALUES (?,?,?,?,?,?,?,?,?,?)", + (new_id("bkp", _dt(r["requested_at"])), dev_id(_int(r["device_id"])), name, None, None, r["requested_at"], + r["status"], r["key_binary"], r["key_rsc"], r["error"])) + for r in old["jobs"]: + con.execute( + "INSERT INTO jobs (id, device_id, device_name, type, status, message, created_at, finished_at)" + " VALUES (?,?,?,?,?,?,?,?)", + (new_id("job", _dt(r["created_at"])), dev_id(_int(r["device_id"])), r["device_name"], r["type"], r["status"], + r["message"], r["created_at"], r["finished_at"])) + + counts = {t: len(old[t]) for t in TABLES} + for t in TABLES: # сверка числа строк до коммита + n = con.execute(f"SELECT count(*) FROM {t}").fetchone()[0] + if n != counts[t]: + raise RuntimeError(f"Миграция ID: в {t} было {counts[t]} строк, стало {n}") + for t in TABLES: + con.execute(f"DROP TABLE {t}_old") + con.execute( + "INSERT INTO events (id, ts, type, entity_type, entity_id, device_id, job_id, actor, message, data)" + " VALUES (?,?,?,?,?,?,?,?,?,?)", + (new_id("evt"), datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S.%f"), "system.migrated", "system", None, + None, None, "system", "Числовые ID заменены на глобально уникальные (dev_/grp_/job_/bkp_)", + json.dumps({"rows": counts, "backup_file": bak.name}))) + con.execute("COMMIT") + except Exception: + con.execute("ROLLBACK") + log.exception("Миграция ID не выполнена, БД не изменена (копия: %s)", bak) + raise + log.warning("Миграция ID выполнена: %s", counts) diff --git a/app/models.py b/app/models.py index d7b22b5..ed735c3 100644 --- a/app/models.py +++ b/app/models.py @@ -5,23 +5,29 @@ from sqlalchemy import Boolean, DateTime, ForeignKey, Integer, String, Text from sqlalchemy.orm import Mapped, mapped_column from app.db import Base +from app.ids import new_id def now() -> datetime: return datetime.now(timezone.utc) +def _pk(prefix: str): + """Первичный ключ-строка: глобально уникальный ID _ (см. app/ids.py).""" + return mapped_column(String(40), primary_key=True, default=lambda: new_id(prefix)) + + class Group(Base): __tablename__ = "device_groups" - id: Mapped[int] = mapped_column(primary_key=True) + id: Mapped[str] = _pk("grp") name: Mapped[str] = mapped_column(String(64), unique=True) class Device(Base): __tablename__ = "devices" - id: Mapped[int] = mapped_column(primary_key=True) + id: Mapped[str] = _pk("dev") name: Mapped[str] = mapped_column(String(64), unique=True) host: Mapped[str] = mapped_column(String(255)) port: Mapped[int] = mapped_column(Integer, default=443) @@ -29,7 +35,7 @@ class Device(Base): password_enc: Mapped[str] = mapped_column(Text) verify_tls: Mapped[bool] = mapped_column(Boolean, default=False) use_tls: Mapped[bool] = mapped_column(Boolean, default=True, server_default="1") - group_id: Mapped[int | None] = mapped_column(Integer, nullable=True) # FK не включены — см. groups.delete_group + group_id: Mapped[str | None] = mapped_column(String(40), nullable=True) # FK не включены — см. groups.delete_group note: Mapped[str | None] = mapped_column(Text, nullable=True) # свободное примечание к устройству created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=now) @@ -48,10 +54,15 @@ class Device(Base): class Backup(Base): + """Резервная копия = пара файлов (.backup + .rsc). Строка не удаляется: при удалении файлов из бакета + ставится deleted_at, ID остаётся в истории.""" __tablename__ = "backups" - id: Mapped[int] = mapped_column(primary_key=True) - device_id: Mapped[int] = mapped_column(ForeignKey("devices.id", ondelete="CASCADE")) + id: Mapped[str] = _pk("bkp") + device_id: Mapped[str | None] = mapped_column(String(40), ForeignKey("devices.id", ondelete="SET NULL"), nullable=True) + device_name: Mapped[str] = mapped_column(String(64), default="") # снимок имени: устройство могут удалить + job_id: Mapped[str | None] = mapped_column(String(40), nullable=True) # задача, создавшая копию + deleted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) requested_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=now) status: Mapped[str] = mapped_column(String(16), default="running") # running|done|failed key_binary: Mapped[str | None] = mapped_column(String(512), nullable=True) @@ -62,11 +73,36 @@ class Backup(Base): class Job(Base): __tablename__ = "jobs" - id: Mapped[int] = mapped_column(primary_key=True) - device_id: Mapped[int | None] = mapped_column(Integer, nullable=True) + id: Mapped[str] = _pk("job") + device_id: Mapped[str | None] = mapped_column(String(40), nullable=True) device_name: Mapped[str] = mapped_column(String(64), default="") type: Mapped[str] = mapped_column(String(32)) # backup|ros_update|fw_update status: Mapped[str] = mapped_column(String(16), default="pending") # pending|running|done|failed message: Mapped[str] = mapped_column(Text, default="") created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=now) finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + +class AppSetting(Base): + """Настройки, редактируемые из UI (ключ → JSON-значение); значения по умолчанию — из .env.""" + __tablename__ = "app_settings" + + key: Mapped[str] = mapped_column(String(64), primary_key=True) + value: Mapped[str] = mapped_column(Text) + updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=now) + + +class Event(Base): + """Журнал событий: каждая запись имеет свой ID и ссылается на ID сущности.""" + __tablename__ = "events" + + id: Mapped[str] = _pk("evt") + ts: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=now, index=True) + type: Mapped[str] = mapped_column(String(48), index=True) # device.created, job.failed, … + entity_type: Mapped[str] = mapped_column(String(16)) # device|group|job|backup|session|system + entity_id: Mapped[str | None] = mapped_column(String(40), nullable=True, index=True) + device_id: Mapped[str | None] = mapped_column(String(40), nullable=True, index=True) + job_id: Mapped[str | None] = mapped_column(String(40), nullable=True) + actor: Mapped[str] = mapped_column(String(80), default="system") # ui:<пользователь> | api | poller | system + message: Mapped[str] = mapped_column(Text, default="") + data: Mapped[str | None] = mapped_column(Text, nullable=True) # JSON с подробностями diff --git a/app/s3.py b/app/s3.py index 1a27480..c6f0c91 100644 --- a/app/s3.py +++ b/app/s3.py @@ -50,9 +50,11 @@ async def list_backups(device_name: str | None = None) -> list[dict]: return await asyncio.to_thread(_list, prefix) -async def upload_file(path, key: str) -> None: - """PutObject: загрузка файла в бакет от имени Control Server.""" - await asyncio.to_thread(_client().upload_file, str(path), get_settings().s3_bucket, key) +async def upload_file(path, key: str, metadata: dict[str, str] | None = None) -> None: + """PutObject: загрузка файла в бакет от имени Control Server. metadata — пользовательские метаданные объекта + (x-amz-meta-*): сюда кладутся ID резервной копии и устройства.""" + extra = {"ExtraArgs": {"Metadata": metadata}} if metadata else {} + await asyncio.to_thread(_client().upload_file, str(path), get_settings().s3_bucket, key, **extra) async def delete_object(key: str) -> None: diff --git a/app/security.py b/app/security.py index 7c23f99..17cde2e 100644 --- a/app/security.py +++ b/app/security.py @@ -1,4 +1,6 @@ import hmac +import threading +import time from cryptography.fernet import Fernet from fastapi import Depends, HTTPException @@ -31,7 +33,49 @@ def check_admin(user: str, password: str) -> bool: return ok_user and ok_pass -def require_api_token(cred: HTTPAuthorizationCredentials | None = Depends(_bearer)) -> None: +# --- подтверждение действий паролем (очистка журнала) и защита от перебора --- +FAIL_LIMIT, WINDOW_S, LOCK_S = 5, 600, 600 # 5 неверных за 10 минут → блокировка на 10 минут +_guard = threading.Lock() +_fails: dict[str, list[float]] = {} +_locked_until: dict[str, float] = {} + + +def verify_password(user: str, password: str) -> bool: + """Пароль пользователя сессии (единственный пользователь — ADMIN_USER); сравнение за постоянное время.""" + s = get_settings() + return (hmac.compare_digest(user.encode(), s.admin_user.encode()) + and hmac.compare_digest(password.encode(), s.admin_password.encode())) + + +def lockout_remaining(user: str) -> int: + """Сколько секунд осталось до конца блокировки (0 — не заблокирован).""" + with _guard: + return max(0, int(_locked_until.get(user, 0) - time.monotonic() + 0.999)) + + +def register_failure(user: str) -> int: + """Учитывает неверный пароль; возвращает число оставшихся попыток (0 — пользователь заблокирован).""" + now = time.monotonic() + with _guard: + recent = [t for t in _fails.get(user, []) if now - t < WINDOW_S] + [now] + _fails[user] = recent + if len(recent) >= FAIL_LIMIT: + _locked_until[user] = now + LOCK_S + _fails[user] = [] + return 0 + return FAIL_LIMIT - len(recent) + + +def reset_failures(user: str) -> None: + with _guard: + _fails.pop(user, None) + _locked_until.pop(user, None) + + +async def require_api_token(cred: HTTPAuthorizationCredentials | None = Depends(_bearer)) -> None: expected = get_settings().api_token if cred is None or not hmac.compare_digest(cred.credentials.encode(), expected.encode()): raise HTTPException(status_code=401, detail="Invalid or missing API token") + from app.services import events # локальный импорт: security загружается раньше сервисов + + events.set_actor("api") diff --git a/app/services/backups.py b/app/services/backups.py index f352eaf..1a3d623 100644 --- a/app/services/backups.py +++ b/app/services/backups.py @@ -1,38 +1,108 @@ -"""Список бэкапов из бакета с фильтрами (устройство, группа, тип, даты, текст).""" +"""Список бэкапов из бакета с фильтрами и связью с метаданными (backups.id = bkp_…).""" import asyncio +import logging +import threading from datetime import date -from app import s3 -from app.services import devices, groups +from app import ids, s3 +from app.config import get_settings +from app.db import session_scope +from app.models import Backup, Device, now +from app.services import devices, events, groups + +log = logging.getLogger("ros_control.backups") +_sync_lock = threading.Lock() # один процесс: не создаём дубликаты «сирот» при параллельных запросах def _kind(key: str) -> str: return "backup" if key.endswith(".backup") else "rsc" if key.endswith(".rsc") else "other" +def sync_rows(objects: list[dict]) -> dict[str, dict]: + """Сопоставляет файлы бакета со строками backups (по ключу) и возвращает key -> {backup_id, device_id}. + + - файлы без строки (загруженные до появления ID или вручную) получают строку bkp_… по паре <устройство>/<метка>; + - копия, у которой не осталось файлов в бакете, помечается удалённой (строка и ID остаются в истории); + - вернувшиеся файлы снимают пометку.""" + keys = {o["key"]: o for o in objects} + with _sync_lock, session_scope() as s: + by_key: dict[str, dict] = {} + for b in s.query(Backup).all(): + present = [k for k in (b.key_binary, b.key_rsc) if k and k in keys] + for k in present: + by_key[k] = {"backup_id": b.id, "device_id": b.device_id} + if b.status == "done" and b.deleted_at is None and not present: + b.deleted_at = now() + events.record("backup.deleted", "backup", b.id, f"Файлы резервной копии {b.device_name} удалены из бакета", + device_id=b.device_id, data={"key_binary": b.key_binary, "key_rsc": b.key_rsc}, s=s) + elif present and b.deleted_at is not None: + b.deleted_at = None + + orphans: dict[tuple[str, str], dict] = {} + for key, o in keys.items(): + parts = key.split("/") + if key in by_key or len(parts) < 3 or _kind(key) == "other": + continue + stem = parts[-1].rsplit(".", 1)[0] + orphans.setdefault((parts[1], stem), {})[_kind(key)] = (key, o["last_modified"]) + dev_by_name = {d.name: d.id for d in s.query(Device)} + for (name, _stem), files in orphans.items(): + at = max(t for _, t in files.values()) + b = Backup(id=ids.new_id("bkp", at), device_id=dev_by_name.get(name), device_name=name, status="done", + requested_at=at, key_binary=files.get("backup", (None,))[0], key_rsc=files.get("rsc", (None,))[0]) + s.add(b) + s.flush() + events.record("backup.imported", "backup", b.id, f"Файлы бакета связаны с резервной копией {name}", + device_id=b.device_id, data={"key_binary": b.key_binary, "key_rsc": b.key_rsc}, s=s) + for k in (b.key_binary, b.key_rsc): + if k: + by_key[k] = {"backup_id": b.id, "device_id": b.device_id} + return by_key + + +async def reconcile() -> None: + """При старте: привести метаданные в соответствие с бакетом (не критично — сбой не мешает работе).""" + s = get_settings() + if not (s.s3_bucket and s.s3_access_key): + return + try: + objects = await s3.list_backups(None) + events.set_actor("system") + sync_rows(objects) + except Exception: # noqa: BLE001 + log.exception("reconcile бэкапов не выполнен") + + async def search(device: str = "", group: str = "", kind: str = "", date_from: date | None = None, date_to: date | None = None, q: str = "") -> tuple[list[dict], dict]: """Отфильтрованные бэкапы и итоги по всему бакету ({"total": N, "size": байты}). Ключи вида backups/<устройство>/<файл>: устройство берётся из ключа, группа — по имени устройства. - group: "" — все, "none" — без группы (в т.ч. бэкапы удалённых устройств), иначе id группы. - Даты — по времени изменения объекта (UTC), включительно.""" + group: "" — все, "none" — без группы (в т.ч. бэкапы удалённых устройств), иначе ID группы (grp_…). + Даты — по времени изменения объекта (UTC), включительно. Каждому файлу сопоставлен backup_id.""" group_names = {g["id"]: g["name"] for g in groups.list_groups()} device_group = {d.name: group_names.get(d.group_id) for d in devices.list_devices()} - want_group = group_names.get(int(group)) if group.isdigit() else None + want_group = group_names.get(group) if ids.is_id(group, "grp") else None everything = await s3.list_backups(None) + try: + links = sync_rows(everything) + except Exception: # noqa: BLE001 — список важнее связывания + log.exception("sync_rows не выполнен") + links = {} stats = {"total": len(everything), "size": sum(i["size"] for i in everything)} q = q.strip().lower() out = [] for item in everything: parts = item["key"].split("/") name = parts[1] if len(parts) >= 3 else "" - item = {**item, "device": name, "group": device_group.get(name), "kind": _kind(item["key"])} + link = links.get(item["key"], {}) + item = {**item, "device": name, "group": device_group.get(name), "kind": _kind(item["key"]), + "backup_id": link.get("backup_id"), "device_id": link.get("device_id")} if device and name != device: continue if group == "none" and item["group"] is not None: continue - if group.isdigit() and (want_group is None or item["group"] != want_group): + if ids.is_id(group, "grp") and (want_group is None or item["group"] != want_group): continue if kind and item["kind"] != kind: continue @@ -74,4 +144,8 @@ async def delete_many(keys: list[str]) -> tuple[int, int]: return False results = await asyncio.gather(*(one(k) for k in keys)) + try: # метаданные: копии, у которых не осталось файлов, помечаются удалёнными (с событием) + sync_rows(await s3.list_backups(None)) + except Exception: # noqa: BLE001 + log.exception("sync_rows после удаления не выполнен") return sum(results), len(results) - sum(results) diff --git a/app/services/devices.py b/app/services/devices.py index d76d4c4..caa5dfe 100644 --- a/app/services/devices.py +++ b/app/services/devices.py @@ -1,19 +1,20 @@ import re from dataclasses import dataclass -from app import security +from app import ids, security from app.config import get_settings from app.db import session_scope from app.models import Device, Group -from app.ros.operations import version_newer from app.ros.client import RosClient +from app.ros.operations import version_newer +from app.services import events NAME_RE = re.compile(r"^[A-Za-z0-9._-]{1,64}$") # имя попадает в ключи S3 — только безопасные символы @dataclass class Conn: - id: int + id: str name: str host: str port: int @@ -33,7 +34,8 @@ def list_devices() -> list[Device]: return list(s.query(Device).order_by(Device.name)) -def get_device(device_id: int) -> Device: +def get_device(device_id: str) -> Device: + ids.check(device_id, "dev") with session_scope() as s: d = s.get(Device, device_id) if d is None: @@ -41,7 +43,9 @@ def get_device(device_id: int) -> Device: return d -def _check_group(s, group_id: int | None) -> None: +def _check_group(s, group_id: str | None) -> None: + if group_id is not None: + ids.check(group_id, "grp") if group_id is not None and s.get(Group, group_id) is None: raise LookupError(f"Группа {group_id} не найдена") @@ -56,7 +60,7 @@ def _clean_note(note: str | None) -> str | None: return note or None -def _resolve_group(s, group_id: int | None, new_group: str | None) -> int | None: +def _resolve_group(s, group_id: str | None, new_group: str | None) -> str | None: """Группа устройства. new_group создаёт группу (или берёт существующую с таким именем) в той же транзакции, что и устройство: при ошибке не остаётся лишней группы.""" name = (new_group or "").strip() @@ -68,13 +72,14 @@ def _resolve_group(s, group_id: int | None, new_group: str | None) -> int | None g = Group(name=name) s.add(g) s.flush() + events.record("group.created", "group", g.id, f"Группа «{name}» создана", data={"name": name}, s=s) return g.id _check_group(s, group_id) return group_id def create_device(name, host, port, username, password, verify_tls=False, use_tls=True, - group_id: int | None = None, new_group: str | None = None, note: str | None = None) -> Device: + group_id: str | None = None, new_group: str | None = None, note: str | None = None) -> Device: validate_name(name) note = _clean_note(note) with session_scope() as s: @@ -86,37 +91,55 @@ def create_device(name, host, port, username, password, verify_tls=False, use_tl group_id=group_id, note=note) s.add(d) s.flush() + events.record("device.created", "device", d.id, f"Устройство {name} добавлено", device_id=d.id, + data={"name": name, "host": host, "port": port, "group_id": group_id}, s=s) return d -def update_device(device_id: int, **fields) -> Device: +def update_device(device_id: str, **fields) -> Device: + ids.check(device_id, "dev") with session_scope() as s: d = s.get(Device, device_id) if d is None: raise LookupError(f"Устройство {device_id} не найдено") if fields.get("name") and fields["name"] != d.name: # имя — часть ключей бэкапов в S3 raise ValueError("Имя устройства нельзя изменить") + changed = [] for k in ("host", "port", "username", "verify_tls", "use_tls"): - if fields.get(k) is not None: + if fields.get(k) is not None and getattr(d, k) != fields[k]: setattr(d, k, fields[k]) + changed.append(k) if "group_id" in fields or fields.get("new_group"): # group_id=None — открепить от группы - d.group_id = _resolve_group(s, fields.get("group_id"), fields.get("new_group")) + gid = _resolve_group(s, fields.get("group_id"), fields.get("new_group")) + if gid != d.group_id: + d.group_id = gid + changed.append("group_id") if "note" in fields: # пустое примечание — очистить - d.note = _clean_note(fields["note"]) + note = _clean_note(fields["note"]) + if note != d.note: + d.note = note + changed.append("note") if fields.get("password"): # пустой пароль = не менять d.password_enc = security.encrypt(fields["password"]) + changed.append("password") # только факт смены, значение в журнал не пишется + if changed: + events.record("device.updated", "device", d.id, f"Устройство {d.name} изменено: {', '.join(changed)}", + device_id=d.id, data={"fields": changed}, s=s) return d -def delete_device(device_id: int) -> None: +def delete_device(device_id: str) -> None: + ids.check(device_id, "dev") with session_scope() as s: d = s.get(Device, device_id) if d is None: raise LookupError(f"Устройство {device_id} не найдено") + events.record("device.deleted", "device", d.id, f"Устройство {d.name} удалено", device_id=d.id, + data={"name": d.name, "host": d.host}, s=s) s.delete(d) -def get_conn(device_id: int) -> Conn: +def get_conn(device_id: str) -> Conn: d = get_device(device_id) return Conn(d.id, d.name, d.host, d.port, d.username, security.decrypt(d.password_enc), d.verify_tls, d.use_tls) diff --git a/app/services/events.py b/app/services/events.py new file mode 100644 index 0000000..088698f --- /dev/null +++ b/app/services/events.py @@ -0,0 +1,143 @@ +"""Журнал событий: каждая запись имеет свой 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 diff --git a/app/services/groups.py b/app/services/groups.py index fe9f8dc..db5eddc 100644 --- a/app/services/groups.py +++ b/app/services/groups.py @@ -1,7 +1,9 @@ from sqlalchemy import func, update +from app import ids from app.db import session_scope from app.models import Device, Group +from app.services import events def _clean(name: str) -> str: @@ -19,7 +21,8 @@ def list_groups() -> list[dict]: for g in s.query(Group).order_by(Group.name)] -def get_group(group_id: int) -> Group: +def get_group(group_id: str) -> Group: + ids.check(group_id, "grp") with session_scope() as s: g = s.get(Group, group_id) if g is None: @@ -35,10 +38,12 @@ def create_group(name: str) -> Group: g = Group(name=name) s.add(g) s.flush() + events.record("group.created", "group", g.id, f"Группа «{name}» создана", data={"name": name}, s=s) return g -def rename_group(group_id: int, name: str) -> Group: +def rename_group(group_id: str, name: str) -> Group: + ids.check(group_id, "grp") name = _clean(name) with session_scope() as s: g = s.get(Group, group_id) @@ -46,23 +51,35 @@ def rename_group(group_id: int, name: str) -> Group: raise LookupError(f"Группа {group_id} не найдена") if s.query(Group).filter(Group.name == name, Group.id != group_id).first(): raise ValueError(f"Группа «{name}» уже существует") - g.name = name + old, g.name = g.name, name + if old != name: + events.record("group.renamed", "group", g.id, f"Группа «{old}» переименована в «{name}»", + data={"old": old, "new": name}, s=s) return g -def delete_group(group_id: int) -> None: +def delete_group(group_id: str) -> None: """Удаляет группу; устройства остаются и становятся «Без группы» (FK в SQLite не включены).""" + ids.check(group_id, "grp") with session_scope() as s: g = s.get(Group, group_id) if g is None: raise LookupError(f"Группа {group_id} не найдена") + detached = [d for (d,) in s.query(Device.id).filter(Device.group_id == group_id)] s.execute(update(Device).where(Device.group_id == group_id).values(group_id=None)) + events.record("group.deleted", "group", g.id, f"Группа «{g.name}» удалена", data={"name": g.name, "detached_devices": detached}, s=s) s.delete(g) -def move_devices(device_ids: list[int], group_id: int | None) -> None: +def move_devices(device_ids: list[str], group_id: str | None) -> None: """Переносит устройства в группу (None — «Без группы»).""" + for i in device_ids: + ids.check(i, "dev") with session_scope() as s: - if group_id is not None and s.get(Group, group_id) is None: - raise LookupError(f"Группа {group_id} не найдена") + if group_id is not None: + ids.check(group_id, "grp") + if s.get(Group, group_id) is None: + raise LookupError(f"Группа {group_id} не найдена") s.execute(update(Device).where(Device.id.in_(device_ids)).values(group_id=group_id)) + events.record("devices.moved", "group", group_id, f"Перенесено устройств: {len(device_ids)}", + data={"device_ids": device_ids, "group_id": group_id}, s=s) diff --git a/app/services/jobs.py b/app/services/jobs.py index e34de02..2bf0fc6 100644 --- a/app/services/jobs.py +++ b/app/services/jobs.py @@ -1,10 +1,11 @@ """Фоновые задачи: запись в таблицу jobs, параллелизм ограничен MAX_CONCURRENCY.""" import asyncio +from app import ids from app.config import get_settings from app.db import session_scope from app.models import Device, Job, now -from app.services import ops +from app.services import events, ops JOB_TYPES = { "backup": ops.run_backup, @@ -23,15 +24,19 @@ def _semaphore() -> asyncio.Semaphore: return _sem -def _finish(job_id: int, status: str, message: str) -> None: +def _finish(job_id: str, status: str, message: str) -> None: with session_scope() as s: j = s.get(Job, job_id) j.status, j.message = status, message if status in ("done", "failed"): j.finished_at = now() + etype = {"running": "job.started", "done": "job.done", "failed": "job.failed"}[status] + events.record(etype, "job", job_id, message or f"Задача {j.type}: {status}", device_id=j.device_id, + job_id=job_id, data={"type": j.type, "status": status}, s=s) -async def _run(job_id: int, job_type: str, device_id: int) -> None: +async def _run(job_id: str, job_type: str, device_id: str) -> None: + events.set_job(job_id) # события и бэкап, созданные внутри задачи, ссылаются на её ID async with _semaphore(): _finish(job_id, "running", "") try: @@ -40,11 +45,13 @@ async def _run(job_id: int, job_type: str, device_id: int) -> None: _finish(job_id, "failed", str(e)) -def start_jobs(job_type: str, device_ids: list[int]) -> list[int]: - """Создаёт задачи для списка устройств и запускает их в фоне. Возвращает id задач.""" +def start_jobs(job_type: str, device_ids: list[str]) -> list[str]: + """Создаёт задачи для списка устройств и запускает их в фоне. Возвращает ID задач.""" if job_type not in JOB_TYPES: raise ValueError(f"Неизвестный тип задачи: {job_type}") - job_ids = [] + for did in device_ids: + ids.check(did, "dev") + started = [] with session_scope() as s: for did in device_ids: d = s.get(Device, did) @@ -53,12 +60,14 @@ def start_jobs(job_type: str, device_ids: list[int]) -> list[int]: j = Job(device_id=did, device_name=d.name, type=job_type) s.add(j) s.flush() - job_ids.append((j.id, did)) - for jid, did in job_ids: - task = asyncio.create_task(_run(jid, job_type, did)) + events.record("job.created", "job", j.id, f"Задача {job_type} для {d.name} создана", device_id=did, + job_id=j.id, data={"type": job_type}, s=s) + started.append((j.id, did)) + for jid, did in started: + task = asyncio.create_task(_run(jid, job_type, did)) # наследует контекст (актор) _tasks.add(task) task.add_done_callback(_tasks.discard) - return [jid for jid, _ in job_ids] + return [jid for jid, _ in started] def list_jobs(limit: int = 50) -> list[Job]: @@ -66,7 +75,8 @@ def list_jobs(limit: int = 50) -> list[Job]: return list(s.query(Job).order_by(Job.id.desc()).limit(limit)) -def get_job(job_id: int) -> Job: +def get_job(job_id: str) -> Job: + ids.check(job_id, "job") with session_scope() as s: j = s.get(Job, job_id) if j is None: @@ -79,3 +89,5 @@ def fail_stale_jobs() -> None: with session_scope() as s: for j in s.query(Job).filter(Job.status.in_(("pending", "running"))): j.status, j.message, j.finished_at = "failed", "Прервано перезапуском сервера", now() + events.record("job.failed", "job", j.id, j.message, device_id=j.device_id, job_id=j.id, + data={"type": j.type, "status": "failed"}, s=s) diff --git a/app/services/ops.py b/app/services/ops.py index d90faac..f174742 100644 --- a/app/services/ops.py +++ b/app/services/ops.py @@ -11,19 +11,25 @@ from app.db import session_scope from app.models import Backup, Device, now from app.ros import operations as ros from app.ros.client import RosError -from app.services import devices +from app.services import devices, events -def _save_status(device_id: int, fields: dict) -> None: +def _save_status(device_id: str, fields: dict) -> None: with session_scope() as s: d = s.get(Device, device_id) if d is not None: # устройство могли удалить во время опроса + was = d.online for k, v in fields.items(): setattr(d, k, v) d.status_at = now() + if "online" in fields and fields["online"] is not was: # в журнал — только смены состояния, не каждый опрос + up = fields["online"] + events.record("device.online" if up else "device.offline", "device", d.id, + f"Устройство {d.name} {'доступно' if up else 'недоступно'}", device_id=d.id, + data={"error": (fields.get("last_error") or "")[:300]} if not up else None, s=s) -async def refresh_status(device_id: int) -> None: +async def refresh_status(device_id: str) -> None: """Полный опрос: статус, версии и проверка обновлений. Недоступность — не исключение.""" conn = devices.get_conn(device_id) try: @@ -35,7 +41,7 @@ async def refresh_status(device_id: int) -> None: _save_status(device_id, fields) -async def poll_device(device_id: int, full: bool = False) -> None: +async def poll_device(device_id: str, full: bool = False) -> None: """Фоновый опрос. Лёгкий (один запрос system/resource) обновляет online, uptime и версию ROS; полный — если нужна проверка обновлений (full) либо устройство было недоступно/ещё не опрошено (после возвращения версии могли измениться).""" @@ -51,18 +57,18 @@ async def poll_device(device_id: int, full: bool = False) -> None: _save_status(device_id, dict(online=True, status_json=json.dumps(status), last_error=None)) -async def refresh_many(device_ids: list[int] | None = None) -> None: - ids = device_ids if device_ids is not None else [d.id for d in devices.list_devices()] +async def refresh_many(device_ids: list[str] | None = None) -> None: + targets = device_ids if device_ids is not None else [d.id for d in devices.list_devices()] sem = asyncio.Semaphore(get_settings().max_concurrency) - async def one(i: int): + async def one(i: str): async with sem: await refresh_status(i) - await asyncio.gather(*(one(i) for i in ids)) + await asyncio.gather(*(one(i) for i in targets)) -async def run_backup(device_id: int) -> str: +async def run_backup(device_id: str) -> str: """Бэкап .backup + .rsc: создать на устройстве → скачать через REST во временную папку → загрузить в S3 от имени сервера → удалить файлы с устройства.""" conn = devices.get_conn(device_id) @@ -75,11 +81,15 @@ async def run_backup(device_id: int) -> str: key_bin, key_rsc = files[f"{base}.backup"], files[f"{base}.rsc"] with session_scope() as s: - b = Backup(device_id=device_id, key_binary=key_bin, key_rsc=key_rsc) + b = Backup(device_id=device_id, device_name=conn.name, job_id=events.current_job(), + key_binary=key_bin, key_rsc=key_rsc) s.add(b) s.get(Device, device_id).last_backup_requested_at = b.requested_at = now() s.flush() backup_id = b.id + events.record("backup.created", "backup", backup_id, f"Резервная копия {conn.name} запрошена", device_id=device_id, + data={"key_binary": key_bin, "key_rsc": key_rsc}, s=s) + meta = {"backup-id": backup_id, "device-id": device_id} # связь файлов в бакете с метаданными status, error = "failed", "прервано" try: @@ -90,7 +100,7 @@ async def run_backup(device_id: int) -> str: for name, key in files.items(): local = Path(tmp) / name await ros.download_file(c, name, local) - await s3.upload_file(local, key) + await s3.upload_file(local, key, meta) finally: for name in files: # на устройстве файлы не остаются ни при успехе, ни при сбое try: @@ -105,20 +115,23 @@ async def run_backup(device_id: int) -> str: with session_scope() as s: b = s.get(Backup, backup_id) b.status, b.error = status, error + events.record("backup.done" if status == "done" else "backup.failed", "backup", backup_id, + "Резервная копия загружена в S3" if status == "done" else f"Резервная копия не создана: {error}", + device_id=device_id, data={"key_binary": key_bin, "key_rsc": key_rsc} if status == "done" else {"error": error}, s=s) return f"Бэкап загружен в S3: {key_bin}, {key_rsc}" -async def set_channel(device_id: int, channel: str) -> None: +async def set_channel(device_id: str, channel: str) -> None: async with devices.open_client(devices.get_conn(device_id)) as c: await ros.set_channel(c, channel) await refresh_status(device_id) -async def run_ros_update(device_id: int) -> str: +async def run_ros_update(device_id: str) -> str: async with devices.open_client(devices.get_conn(device_id)) as c: return await ros.install_ros_update(c) -async def run_fw_update(device_id: int) -> str: +async def run_fw_update(device_id: str) -> str: async with devices.open_client(devices.get_conn(device_id)) as c: return await ros.upgrade_firmware(c) diff --git a/app/services/poller.py b/app/services/poller.py index a4296cb..2c61a0c 100644 --- a/app/services/poller.py +++ b/app/services/poller.py @@ -4,12 +4,12 @@ import logging import time from app.config import get_settings -from app.services import devices, ops +from app.services import devices, events, ops log = logging.getLogger("ros_control.poller") -async def cycle(last_full: dict[int, float]) -> None: +async def cycle(last_full: dict[str, float]) -> None: """Один проход по всем устройствам (параллельно, не более MAX_CONCURRENCY одновременно).""" s = get_settings() ids = [d.id for d in devices.list_devices()] @@ -18,7 +18,7 @@ async def cycle(last_full: dict[int, float]) -> None: sem = asyncio.Semaphore(s.max_concurrency) started = time.monotonic() - async def one(device_id: int) -> None: + async def one(device_id: str) -> None: async with sem: full = started - last_full.get(device_id, float("-inf")) >= s.update_check_interval try: @@ -35,7 +35,8 @@ async def cycle(last_full: dict[int, float]) -> None: async def run() -> None: - last_full: dict[int, float] = {} + events.set_actor("poller") + last_full: dict[str, float] = {} interval = get_settings().poll_interval while True: try: diff --git a/app/services/rotation.py b/app/services/rotation.py new file mode 100644 index 0000000..c174f0b --- /dev/null +++ b/app/services/rotation.py @@ -0,0 +1,20 @@ +"""Ротация журнала: при старте и далее раз в час (настройки — services.settings).""" +import asyncio +import logging + +from app.services import events + +log = logging.getLogger("ros_control.rotation") +INTERVAL_S = 3600 + + +async def run() -> None: + events.set_actor("system") + while True: + try: + done = await asyncio.to_thread(events.rotate) + if done["by_age"] or done["by_count"]: + log.warning("ротация журнала: %s", done) + except Exception: # noqa: BLE001 — сбой ротации не должен останавливать приложение + log.exception("ротация журнала не выполнена") + await asyncio.sleep(INTERVAL_S) diff --git a/app/services/settings.py b/app/services/settings.py new file mode 100644 index 0000000..9600127 --- /dev/null +++ b/app/services/settings.py @@ -0,0 +1,53 @@ +"""Настройки, редактируемые из UI и хранимые в БД (app_settings); значения по умолчанию — из .env.""" +import json + +from app.config import get_settings +from app.db import session_scope +from app.models import AppSetting, now +from app.services import events + +RETENTION_MAX_DAYS = 3650 +ROWS_MIN, ROWS_MAX = 100, 10_000_000 + + +def _get(key: str, default): + with session_scope() as s: + row = s.get(AppSetting, key) + return json.loads(row.value) if row is not None else default + + +def journal() -> dict: + """Действующие настройки ротации: {'retention_days': int, 'max_rows': int}; 0 — без ограничения.""" + cfg = get_settings() + return {"retention_days": _get("events.retention_days", cfg.events_retention_days), + "max_rows": _get("events.max_rows", cfg.events_max_rows)} + + +def validate_journal(retention_days, max_rows) -> dict: + try: + days, rows = int(str(retention_days).strip()), int(str(max_rows).strip()) + except ValueError: + raise ValueError("Укажите целые числа") from None + if not 0 <= days <= RETENTION_MAX_DAYS: + raise ValueError(f"Срок хранения: целое число от 0 до {RETENTION_MAX_DAYS} (0 — без ограничения)") + if rows != 0 and not ROWS_MIN <= rows <= ROWS_MAX: + raise ValueError(f"Максимум записей: 0 (без ограничения) или от {ROWS_MIN} до {ROWS_MAX:,}".replace(",", " ")) + return {"retention_days": days, "max_rows": rows} + + +def save_journal(retention_days, max_rows) -> dict: + """Сохраняет настройки ротации, пишет событие и сразу применяет ротацию.""" + new = validate_journal(retention_days, max_rows) + old = journal() + with session_scope() as s: + for key, value in (("events.retention_days", new["retention_days"]), ("events.max_rows", new["max_rows"])): + row = s.get(AppSetting, key) + if row is None: + s.add(AppSetting(key=key, value=json.dumps(value))) + else: + row.value, row.updated_at = json.dumps(value), now() + if new != old: + events.record("journal.settings_changed", "journal", None, "Изменены настройки ротации журнала", + data={"old": old, "new": new}, s=s) + events.rotate() + return new diff --git a/app/ui/icons.py b/app/ui/icons.py index 5840014..2bf4200 100644 --- a/app/ui/icons.py +++ b/app/ui/icons.py @@ -12,5 +12,7 @@ ICONS = { "download": '', "router": '', "trash": '', + "settings": '', + "copy": '', "arrow": '', } diff --git a/app/ui/routes.py b/app/ui/routes.py index 8f55ae6..dd490a3 100644 --- a/app/ui/routes.py +++ b/app/ui/routes.py @@ -1,4 +1,5 @@ """WEB UI: Jinja2 + HTMX. Использует те же сервисы, что и JSON API.""" +import json from datetime import date from pathlib import Path from urllib.parse import urlencode @@ -7,10 +8,10 @@ from fastapi import APIRouter, Depends, Form, Request, Response from fastapi.responses import HTMLResponse, RedirectResponse from fastapi.templating import Jinja2Templates -from app import s3, security +from app import ids, s3, security from app.config import get_settings from app.ros.operations import CHANNELS, version_newer -from app.services import backups, devices, groups, jobs, ops +from app.services import backups, devices, events, groups, jobs, ops, settings from app.ui.icons import ICONS @@ -28,6 +29,20 @@ templates = Jinja2Templates(directory=Path(__file__).parent / "templates") templates.env.filters["dt"] = lambda v: v.strftime("%Y-%m-%d %H:%M") if v else "—" templates.env.filters["size"] = lambda n: f"{n / 1024 / 1024:.1f} МБ" if n >= 1024 * 1024 else f"{n / 1024:.0f} КБ" templates.env.filters["plural"] = plural +templates.env.filters["short_id"] = ids.short_id +templates.env.filters["num"] = lambda n: f"{int(n):,}".replace(",", "\u00a0") +templates.env.filters["dts"] = lambda v: v.strftime("%Y-%m-%d %H:%M:%S") if v else "—" + +# смысловая окраска типов событий (как в утверждённом макете) +_EV_KIND = {"bad": ("device.offline", "job.failed", "auth.failed", "backup.failed", "journal.clear_denied", "journal.clear_locked"), + "ok": ("device.online", "job.done", "backup.done", "auth.login"), + "run": ("job.started",), + "warn": ("backup.deleted", "device.deleted", "group.deleted", "system.migrated", "journal.cleared", "journal.rotated")} +_EV_KIND_BY_TYPE = {t: k for k, ts in _EV_KIND.items() for t in ts} +_EV_ENTITY = {"device": "Устройство", "group": "Группа", "job": "Задача", "backup": "Копия", "session": "Сессия", + "system": "Система", "journal": "Журнал"} +templates.env.globals["ev_kind"] = lambda t: _EV_KIND_BY_TYPE.get(t, "plain") +templates.env.globals["ev_entity"] = lambda t: _EV_ENTITY.get(t, t) templates.env.globals["ICONS"] = ICONS templates.env.globals["ros_newer"] = version_newer @@ -36,9 +51,11 @@ class LoginRequired(Exception): pass -def require_login(request: Request) -> None: - if not request.session.get("user"): +async def require_login(request: Request) -> None: + user = request.session.get("user") + if not user: raise LoginRequired() + events.set_actor(f"ui:{user}") # async: ContextVar остаётся в контексте обработчика запроса router = APIRouter(include_in_schema=False) @@ -59,8 +76,10 @@ def _done(request: Request, fallback: str = "/"): return RedirectResponse(fallback, status_code=303) -def _opt_int(value: str) -> int | None: - return int(value) if str(value).strip().isdigit() else None +def _opt_id(value: str, prefix: str) -> str | None: + """ID нужного типа или None («без группы», пустое значение, чужой тип, мусор).""" + value = str(value).strip() + return value if ids.is_id(value, prefix) else None FILTER_FIELDS = ("group", "q", "status", "updates", "channel") # в форме и URL — с префиксом f_ @@ -73,6 +92,8 @@ def _flt(params) -> dict: for k, allowed in FILTER_VALUES.items(): if flt[k] not in allowed: flt[k] = "" + if flt["group"] not in ("", "none") and not ids.is_id(flt["group"], "grp"): + flt["group"] = "" return flt @@ -117,8 +138,12 @@ async def login_form(request: Request): @router.post("/login") async def login(request: Request, username: str = Form(), password: str = Form()): if not security.check_admin(username, password): + events.set_actor("anonymous") + events.record("auth.failed", "session", None, "Неудачная попытка входа в UI", data={"user": username[:64]}) return _render(request, "login.html", error="Неверный логин или пароль") request.session["user"] = username + events.set_actor(f"ui:{username}") + events.record("auth.login", "session", None, f"Вход в UI: {username}", data={"user": username}) return RedirectResponse("/", status_code=303) @@ -153,7 +178,7 @@ async def jobs_table(request: Request): @router.post("/ui/refresh", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def refresh(request: Request, device_ids: list[int] = Form(default=[])): +async def refresh(request: Request, device_ids: list[str] = Form(default=[])): flt = _flt(await request.form()) # без выбора обновляем все устройства текущего представления (группа + фильтры) await ops.refresh_many(device_ids or [d.id for d in _visible(flt)]) @@ -161,7 +186,7 @@ async def refresh(request: Request, device_ids: list[int] = Form(default=[])): @router.post("/ui/batch/{action}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def batch(request: Request, action: str, device_ids: list[int] = Form(default=[]), +async def batch(request: Request, action: str, device_ids: list[str] = Form(default=[]), channel: str = Form(default="")): if not device_ids: return _render(request, "_jobs.html", jobs=jobs.list_jobs(15), flash="Выберите устройства") @@ -175,7 +200,7 @@ async def batch(request: Request, action: str, device_ids: list[int] = Form(defa @router.post("/ui/devices/{device_id}/{action}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def device_action(request: Request, device_id: int, action: str, channel: str = Form(default="")): +async def device_action(request: Request, device_id: str, action: str, channel: str = Form(default="")): """Действие над одним устройством (пункты меню «⋯» в строке таблицы).""" if action == "refresh": await ops.refresh_status(device_id) @@ -188,10 +213,10 @@ async def device_action(request: Request, device_id: int, action: str, channel: @router.post("/ui/move", dependencies=[Depends(require_login)]) -async def move(device_ids: list[int] = Form(default=[]), group_id: str = Form(""), next: str = Form("/")): +async def move(device_ids: list[str] = Form(default=[]), group_id: str = Form(""), next: str = Form("/")): """Перенос отмеченных устройств в группу (обычный POST + редирект: счётчики вкладок обновляются).""" if device_ids: - groups.move_devices(device_ids, _opt_int(group_id)) + groups.move_devices(device_ids, _opt_id(group_id, "grp")) return RedirectResponse(next if next.startswith("/") and not next.startswith("//") else "/", status_code=303) @@ -223,12 +248,12 @@ def _group_form(request: Request, group, name: str, error: str | None = None): @router.get("/ui/dialog/device", response_class=HTMLResponse, dependencies=[Depends(require_login)]) async def dialog_device_new(request: Request, group: str = ""): v = _device_values() - v["group_id"] = group if group.isdigit() else "" + v["group_id"] = group if ids.is_id(group, "grp") else "" return _render(request, "_device_form.html", in_modal=True, device=None, v=v, error=None, groups=groups.list_groups()) @router.get("/ui/dialog/device/{device_id}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_device_edit(request: Request, device_id: int): +async def dialog_device_edit(request: Request, device_id: str): d = devices.get_device(device_id) return _render(request, "_device_form.html", in_modal=True, device=d, v=_device_values(d), error=None, groups=groups.list_groups()) @@ -240,7 +265,7 @@ async def dialog_group_new(request: Request): @router.get("/ui/dialog/group/{group_id}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_group_edit(request: Request, group_id: int): +async def dialog_group_edit(request: Request, group_id: str): g = groups.get_group(group_id) return _render(request, "_group_form.html", in_modal=True, group=g, gname=g.name, error=None) @@ -270,30 +295,30 @@ async def group_create(request: Request, name: str = Form()): @router.post("/groups/{group_id}/rename", dependencies=[Depends(require_login)]) -async def group_rename(request: Request, group_id: int, name: str = Form()): +async def group_rename(request: Request, group_id: str, name: str = Form()): return _group_action(request, groups.rename_group, group_id, name, name=name, group=groups.get_group(group_id)) @router.post("/groups/{group_id}/delete", dependencies=[Depends(require_login)]) -async def group_delete(request: Request, group_id: int): +async def group_delete(request: Request, group_id: str): return _group_action(request, groups.delete_group, group_id) # --- управление списком устройств --- -def _group_choice(group_id: str, new_group: str) -> tuple[int | None, str | None]: +def _group_choice(group_id: str, new_group: str) -> tuple[str | None, str | None]: """Выбор в списке «Группа»: id, «без группы» или «Новая группа…» (тогда нужно название).""" if group_id == "__new__": if not new_group.strip(): raise ValueError("Введите название новой группы") return None, new_group.strip() - return _opt_int(group_id), None + return _opt_id(group_id, "grp"), None @router.get("/devices/new", response_class=HTMLResponse, dependencies=[Depends(require_login)]) async def device_new(request: Request, group: str = ""): v = _device_values() - v["group_id"] = group if group.isdigit() else "" + v["group_id"] = group if ids.is_id(group, "grp") else "" return _device_form(request, None, v) @@ -314,13 +339,13 @@ async def device_create(request: Request, name: str = Form(), host: str = Form() @router.get("/devices/{device_id}/edit", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def device_edit(request: Request, device_id: int): +async def device_edit(request: Request, device_id: str): d = devices.get_device(device_id) return _device_form(request, d, _device_values(d)) @router.post("/devices/{device_id}/edit", dependencies=[Depends(require_login)]) -async def device_update(request: Request, device_id: int, host: str = Form(), port: int = Form(443), +async def device_update(request: Request, device_id: str, host: str = Form(), port: int = Form(443), username: str = Form(), password: str = Form(""), verify_tls: bool = Form(False), use_tls: bool = Form(False), group_id: str = Form(""), new_group: str = Form(""), note: str = Form("")): @@ -338,7 +363,7 @@ async def device_update(request: Request, device_id: int, host: str = Form(), po @router.post("/devices/{device_id}/delete", dependencies=[Depends(require_login)]) -async def device_delete(request: Request, device_id: int): +async def device_delete(request: Request, device_id: str): devices.delete_device(device_id) return _done(request) @@ -390,3 +415,137 @@ async def backup_delete(key: str = Form()): raise ValueError("Недопустимый ключ") await s3.delete_object(key) return RedirectResponse("/backups", status_code=303) + + +# --- журнал событий --- + +EVENTS_PAGE = 100 +EV_ACTORS = ("ui", "api", "poller", "system", "anonymous") + + +def _ev_filters(params) -> dict: + p = lambda k: str(params.get(k) or "").strip() # noqa: E731 + f = dict(type=p("type"), actor=p("actor"), device=p("device"), date_from=p("date_from"), date_to=p("date_to"), q=p("q")) + if f["device"] and not ids.is_id(f["device"], "dev"): + f["device"] = "" + if f["actor"] not in EV_ACTORS: + f["actor"] = "" + return f + + +def _ev_query(f: dict) -> dict: + return dict(type_=f["type"] or None, actor=f["actor"] or None, device_id=f["device"] or None, + date_from=_opt_date(f["date_from"]), date_to=_opt_date(f["date_to"]), text=f["q"] or None) + + +def _ev_page(f: dict, before: str | None = None, shown: int = 0) -> dict: + query = _ev_query(f) + rows = events.list_events(before=before, limit=EVENTS_PAGE + 1, **query) + more = len(rows) > EVENTS_PAGE + rows = rows[:EVENTS_PAGE] + return dict(rows=rows, more=more, next_before=rows[-1].id if rows else None, total=events.count(**query), + shown=shown + len(rows), qs=urlencode({k: v for k, v in f.items() if v})) + + +def _rotation_text(cfg: dict) -> str: + days, rows = cfg["retention_days"], cfg["max_rows"] + if not days and not rows: + return "без ограничений" + d = f"{days} {plural(days, 'день', 'дня', 'дней')}" if days else "без ограничения по сроку" + r = f"{rows:,}".replace(",", "\u00a0") + " записей" if rows else "без ограничения по числу" + return f"{d} / {r}" + + +@router.get("/events", response_class=HTMLResponse, dependencies=[Depends(require_login)]) +async def events_page(request: Request): + f = _ev_filters(request.query_params) + st = events.stats() + all_types = events.types() + return _render(request, "events.html", section="events", flt=f, flt_active=sum(1 for v in f.values() if v), + st=st, cfg=settings.journal(), rotation=_rotation_text(settings.journal()), oob=False, + categories=sorted({t.split(".")[0] for t in all_types}), types=all_types, actors=EV_ACTORS, + devices=devices.list_devices(), **_ev_page(f)) + + +@router.get("/ui/events", response_class=HTMLResponse, dependencies=[Depends(require_login)]) +async def events_more(request: Request, before: str = "", shown: int = 0): + """Следующие записи журнала (кнопка «Показать ещё»): строки + обновлённая кнопка и подвал (OOB).""" + f = _ev_filters(request.query_params) + return _render(request, "_events_chunk.html", oob=True, **_ev_page(f, before=before or None, shown=shown)) + + +def _event_ctx(event_id: str) -> dict: + ids.check(event_id, "evt") + e = events.get_event(event_id) + device = None + if e.device_id: + try: + device = devices.get_device(e.device_id) + except LookupError: + device = None + data = json.dumps(json.loads(e.data), ensure_ascii=False, indent=2) if e.data else "" + return dict(e=e, device=device, data_json=data) + + +@router.get("/ui/dialog/event/{event_id}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) +async def dialog_event(request: Request, event_id: str): + return _render(request, "_event_dialog.html", in_modal=True, **_event_ctx(event_id)) + + +@router.get("/events/{event_id}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) +async def event_page(request: Request, event_id: str): + return _render(request, "event_page.html", section="events", in_modal=False, **_event_ctx(event_id)) + + +def _settings_dialog(request: Request, error: str | None = None, values: dict | None = None): + return _render(request, "_events_settings.html", in_modal=_is_htmx(request), error=error, + values=values or settings.journal(), st=events.stats()) + + +@router.get("/ui/dialog/events-settings", response_class=HTMLResponse, dependencies=[Depends(require_login)]) +async def dialog_events_settings(request: Request): + return _settings_dialog(request) + + +@router.post("/events/settings", dependencies=[Depends(require_login)]) +async def events_settings_save(request: Request, retention_days: str = Form(""), max_rows: str = Form("")): + try: + settings.save_journal(retention_days, max_rows) + except ValueError as e: + return _settings_dialog(request, str(e), {"retention_days": retention_days, "max_rows": max_rows}) + return _done(request, "/events") + + +def _clear_dialog(request: Request, error: str | None = None): + user = request.session.get("user", "") + lock = security.lockout_remaining(user) + return _render(request, "_events_clear.html", in_modal=_is_htmx(request), error=error, user=user, + locked=lock > 0, minutes=max(1, -(-lock // 60)), st=events.stats()) + + +@router.get("/ui/dialog/events-clear", response_class=HTMLResponse, dependencies=[Depends(require_login)]) +async def dialog_events_clear(request: Request): + return _clear_dialog(request) + + +@router.post("/events/clear", dependencies=[Depends(require_login)]) +async def events_clear(request: Request, password: str = Form("")): + """Очистка журнала — только через окно подтверждения паролем пользователя сессии (через API её нет).""" + if not _is_htmx(request): + raise ValueError("Очистка журнала выполняется только через окно подтверждения") + user = request.session["user"] + if security.lockout_remaining(user): + return _clear_dialog(request) + if not password: + return _clear_dialog(request, "Введите пароль") + if not security.verify_password(user, password): + left = security.register_failure(user) + events.record("journal.clear_denied", "journal", None, "Неверный пароль при очистке журнала", data={"attempts_left": left}) + if left == 0: + events.record("journal.clear_locked", "journal", None, "Очистка журнала заблокирована после неверных попыток", + data={"lock_seconds": security.LOCK_S}) + return _clear_dialog(request) + return _clear_dialog(request, f"Неверный пароль. Осталось попыток: {left}") + security.reset_failures(user) + events.clear(user) + return _done(request, "/events") diff --git a/app/ui/static/app.js b/app/ui/static/app.js index d7b1ad5..a443400 100644 --- a/app/ui/static/app.js +++ b/app/ui/static/app.js @@ -85,7 +85,7 @@ $("#modal-body [autofocus], #modal-body input:not([type=hidden]):not([readonly])")?.focus(); const sel = $("select[name=group_id][data-newgroup]"); if (sel) $("#new-group-field").hidden = sel.value !== "__new__"; } - syncSelection(); syncFilters(); + syncSelection(); syncFilters(); syncEnables(); }); // --- автообновление статусов: сервер опрашивает устройства сам, страница подтягивает результат --- @@ -100,6 +100,28 @@ } setInterval(autoRefresh, REFRESH_MS); - const init = () => { syncSelection(); syncFilters(); }; + // --- журнал: копирование ID, подтверждение очистки, строка таблицы открывает запись --- + function legacyCopy(text) { // запасной путь: http вне localhost или нет разрешения на буфер обмена + const ta = document.createElement("textarea"); ta.value = text; ta.style.position = "fixed"; ta.style.opacity = "0"; + document.body.appendChild(ta); ta.select(); + try { document.execCommand("copy"); } finally { ta.remove(); } + } + function copyText(text) { + if (navigator.clipboard && window.isSecureContext) { + return navigator.clipboard.writeText(text).catch(() => legacyCopy(text)); + } + legacyCopy(text); + return Promise.resolve(); + } + document.addEventListener("click", (e) => { + const c = e.target.closest("[data-copy]"); + if (c) { copyText(c.dataset.copy).then(() => { const t = c.title; c.title = "Скопировано"; c.setAttribute("aria-label", "Скопировано"); setTimeout(() => { c.title = t; c.setAttribute("aria-label", "Копировать"); }, 1200); }); return; } + const tr = e.target.closest("tr.clickable"); + if (tr && !e.target.closest("a, button")) tr.querySelector("a.row-link")?.click(); + }); + const syncEnables = () => $$("input[data-enables]").forEach((i) => { const t = $(i.dataset.enables); if (t) t.disabled = i.disabled || !i.value; }); + document.addEventListener("input", (e) => { if (e.target.matches("input[data-enables]")) syncEnables(); }); + + const init = () => { syncSelection(); syncFilters(); syncEnables(); }; document.readyState === "loading" ? document.addEventListener("DOMContentLoaded", init) : init(); })(); diff --git a/app/ui/static/style.css b/app/ui/static/style.css index 42f789e..b413227 100644 --- a/app/ui/static/style.css +++ b/app/ui/static/style.css @@ -205,3 +205,28 @@ dialog.modal::backdrop { background: var(--backdrop); } .card > table.grid:first-child thead th:first-child { border-top-left-radius: 11px; } .card > table.grid:first-child thead th:last-child { border-top-right-radius: 11px; } .card > table.grid:last-child tbody tr:last-child td { border-bottom: 0; } + +/* ---------- Журнал ---------- */ +.badge.plain { background: var(--neu-b); color: var(--neu-t); } +table.grid.events { min-width: 1000px; } +.grid.events td { height: 48px; } +.grid.events td.c-small:first-child { white-space: nowrap; } +.grid.events td.msg-cell { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.grid.events tr.clickable { cursor: pointer; } +.grid.events tr.clickable:hover td { background: var(--acc50); } +.grid.events .ent { margin-right: 8px; } +.grid.events a.row-link { color: var(--muted); } +.more { display: flex; justify-content: center; align-items: center; min-height: 0; } +.more:not(:empty) { height: 64px; border-bottom: 1px solid var(--line); } +dialog.modal:has(.w640) { width: 640px; } +.kv { display: grid; grid-template-columns: 130px 1fr; gap: 12px; align-items: center; min-height: 32px; } +.kv > span:first-child { color: var(--muted); font-size: 13px; } +.kv > div { display: flex; align-items: center; gap: 8px; min-width: 0; flex-wrap: wrap; } +.kv .mono { overflow-wrap: anywhere; } +pre.json { margin: 0; padding: 12px; border-radius: 8px; background: var(--head); border: 1px solid var(--line); font: 400 13px/1.4 var(--mono); white-space: pre-wrap; overflow-wrap: anywhere; } +.sum { padding: 12px; border-radius: 8px; background: var(--head); border: 1px solid var(--line); font: 400 13px/1.5 var(--font); color: var(--muted); } +.sum b { color: var(--text); font-weight: 600; } +.note.info { background: var(--acc50); color: var(--accd); } +.input.invalid { border-color: var(--bad-t); } +.input:disabled { background: var(--disabled-bg); color: var(--disabled-t); cursor: not-allowed; } +.field-err { font: 500 13px/1.4 var(--font); color: var(--bad-t); } diff --git a/app/ui/templates/_device_form.html b/app/ui/templates/_device_form.html index cc69e57..dc2d51e 100644 --- a/app/ui/templates/_device_form.html +++ b/app/ui/templates/_device_form.html @@ -26,7 +26,7 @@ diff --git a/app/ui/templates/_event_dialog.html b/app/ui/templates/_event_dialog.html new file mode 100644 index 0000000..e64a29e --- /dev/null +++ b/app/ui/templates/_event_dialog.html @@ -0,0 +1,23 @@ +{% import "_ui.html" as ui %} +
+ + + +
diff --git a/app/ui/templates/_events_chunk.html b/app/ui/templates/_events_chunk.html new file mode 100644 index 0000000..6b55b9c --- /dev/null +++ b/app/ui/templates/_events_chunk.html @@ -0,0 +1,2 @@ +{% include "_events_rows.html" %} +{% include "_events_more.html" %} diff --git a/app/ui/templates/_events_clear.html b/app/ui/templates/_events_clear.html new file mode 100644 index 0000000..1b7c6b2 --- /dev/null +++ b/app/ui/templates/_events_clear.html @@ -0,0 +1,20 @@ +{% import "_ui.html" as ui %} +
+ + +
+ + {% if in_modal %}{% else %}Отмена{% endif %} + +
+
diff --git a/app/ui/templates/_events_more.html b/app/ui/templates/_events_more.html new file mode 100644 index 0000000..8496e00 --- /dev/null +++ b/app/ui/templates/_events_more.html @@ -0,0 +1,2 @@ +
{% if more %}{% endif %}
+
Показано {{ shown }} из {{ total|num }} {{ total|plural("записи", "записей", "записей") }}
diff --git a/app/ui/templates/_events_rows.html b/app/ui/templates/_events_rows.html new file mode 100644 index 0000000..1231358 --- /dev/null +++ b/app/ui/templates/_events_rows.html @@ -0,0 +1,12 @@ +{% for e in rows %} + + {{ e.ts|dts }} + {{ e.type }} + {{ ev_entity(e.entity_type) }}{{ (e.entity_id|short_id) if e.entity_id else "—" }} + {{ e.message }} + {{ e.actor }} + {{ e.id|short_id }} + +{% else %} +{% if not oob %}Записей не найдено.{% endif %} +{% endfor %} diff --git a/app/ui/templates/_events_settings.html b/app/ui/templates/_events_settings.html new file mode 100644 index 0000000..5f6fe72 --- /dev/null +++ b/app/ui/templates/_events_settings.html @@ -0,0 +1,21 @@ +{% import "_ui.html" as ui %} +
+ + +
+ + {% if in_modal %}{% else %}Отмена{% endif %} + +
+
diff --git a/app/ui/templates/_jobs.html b/app/ui/templates/_jobs.html index 87567bb..344eb71 100644 --- a/app/ui/templates/_jobs.html +++ b/app/ui/templates/_jobs.html @@ -11,12 +11,12 @@ {% if flash %}

{{ flash }}

{% endif %}
- - + + {% for j in jobs %} - + diff --git a/app/ui/templates/backups.html b/app/ui/templates/backups.html index 5b3c434..eb0bb78 100644 --- a/app/ui/templates/backups.html +++ b/app/ui/templates/backups.html @@ -20,7 +20,7 @@ - +
#УстройствоОперацияСтатусСообщениеСоздана (UTC)
IDУстройствоОперацияСтатусСообщениеСоздана (UTC)
{{ j.id }}{{ j.id|short_id }} {{ j.device_name }} {{ {"backup": "Бэкап", "ros_update": "Обновление ROS", "fw_update": "Обновление FW"}.get(j.type, j.type) }} {{ {"pending": "В очереди", "running": "Выполняется", "done": "Готово", "failed": "Ошибка"}.get(j.status, j.status) }} {{ i.device or "—" }} {{ i.group or "—" }} {{ {"backup": ".backup", "rsc": ".rsc"}.get(i.kind, "—") }}{{ i.key.split("/")[-1] }}{{ i.key.split("/")[-1] }} {{ i.size|size }} {{ i.last_modified|dt }} diff --git a/app/ui/templates/base.html b/app/ui/templates/base.html index 5bf1446..4302f70 100644 --- a/app/ui/templates/base.html +++ b/app/ui/templates/base.html @@ -11,7 +11,7 @@ {% if request.session.get("user") %} - {% set nav = [("Устройства", "/", "devices"), ("Группы", "/groups", "groups"), ("Резервные копии", "/backups", "backups")] %} + {% set nav = [("Устройства", "/", "devices"), ("Группы", "/groups", "groups"), ("Резервные копии", "/backups", "backups"), ("Журнал", "/events", "events")] %}
ros_control