diff --git a/.env.example b/.env.example index 6aba488..d80a32c 100644 --- a/.env.example +++ b/.env.example @@ -20,6 +20,9 @@ S3_ACCESS_KEY= S3_SECRET_KEY= S3_PREFIX=backups +# Кэш списка бакета (страница «Бэкапы», GET /api/v1/backups): секунд без обращения к S3, 0 — выключить кэш +BACKUPS_CACHE_TTL=60 + # Ограничение параллельных обращений к устройствам MAX_CONCURRENCY=10 diff --git a/README.md b/README.md index 8f93c3d..de4859c 100644 --- a/README.md +++ b/README.md @@ -55,6 +55,11 @@ Admin Dashboard ⇄ Control Server ⇄ RouterOS REST (на каждом устр - Пароли устройств шифруются ключом `SECRET_KEY` (Fernet); потеря ключа = потеря доступа к сохранённым паролям. - **Секреты в бэкапах:** `.rsc` создаётся с `show-sensitive`, а `.backup` — без шифрования, поэтому файлы содержат пароли и ключи открытым текстом. Ограничьте доступ к бакету и ссылкам скачивания. +## Ограничения + +- **Один процесс на БД**: сервер держит файловую блокировку `<файл БД>.lock` рядом с БД (снимается при остановке); второй процесс на той же БД (`--workers 2+`, вторая копия контейнера на том же томе) не стартует — понятная ошибка вместо молчаливой порчи данных. В памяти процесса (не переживает перезапуск и не разделяется между процессами) — семафор фоновых задач, блокировка синхронизации бэкапов, счётчики неудачных попыток очистки журнала и кэш списка бакета. +- **Кэш списка бакета**: страница «Бэкапы» и `GET /api/v1/backups` не перечитывают бакет на каждый просмотр — список живёт `BACKUPS_CACHE_TTL` секунд (по умолчанию 60; `0` — кэш выключен). Изменения бакета, сделанные не через это приложение, видны не позже TTL или сразу — кнопкой «Обновить список» (`refresh=1`). Собственные изменения (бэкап, удаление) сбрасывают кэш сами. + ## Запуск ```bash @@ -76,6 +81,7 @@ docker compose up -d --build # UI: http://localhost:8000, OpenAPI: /docs | `API_TOKEN` | `change-me` | Bearer-токен API | | `DATABASE_URL` | `sqlite:///./data/ros_control.db` | база метаданных (в compose — том) | | `S3_ENDPOINT`, `S3_REGION`, `S3_BUCKET`, `S3_ACCESS_KEY`, `S3_SECRET_KEY`, `S3_PREFIX` | Yandex Object Storage, `backups` | бакет для резервных копий | +| `BACKUPS_CACHE_TTL` | `60` | кэш списка бакета, с; `0` — выключить (см. «Ограничения») | | `MAX_CONCURRENCY` | `10` | одновременные обращения к устройствам | | `ROS_CONNECT_TIMEOUT` | `4` | секунд на установление соединения (скорость обнаружения недоступности) | | `ROS_TIMEOUT` | `30` | секунд на ответ устройства (долгие операции) | @@ -89,7 +95,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 # 22 теста, фоновый опрос в тестах выключен +./venv/bin/python -m pytest # 24 теста, фоновый опрос в тестах выключен ``` ## API v1 @@ -115,7 +121,7 @@ curl -s -H "Authorization: Bearer $API_TOKEN" http://localhost:8000/api/v1/devic | PATCH/DELETE | `/api/v1/groups/{id}` | переименовать / удалить (устройства остаются без группы) | | POST | `/api/v1/batch/{backup\|ros_update\|fw_update}` | `{"device_ids": [...]}` и/или `{"group_id": N}` — групповая операция (задача) | | PUT | `/api/v1/batch/channel` | `{"device_ids": [...] или "group_id": N, "channel": "..."}` — групповая смена канала (задача `set_channel`) → 202 `{"job_ids": [...]}` | -| GET/DELETE | `/api/v1/backups`, `/backups/download?key=` | бэкапы в бакете (фильтры `device_id`, `group`, `kind`, `date_from`, `date_to`, `q`) | +| GET/DELETE | `/api/v1/backups`, `/backups/download?key=` | бэкапы в бакете (фильтры `device_id`, `group`, `kind`, `date_from`, `date_to`, `q`; `refresh=1` — минуя кэш) | | 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`) | @@ -154,3 +160,4 @@ curl -s -H "Authorization: Bearer $API_TOKEN" http://localhost:8000/api/v1/devic - `016-unique-ids-and-event-log` — глобально уникальные ID (префикс + UUIDv7) для всех сущностей, журнал событий в БД, миграция числовых ID. - `017-events-ui-rotation` — страница «Журнал» в UI, настройки ротации, очистка журнала с подтверждением паролем. - `018-correctness-consistency` — одиночное удаление бэкапа в UI через общий сервис `delete_many`, единая система миграций (колонки старой схемы — в `migrations.run` до миграции ID), групповая смена канала фоновыми задачами (`set_channel`). +- `019-performance-scaling` — кэш списка бакета (`BACKUPS_CACHE_TTL`, «Обновить список»); обработчики без обращений к event loop — обычные функции (пул потоков FastAPI), запись статуса и тяжёлые операции с БД в фоне — через `asyncio.to_thread`; SQLite — WAL и `busy_timeout`; файловая блокировка БД — один процесс на БД. diff --git a/app/api/v1.py b/app/api/v1.py index 078a7c0..d87d907 100644 --- a/app/api/v1.py +++ b/app/api/v1.py @@ -110,7 +110,7 @@ class BatchChannelIn(BatchIn): # --- устройства --- @router.get("/devices", response_model=list[DeviceOut]) -async def list_devices( +def list_devices( group: str = "", # "" — все, "none" — без группы, иначе ID группы (grp_…) q: str = "", status: Literal["", "online", "offline"] = "", @@ -124,44 +124,44 @@ async def list_devices( # --- группы --- @router.get("/groups", response_model=list[GroupOut]) -async def list_groups(): +def list_groups(): return groups.list_groups() @router.post("/groups", response_model=GroupOut, status_code=201) -async def create_group(body: GroupIn): +def create_group(body: GroupIn): g = groups.create_group(body.name) return GroupOut(id=g.id, name=g.name) @router.patch("/groups/{group_id}", response_model=GroupOut) -async def rename_group(group_id: str, body: GroupIn): +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: str): +def delete_group(group_id: str): groups.delete_group(group_id) @router.post("/devices", response_model=DeviceOut, status_code=201) -async def create_device(body: DeviceIn): +def create_device(body: DeviceIn): return DeviceOut.of(devices.create_device(**body.model_dump())) @router.get("/devices/{device_id}", response_model=DeviceOut) -async def get_device(device_id: str): +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: str, body: DevicePatch): +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: str): +def delete_device(device_id: str): devices.delete_device(device_id) @@ -224,13 +224,14 @@ async def list_backups( date_from: date | None = None, date_to: date | None = None, q: str = "", + refresh: bool = False, # принудительно перечитать бакет, минуя кэш (BACKUPS_CACHE_TTL) ): name = devices.get_device(device_id).name if device_id else "" - return await backups.list_backups(name, group, kind, date_from, date_to, q) + return await backups.list_backups(name, group, kind, date_from, date_to, q, refresh=refresh) @router.get("/backups/download") -async def download_backup(key: str): +def download_backup(key: str): if not s3.key_allowed(key): raise ValueError("Недопустимый ключ") return RedirectResponse(s3.presign_get(key)) @@ -282,22 +283,22 @@ class JournalSettingsIn(BaseModel): @router.get("/events/settings") -async def get_journal_settings(): +def get_journal_settings(): """Настройки ротации журнала и сводка (число записей, самая старая запись).""" return {**settings.journal(), **events.stats()} @router.put("/events/settings") -async def put_journal_settings(body: JournalSettingsIn): +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): +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, @@ -305,7 +306,7 @@ async def list_events(entity_id: str | None = None, type: str | None = None, dev @router.get("/events/{event_id}", response_model=EventOut) -async def get_event(event_id: str): +def get_event(event_id: str): ids.check(event_id, "evt") return events.get_event(event_id) @@ -313,10 +314,10 @@ async def get_event(event_id: str): # --- задачи --- @router.get("/jobs", response_model=list[JobOut]) -async def list_jobs(): +def list_jobs(): return jobs.list_jobs() @router.get("/jobs/{job_id}", response_model=JobOut) -async def get_job(job_id: str): +def get_job(job_id: str): return jobs.get_job(job_id) diff --git a/app/config.py b/app/config.py index c4a3e3b..26eab19 100644 --- a/app/config.py +++ b/app/config.py @@ -21,6 +21,8 @@ class Settings(BaseSettings): s3_secret_key: str = "" s3_prefix: str = "backups" + backups_cache_ttl: int = 60 # секунд кэш списка бакета живёт без обращения к S3; 0 — кэш выключен + max_concurrency: int = 10 ros_timeout: float = 30.0 # секунд на чтение ответа устройства (долгие операции) ros_connect_timeout: float = 4.0 # секунд на установление соединения: так быстро определяется недоступность diff --git a/app/db.py b/app/db.py index d34873d..800e44a 100644 --- a/app/db.py +++ b/app/db.py @@ -1,7 +1,7 @@ from contextlib import contextmanager from pathlib import Path -from sqlalchemy import create_engine +from sqlalchemy import create_engine, event from sqlalchemy.orm import DeclarativeBase, sessionmaker from app.config import get_settings @@ -19,9 +19,18 @@ def init_db(url: str | None = None) -> None: """Создаёт engine и таблицы. Вызывается при старте приложения (и в тестах).""" global _engine, _SessionLocal url = url or get_settings().database_url - if url.startswith("sqlite:///") and ":memory:" not in url: + is_file_db = url.startswith("sqlite:///") and ":memory:" not in url + if is_file_db: Path(url.removeprefix("sqlite:///")).parent.mkdir(parents=True, exist_ok=True) _engine = create_engine(url, connect_args={"check_same_thread": False}) + if is_file_db: # :memory: — отдельная БД на соединение, WAL и busy_timeout ей не нужны + @event.listens_for(_engine, "connect") + def _pragmas(dbapi_conn, _rec): + cur = dbapi_conn.cursor() + cur.execute("PRAGMA journal_mode=WAL") # несколько читателей не блокируют друг друга и писателя + cur.execute("PRAGMA synchronous=NORMAL") # безопасно с WAL, быстрее FULL + cur.execute("PRAGMA busy_timeout=5000") # ждать снятия блокировки БД, а не падать сразу + cur.close() _SessionLocal = sessionmaker(_engine, expire_on_commit=False) from app import models # noqa: F401 (регистрация моделей) diff --git a/app/main.py b/app/main.py index 36b218b..25d3e45 100644 --- a/app/main.py +++ b/app/main.py @@ -8,6 +8,7 @@ from fastapi.responses import JSONResponse, RedirectResponse from fastapi.staticfiles import StaticFiles from starlette.middleware.sessions import SessionMiddleware +from app import process_lock from app.api import v1 from app.config import get_settings from app.db import init_db @@ -17,17 +18,21 @@ from app.ui import routes as ui @asynccontextmanager 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 - for t in (task, reconcile, rotator): - if t: - t.cancel() - with suppress(asyncio.CancelledError): - await t + process_lock.acquire(get_settings().database_url) # один процесс на файловую БД (--workers 1) + try: + 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 + for t in (task, reconcile, rotator): + if t: + t.cancel() + with suppress(asyncio.CancelledError): + await t + finally: + process_lock.release() def create_app() -> FastAPI: diff --git a/app/process_lock.py b/app/process_lock.py new file mode 100644 index 0000000..77da523 --- /dev/null +++ b/app/process_lock.py @@ -0,0 +1,34 @@ +"""Блокировка файла БД: гарантирует, что с одной файловой БД работает не больше одного процесса +(семафор задач, _sync_lock, кэш бакета и счётчики неудачных паролей — состояние процесса, не БД). +Берётся в lifespan до init_db, снимается при остановке.""" +import fcntl +from pathlib import Path + +_lock_file = None # открытый дескриптор блокировки текущего процесса (None — не бралась) + + +def acquire(database_url: str) -> None: + global _lock_file + if not (database_url.startswith("sqlite:///") and ":memory:" not in database_url): + return # :memory: и не-SQLite URL — блокировка не нужна + db_path = Path(database_url.removeprefix("sqlite:///")) + db_path.parent.mkdir(parents=True, exist_ok=True) + lock_path = db_path.with_name(db_path.name + ".lock") + f = open(lock_path, "w") + try: + fcntl.flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError: + f.close() + raise RuntimeError( + f"БД {db_path} уже используется другим процессом ros_control: " + "поддерживается только один процесс (--workers 1)" + ) from None + _lock_file = f + + +def release() -> None: + global _lock_file + if _lock_file is not None: + fcntl.flock(_lock_file, fcntl.LOCK_UN) + _lock_file.close() + _lock_file = None diff --git a/app/services/backups.py b/app/services/backups.py index 1a3d623..4493557 100644 --- a/app/services/backups.py +++ b/app/services/backups.py @@ -3,6 +3,7 @@ import asyncio import logging import threading from datetime import date +from time import monotonic from app import ids, s3 from app.config import get_settings @@ -13,6 +14,20 @@ from app.services import devices, events, groups log = logging.getLogger("ros_control.backups") _sync_lock = threading.Lock() # один процесс: не создаём дубликаты «сирот» при параллельных запросах +# кэш списка бакета: sync_rows (вся таблица backups) выполняется только при реальном чтении бакета, +# не при каждом просмотре страницы «Бэкапы» — время жизни настраивается BACKUPS_CACHE_TTL. +# gen — счётчик поколений: invalidate() во время уже идущего чтения бакета не должен потеряться — +# иначе конкурентное чтение (начатое до изменения) сохранит устаревший список со свежей меткой времени. +_cache: dict = {"at": None, "gen": 0, "objects": [], "links": {}} +_cache_lock = asyncio.Lock() + + +def invalidate() -> None: + """Сбрасывает кэш списка бакета: вызывается после собственных изменений бакета + (загрузка бэкапа, удаление), чтобы следующее чтение перечитало S3.""" + _cache["gen"] += 1 + _cache["at"] = None + def _kind(key: str) -> str: return "backup" if key.endswith(".backup") else "rsc" if key.endswith(".rsc") else "other" @@ -60,35 +75,55 @@ def sync_rows(objects: list[dict]) -> dict[str, dict]: return by_key +async def bucket_objects(refresh: bool = False) -> list[dict]: + """Список всего бакета с кэшем (TTL = Settings.backups_cache_ttl секунд; 0 — кэш выключен). + При свежем кэше и без refresh возвращает его; иначе под локом (с повторной проверкой после захвата, + чтобы не читать бакет дважды при параллельных запросах) читает s3.list_backups и связывает файлы + с метаданными (sync_rows — тяжёлый запрос по всей таблице backups, поэтому в потоке).""" + ttl = get_settings().backups_cache_ttl + if not refresh and ttl > 0 and _cache["at"] is not None and monotonic() - _cache["at"] < ttl: + return _cache["objects"] + async with _cache_lock: + if not refresh and ttl > 0 and _cache["at"] is not None and monotonic() - _cache["at"] < ttl: + return _cache["objects"] + gen = _cache["gen"] # invalidate() мог случиться после начала чтения — тогда список уже устарел + objects = await s3.list_backups(None) + try: + links = await asyncio.to_thread(sync_rows, objects) + except Exception: # noqa: BLE001 — список важнее связывания + log.exception("sync_rows не выполнен") + links = {} + _cache["objects"], _cache["links"] = objects, links + # если за время чтения кэш инвалидировали — не помечаем его свежим, следующий вызов перечитает бакет + _cache["at"] = monotonic() if _cache["gen"] == gen else None + return objects + + 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) + await bucket_objects(refresh=True) 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]: + date_to: date | None = None, q: str = "", refresh: bool = False) -> tuple[list[dict], dict]: """Отфильтрованные бэкапы и итоги по всему бакету ({"total": N, "size": байты}). Ключи вида backups/<устройство>/<файл>: устройство берётся из ключа, группа — по имени устройства. group: "" — все, "none" — без группы (в т.ч. бэкапы удалённых устройств), иначе ID группы (grp_…). - Даты — по времени изменения объекта (UTC), включительно. Каждому файлу сопоставлен backup_id.""" + Даты — по времени изменения объекта (UTC), включительно. Каждому файлу сопоставлен backup_id. + refresh — принудительно перечитать бакет, минуя кэш.""" 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(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 = {} + everything = await bucket_objects(refresh) + links = _cache["links"] stats = {"total": len(everything), "size": sum(i["size"] for i in everything)} q = q.strip().lower() out = [] @@ -116,8 +151,8 @@ async def search(device: str = "", group: str = "", kind: str = "", date_from: d async def list_backups(device: str = "", group: str = "", kind: str = "", date_from: date | None = None, - date_to: date | None = None, q: str = "") -> list[dict]: - return (await search(device, group, kind, date_from, date_to, q))[0] + date_to: date | None = None, q: str = "", refresh: bool = False) -> list[dict]: + return (await search(device, group, kind, date_from, date_to, q, refresh=refresh))[0] MAX_DELETE = 500 @@ -144,8 +179,9 @@ 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)) + invalidate() + try: # перечитать бакет и метаданные: копии, у которых не осталось файлов, помечаются удалёнными (с событием) + await bucket_objects(refresh=True) except Exception: # noqa: BLE001 - log.exception("sync_rows после удаления не выполнен") + log.exception("bucket_objects после удаления не выполнен") return sum(results), len(results) - sum(results) diff --git a/app/services/jobs.py b/app/services/jobs.py index 1f65c9f..e7b1307 100644 --- a/app/services/jobs.py +++ b/app/services/jobs.py @@ -39,11 +39,12 @@ def _finish(job_id: str, status: str, message: str) -> None: async def _run(job_id: str, job_type: str, device_id: str, params: dict) -> None: events.set_job(job_id) # события и бэкап, созданные внутри задачи, ссылаются на её ID async with _semaphore(): - _finish(job_id, "running", "") + await asyncio.to_thread(_finish, job_id, "running", "") try: - _finish(job_id, "done", await JOB_TYPES[job_type](device_id, **params)) + result = await JOB_TYPES[job_type](device_id, **params) + await asyncio.to_thread(_finish, job_id, "done", result) except Exception as e: # noqa: BLE001 — любой сбой фиксируем в задаче - _finish(job_id, "failed", str(e)) + await asyncio.to_thread(_finish, job_id, "failed", str(e)) def start_jobs(job_type: str, device_ids: list[str], params: dict | None = None) -> list[str]: diff --git a/app/services/ops.py b/app/services/ops.py index dc4f902..7361393 100644 --- a/app/services/ops.py +++ b/app/services/ops.py @@ -11,7 +11,7 @@ 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, events +from app.services import backups, devices, events def _save_status(device_id: str, fields: dict) -> None: @@ -31,30 +31,31 @@ def _save_status(device_id: str, fields: dict) -> None: async def refresh_status(device_id: str) -> None: """Полный опрос: статус, версии и проверка обновлений. Недоступность — не исключение.""" - conn = devices.get_conn(device_id) + conn = await asyncio.to_thread(devices.get_conn, device_id) try: async with devices.open_client(conn) as c: status = await ros.get_status(c) fields = dict(online=True, status_json=json.dumps(status), last_error=None) except RosError as e: fields = dict(online=False, last_error=str(e)) - _save_status(device_id, fields) + await asyncio.to_thread(_save_status, device_id, fields) async def poll_device(device_id: str, full: bool = False) -> None: """Фоновый опрос. Лёгкий (один запрос system/resource) обновляет online, uptime и версию ROS; полный — если нужна проверка обновлений (full) либо устройство было недоступно/ещё не опрошено (после возвращения версии могли измениться).""" - d = devices.get_device(device_id) + d = await asyncio.to_thread(devices.get_device, device_id) if full or d.online is not True or not d.status: return await refresh_status(device_id) try: - async with devices.open_client(devices.get_conn(device_id)) as c: + conn = await asyncio.to_thread(devices.get_conn, device_id) + async with devices.open_client(conn) as c: res = await c.get("system/resource") except RosError as e: - return _save_status(device_id, dict(online=False, last_error=str(e))) + return await asyncio.to_thread(_save_status, device_id, dict(online=False, last_error=str(e))) status = {**d.status, "uptime": res.get("uptime"), "ros_version": res.get("version")} - _save_status(device_id, dict(online=True, status_json=json.dumps(status), last_error=None)) + await asyncio.to_thread(_save_status, device_id, dict(online=True, status_json=json.dumps(status), last_error=None)) async def refresh_many(device_ids: list[str] | None = None) -> None: @@ -80,15 +81,18 @@ async def run_backup(device_id: str) -> str: } key_bin, key_rsc = files[f"{base}.backup"], files[f"{base}.rsc"] - with session_scope() as s: - 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) + def _create_row() -> str: + with session_scope() as s: + 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() + events.record("backup.created", "backup", b.id, f"Резервная копия {conn.name} запрошена", device_id=device_id, + data={"key_binary": key_bin, "key_rsc": key_rsc}, s=s) + return b.id + + backup_id = await asyncio.to_thread(_create_row) meta = {"backup-id": backup_id, "device-id": device_id} # связь файлов в бакете с метаданными status, error = "failed", "прервано" @@ -112,12 +116,17 @@ async def run_backup(device_id: str) -> str: error = str(e) raise finally: - 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) + backups.invalidate() # часть файлов могла успеть загрузиться в бакет и при неудаче + + def _finish_row() -> None: + 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) + + await asyncio.to_thread(_finish_row) return f"Бэкап загружен в S3: {key_bin}, {key_rsc}" diff --git a/app/services/poller.py b/app/services/poller.py index 2c61a0c..01c2a1b 100644 --- a/app/services/poller.py +++ b/app/services/poller.py @@ -12,7 +12,7 @@ log = logging.getLogger("ros_control.poller") async def cycle(last_full: dict[str, float]) -> None: """Один проход по всем устройствам (параллельно, не более MAX_CONCURRENCY одновременно).""" s = get_settings() - ids = [d.id for d in devices.list_devices()] + ids = [d.id for d in await asyncio.to_thread(devices.list_devices)] for gone in set(last_full) - set(ids): del last_full[gone] sem = asyncio.Semaphore(s.max_concurrency) diff --git a/app/ui/routes.py b/app/ui/routes.py index de4ff7f..52c44cf 100644 --- a/app/ui/routes.py +++ b/app/ui/routes.py @@ -131,12 +131,12 @@ def _tabs(flt: dict) -> list[dict]: # --- вход / выход --- @router.get("/login", response_class=HTMLResponse) -async def login_form(request: Request): +def login_form(request: Request): return _render(request, "login.html", error=None) @router.post("/login") -async def login(request: Request, username: str = Form(), password: str = Form()): +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]}) @@ -148,7 +148,7 @@ async def login(request: Request, username: str = Form(), password: str = Form() @router.post("/logout") -async def logout(request: Request): +def logout(request: Request): request.session.clear() return RedirectResponse("/login", status_code=303) @@ -156,7 +156,7 @@ async def logout(request: Request): # --- дашборд --- @router.get("/", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dashboard(request: Request): +def dashboard(request: Request): flt = _flt(request.query_params) return _render(request, "dashboard.html", section="devices", jobs=jobs.list_jobs(15), flt=flt, tabs=_tabs(flt), groups=groups.list_groups(), page_url=_page_url(flt), @@ -164,7 +164,7 @@ async def dashboard(request: Request): @router.get("/ui/devices", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def devices_table(request: Request): +def devices_table(request: Request): flt = _flt(request.query_params) resp = _render(request, "_devices.html", oob=True, **_devices_ctx(flt)) if request.query_params.get("poll") != "1": # автообновление не должно засорять историю браузера @@ -173,7 +173,7 @@ async def devices_table(request: Request): @router.get("/ui/jobs", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def jobs_table(request: Request): +def jobs_table(request: Request): return _render(request, "_jobs.html", jobs=jobs.list_jobs(15)) @@ -214,7 +214,7 @@ async def device_action(request: Request, device_id: str, action: str, channel: @router.post("/ui/move", dependencies=[Depends(require_login)]) -async def move(device_ids: list[str] = Form(default=[]), group_id: str = Form(""), next: str = Form("/")): +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_id(group_id, "grp")) @@ -247,26 +247,26 @@ 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 = ""): +def dialog_device_new(request: Request, group: str = ""): v = _device_values() 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: str): +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()) @router.get("/ui/dialog/group", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_group_new(request: Request): +def dialog_group_new(request: Request): return _render(request, "_group_form.html", in_modal=True, group=None, gname="", error=None) @router.get("/ui/dialog/group/{group_id}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_group_edit(request: Request, group_id: str): +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) @@ -274,7 +274,7 @@ async def dialog_group_edit(request: Request, group_id: str): # --- группы --- @router.get("/groups", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def groups_page(request: Request, error: str = ""): +def groups_page(request: Request, error: str = ""): everyone = devices.list_devices() return _render(request, "groups.html", section="groups", groups=groups.list_groups(), error=error, total=len(everyone), ungrouped=sum(1 for d in everyone if d.group_id is None)) @@ -291,17 +291,17 @@ def _group_action(request: Request, fn, *args, name: str = "", group=None): @router.post("/groups/new", dependencies=[Depends(require_login)]) -async def group_create(request: Request, name: str = Form()): +def group_create(request: Request, name: str = Form()): return _group_action(request, groups.create_group, name, name=name) @router.post("/groups/{group_id}/rename", dependencies=[Depends(require_login)]) -async def group_rename(request: Request, group_id: str, name: str = Form()): +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: str): +def group_delete(request: Request, group_id: str): return _group_action(request, groups.delete_group, group_id) @@ -317,7 +317,7 @@ def _group_choice(group_id: str, new_group: str) -> tuple[str | None, str | None @router.get("/devices/new", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def device_new(request: Request, group: str = ""): +def device_new(request: Request, group: str = ""): v = _device_values() v["group_id"] = group if ids.is_id(group, "grp") else "" return _device_form(request, None, v) @@ -340,13 +340,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: str): +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: str, host: str = Form(), port: int = Form(443), +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("")): @@ -364,7 +364,7 @@ async def device_update(request: Request, device_id: str, host: str = Form(), po @router.post("/devices/{device_id}/delete", dependencies=[Depends(require_login)]) -async def device_delete(request: Request, device_id: str): +def device_delete(request: Request, device_id: str): devices.delete_device(device_id) return _done(request) @@ -380,17 +380,20 @@ def _opt_date(value: str) -> date | None: @router.get("/backups", response_class=HTMLResponse, dependencies=[Depends(require_login)]) async def backups_page(request: Request, device: str = "", group: str = "", kind: str = "", - date_from: str = "", date_to: str = "", q: str = "", deleted: int = -1, failed: int = 0): + date_from: str = "", date_to: str = "", q: str = "", deleted: int = -1, failed: int = 0, + refresh: bool = False): kind = kind if kind in ("backup", "rsc") else "" items, stats, error = [], {"total": 0, "size": 0}, None try: - items, stats = await backups.search(device, group, kind, _opt_date(date_from), _opt_date(date_to), q) + items, stats = await backups.search(device, group, kind, _opt_date(date_from), _opt_date(date_to), q, refresh=refresh) except Exception as e: # noqa: BLE001 — показываем ошибку S3 на странице error = str(e) flt = dict(device=device, group=group, kind=kind, date_from=date_from, date_to=date_to, q=q) + non_empty = {k: v for k, v in flt.items() if v} return _render(request, "backups.html", section="backups", items=items, stats=stats, error=error, flt=flt, flt_active=sum(1 for v in flt.values() if v), bucket=get_settings().s3_bucket, - page_url="/backups" + ("?" + urlencode({k: v for k, v in flt.items() if v}) if any(flt.values()) else ""), + page_url="/backups" + ("?" + urlencode(non_empty) if non_empty else ""), + refresh_url="/backups?" + urlencode({**non_empty, "refresh": "1"}), deleted=deleted, failed=failed, devices=devices.list_devices(), groups=groups.list_groups()) @@ -404,7 +407,7 @@ async def backups_delete_many(key: list[str] = Form(default=[]), next: str = For @router.get("/backups/download", dependencies=[Depends(require_login)]) -async def backup_download(key: str): +def backup_download(key: str): if not s3.key_allowed(key): raise ValueError("Недопустимый ключ") return RedirectResponse(s3.presign_get(key)) @@ -457,7 +460,7 @@ def _rotation_text(cfg: dict) -> str: @router.get("/events", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def events_page(request: Request): +def events_page(request: Request): f = _ev_filters(request.query_params) st = events.stats() all_types = events.types() @@ -468,7 +471,7 @@ async def events_page(request: Request): @router.get("/ui/events", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def events_more(request: Request, before: str = "", shown: int = 0): +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)) @@ -488,12 +491,12 @@ def _event_ctx(event_id: str) -> dict: @router.get("/ui/dialog/event/{event_id}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_event(request: Request, event_id: str): +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): +def event_page(request: Request, event_id: str): return _render(request, "event_page.html", section="events", in_modal=False, **_event_ctx(event_id)) @@ -503,12 +506,12 @@ def _settings_dialog(request: Request, error: str | None = None, values: dict | @router.get("/ui/dialog/events-settings", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_events_settings(request: Request): +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("")): +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: @@ -524,12 +527,12 @@ def _clear_dialog(request: Request, error: str | None = None): @router.get("/ui/dialog/events-clear", response_class=HTMLResponse, dependencies=[Depends(require_login)]) -async def dialog_events_clear(request: Request): +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("")): +def events_clear(request: Request, password: str = Form("")): """Очистка журнала — только через окно подтверждения паролем пользователя сессии (через API её нет).""" if not _is_htmx(request): raise ValueError("Очистка журнала выполняется только через окно подтверждения") diff --git a/app/ui/templates/backups.html b/app/ui/templates/backups.html index eb0bb78..05675d3 100644 --- a/app/ui/templates/backups.html +++ b/app/ui/templates/backups.html @@ -6,7 +6,7 @@

Резервные копии

Бакет {{ bucket }} · {{ stats.total }} {{ stats.total|plural("файл", "файла", "файлов") }} · {{ stats.size|size }}
- + {% if deleted >= 0 %}
diff --git a/docs/changes/019-performance-scaling/plan.md b/docs/changes/019-performance-scaling/plan.md new file mode 100644 index 0000000..efadca3 --- /dev/null +++ b/docs/changes/019-performance-scaling/plan.md @@ -0,0 +1,86 @@ +# План: 019 — производительность и масштабирование (по ревью 2026-09-27) + +## Context + +Ревью `docs/reviews/2026-09-27-codebase-review.md`, раздел «Производительность и масштабирование», пункты 8–10: + +- **п. 8** — каждое открытие страницы «Бэкапы» (и `GET /api/v1/backups`) перечитывает весь бакет (`s3.list_backups(None)`) + и запускает `backups.sync_rows` по всей таблице `backups`. Нагрузка и стоимость запросов к S3 растут линейно с числом копий. +- **п. 9** — синхронные запросы к SQLite выполняются прямо в async-обработчиках и фоновых задачах и блокируют event loop. +- **п. 10** — состояние хранится в памяти процесса (семафор задач, `_sync_lock`, счётчики неудачных паролей, кэш п. 8); + запуск нескольких процессов на одной БД (`--workers > 1`, вторая копия контейнера на том же томе) молча ломает эти гарантии. + +Решения пользователя: +- п. 9 — **точечно**: обработчики без `await` становятся обычными `def` (FastAPI выполняет их в пуле потоков); + запись статуса в опросе и в задачах — через `asyncio.to_thread`; SQLite — WAL + `busy_timeout`. Переход на async-драйвер не делаем. +- п. 10 — **файловая блокировка + README**: второй процесс на той же БД не стартует с понятной ошибкой. +- Стенд — `ros_control-ros_control-1` с боевыми данными в томе (удалять существующее нельзя; можно создавать новое и удалять только его). + Порт 8000 на хосте занят посторонним процессом, стенд запускается на **8001** через override-файл вне репозитория (как в 018). + +## Изменения + +### п. 8 — кэш списка бакета (`app/services/backups.py`, `app/services/ops.py`, `app/ui/routes.py`, `app/api/v1.py`, `app/config.py`, `.env.example`) +- В `backups.py`: `_cache = {"at": monotonic | None, "objects": list}` и `asyncio.Lock`. Функция `async def bucket_objects(refresh=False) -> list[dict]`: + при свежем кэше (моложе `BACKUPS_CACHE_TTL`) возвращает его; иначе под lock (с повторной проверкой после захвата) читает `s3.list_backups(None)`, + выполняет `sync_rows` (в `asyncio.to_thread`, см. п. 9) и сохраняет результат. **`sync_rows` вызывается только при чтении бакета**, не при каждом просмотре. + Связи `key -> backup_id/device_id` тоже кэшируются вместе со списком (результат `sync_rows`). +- `def invalidate()` — сбрасывает кэш. Вызывается после собственных изменений бакета: `ops.run_backup` после загрузки + (в `finally`, т. к. часть файлов могла успеть загрузиться) и `delete_many` (там `sync_rows` уже выполняется — заменить на `invalidate()` + + `bucket_objects(refresh=True)`, чтобы не дублировать логику). +- `search()` и `reconcile()` используют `bucket_objects()`; фильтрация — по кэшированному списку, как сейчас. +- `Settings.backups_cache_ttl: int = 60` (секунд; `0` — кэш выключен), в `.env.example` с комментарием. +- Принудительное обновление: параметр `refresh=1` у `GET /backups` (UI) и `GET /api/v1/backups`; кнопка «Обновить список» + на странице «Бэкапы» (`backups.html`) передаёт `refresh=1` (сохраняя текущие фильтры). +- Изменения бакета извне (вручную, другим клиентом) видны не позже чем через TTL или по «Обновить список» — отметить в README. + +### п. 9 — БД без блокировки event loop (`app/api/v1.py`, `app/ui/routes.py`, `app/services/*`, `app/db.py`) +- **Обработчики без `await` → `def`**, кроме тех, что вызывают `jobs.start_jobs` (внутри `asyncio.create_task` — нужен работающий + event loop, в пуле потоков его нет): `create_backup`, `install_update`, `upgrade_firmware`, `batch`, `batch_channel` (API), + `batch` (UI) **остаются `async`**. Зависимость `require_login` остаётся `async`: ContextVar актора, выставленный в ней, + копируется в поток (anyio `to_thread.run_sync` работает в копии контекста) — это нужно подтвердить тестом. +- **Фоновые пути** (`asyncio.to_thread`, контекст копируется — актор и `job_id` сохраняются): + - `ops._save_status`, чтение устройства в `ops.poll_device` / `ops.refresh_status` (`devices.get_device`, `devices.get_conn` — включая расшифровку пароля); + - `poller.cycle`: `devices.list_devices()`; + - `jobs._finish`; + - `ops.run_backup`: создание строки `Backup` и финальная запись статуса; + - `backups.sync_rows` (тяжёлая: вся таблица `backups`). + Сигнатуры сервисов не меняются — оборачивается вызов, а не функция. +- **SQLite** (`app/db.py::init_db`): для файловой БД на каждом соединении `PRAGMA journal_mode=WAL` (один раз достаточно — режим хранится в файле) + и `PRAGMA busy_timeout=5000` через `sqlalchemy.event.listens_for(engine, "connect")`; `synchronous=NORMAL` допустим с WAL. + Для `:memory:` — без изменений. Миграции (`migrations.run`, отдельное `sqlite3`-соединение) работают с WAL без изменений; проверить тестом миграции. + +### п. 10 — один процесс на БД (`app/db.py` или новый `app/process_lock.py`, `app/main.py`, `README.md`) +- При старте (в `lifespan`, до `init_db`) — `fcntl.flock(LOCK_EX | LOCK_NB)` на файл `<файл БД>.lock` рядом с БД. + Не удалось → `RuntimeError("БД <путь> уже используется другим процессом ros_control: поддерживается только один процесс (--workers 1)")`, + приложение не стартует. Блокировка снимается при остановке (закрытие дескриптора в `lifespan`). +- Для `:memory:` и не-SQLite URL блокировка не берётся. +- README, раздел «Запуск»/«Безопасность» или новый короткий «Ограничения»: только один процесс на БД, `--workers 1`; + что хранится в памяти процесса; кэш списка бакета (TTL, «Обновить список»). + +## Тесты (минимально) +- п. 8: два вызова `backups.search` подряд → `s3.list_backups` вызван один раз; `refresh=True` или `invalidate()` → повторный вызов. +- п. 9: событие, созданное через синхронный теперь обработчик API (например, `POST /api/v1/devices`), имеет актора `api`; через UI — `ui:<пользователь>` + (если существующие тесты это уже проверяют — достаточно, что они проходят; иначе добавить assert). +- п. 10: вторая попытка взять блокировку на тот же файл БД (в `tmp_path`) → `RuntimeError`; после освобождения — успешно. +- Существующие тесты не должны требовать правок, кроме подмен `s3.list_backups` в тестах бэкапов (сброс кэша между тестами — `backups.invalidate()` в фикстуре или в самих тестах). + +## Документация +README: ограничения (один процесс, кэш бакета), `BACKUPS_CACHE_TTL` в «Настройки», `refresh=1` у `/api/v1/backups`, число тестов, строка 019 в истории изменений. +`summary.md` — оркестратор после проверки. + +## Исполнение +- Код, тесты, пересборка стенда — исполнитель (Sonnet). Тесты не запускает, не коммитит, `summary.md` не создаёт. +- Оркестратор: ревью диффа, полный прогон тестов, проверки на стенде без удаления существующих данных. +- Пользователь: ручная проверка UI. + +## Проверка +- `./venv/bin/python -m pytest -q` — все зелёные. +- Стенд: `docker compose -f docker-compose.yml -f up -d --build --force-recreate`; `/login` → 200; новый код в контейнере + (`grep -c bucket_objects /srv/app/services/backups.py` ≥ 1); `PRAGMA journal_mode` = `wal`; в логах нет трейсбеков. +- Сверка боевых данных до/после — по согласованной копии БД (`sqlite3.Connection.backup` внутри контейнера; при WAL простой `docker cp` файла БД недостаточен). +- Кэш: время ответа `GET /api/v1/backups` первый и повторный запрос (повторный — без обращения к S3); `refresh=1` — снова читает бакет. +- Event loop: во время медленного запроса к недоступному устройству (`POST /api/v1/devices/{id}/refresh` для временного устройства `192.0.2.1`) + параллельный `GET /api/v1/devices` отвечает без ожидания; временное устройство затем удаляется. +- Один процесс: `docker exec ros_control-ros_control-1 python -c "…"`-запуск второго экземпляра приложения на той же БД (или `uvicorn --port 8002` внутри контейнера) + завершается ошибкой блокировки, работающий стенд не затронут. +- Ручная проверка UI — пользователь: страница «Бэкапы» быстро открывается при смене фильтров, «Обновить список» подтягивает изменения бакета. diff --git a/docs/changes/019-performance-scaling/summary.md b/docs/changes/019-performance-scaling/summary.md new file mode 100644 index 0000000..24daab5 --- /dev/null +++ b/docs/changes/019-performance-scaling/summary.md @@ -0,0 +1,37 @@ +# Итоги: 019 — производительность и масштабирование (по ревью 2026-09-27) + +Источник — ревью `docs/reviews/2026-09-27-codebase-review.md`, раздел «Производительность и масштабирование» (пункты 8–10). + +## Сделано +- **п. 8 — кэш списка бакета** (`app/services/backups.py::bucket_objects`): список объектов и связи `key → backup_id` кэшируются на + `BACKUPS_CACHE_TTL` секунд (60; `0` — выключен); `sync_rows` выполняется только при реальном чтении бакета. `asyncio.Lock` с повторной + проверкой — параллельные запросы не читают бакет дважды. `invalidate()` вызывается после загрузки бэкапа и удаления; + счётчик поколений не даёт чтению, начатому до `invalidate()`, пометить кэш свежим. Принудительно — `refresh=1` + (`GET /backups`, `GET /api/v1/backups`, кнопка «Обновить список»). +- **п. 9 — БД не блокирует event loop**: 44 обработчика API/UI без `await` стали `def` (пул потоков FastAPI); обработчики, вызывающие + `jobs.start_jobs` (нужен работающий event loop для `create_task`), остались `async`. Через `asyncio.to_thread`: чтение устройства и запись + статуса в опросе (`ops.poll_device`, `ops.refresh_status`), `poller.cycle`, `jobs._finish`, строки `Backup` в `ops.run_backup`, `sync_rows`. + SQLite: `journal_mode=WAL`, `synchronous=NORMAL`, `busy_timeout=5000` (только файловая БД). +- **п. 10 — один процесс на БД** (`app/process_lock.py`): `flock` на `<файл БД>.lock` в `lifespan` до `init_db`, снимается при остановке; + второй процесс получает `RuntimeError` и не стартует. README — раздел «Ограничения». + +## Найдено на ревью и исправлено +- Гонка кэша: `invalidate()` во время чтения бакета терялся — устаревший список считался свежим весь TTL. Исправлено счётчиком поколений; + мутационная проверка: без исправления новый тест падает (`4 == 5`), с исправлением проходит. +- Ошибка в тесте гонки (кэш перед сценарием был свежим) — исправлена одной строкой оркестратором. + +## Проверено +- `pytest`: 24 из 24 (новые: `test_bucket_list_is_cached_between_reads`, `test_process_lock_blocks_second_process`). + Актор в событиях из теперь синхронных обработчиков (`api`, `ui:admin`) подтверждён существующими тестами. +- Стенд (порт 8001): новый код, `journal_mode=wal`, в логах без ошибок. +- Кэш: 36 файлов бакета; `refresh=1` — 0,138 с, из кэша — 0,014 с. +- Event loop: пока `POST /devices/{id}/refresh` ждал недоступное устройство (4,0 с), параллельный `GET /devices` ответил за 0,012 с. +- Второй `uvicorn` на боевой БД внутри контейнера: `RuntimeError` блокировки, `Application startup failed`; основной стенд работает, прерванных задач нет. +- Боевые данные: группы 4/4, устройства 14/14, бэкапы 45/45, задачи 71/71 — совпадают по ID; добавлены 3 события временного устройства (удалено). + Сверка — по согласованной копии (`sqlite3 backup` внутри контейнера), т. к. при WAL файл БД без `-wal` неполон. + +## Оговорки +- Кэш, семафор задач и счётчики паролей — в памяти процесса; горизонтальное масштабирование по-прежнему не поддерживается (теперь явно). +- Изменения бакета извне видны не позже TTL или по «Обновить список»; `refresh=1` остаётся в адресе страницы после обновления. +- Для копирования БД с хоста теперь нужен и файл `-wal` (или `sqlite3 .backup`). +- Ручная проверка UI пользователем на момент коммита не подтверждена. diff --git a/tests/conftest.py b/tests/conftest.py index 87730f4..03f7fa8 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -3,6 +3,7 @@ from cryptography.fernet import Fernet from app import db, security from app.config import get_settings +from app.services import backups @pytest.fixture(autouse=True) @@ -17,6 +18,7 @@ def env(tmp_path, monkeypatch): monkeypatch.setenv("ADMIN_PASSWORD", "pw") get_settings.cache_clear() security.reset_failures("admin") # блокировки очистки журнала — в памяти процесса + backups.invalidate() # кэш списка бакета — модульное состояние, не должен переживать тест db.init_db() yield get_settings.cache_clear() diff --git a/tests/test_app.py b/tests/test_app.py index e182240..d018f2b 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -9,7 +9,7 @@ import httpx import pytest from fastapi.testclient import TestClient -from app import db, ids, s3, security +from app import db, ids, process_lock, s3, security from app.main import create_app from app.models import Backup, Device, Event, now from app.db import session_scope @@ -520,10 +520,10 @@ async def test_backup_files_get_ids_and_deletion_is_recorded(monkeypatch): found, _ = await backups.search() assert len({i["backup_id"] for i in found}) == 1 and ids.is_id(found[0]["backup_id"], "bkp") # одна копия = пара файлов again, _ = await backups.search() - assert again[0]["backup_id"] == found[0]["backup_id"] # повторный вызов ID не меняет и дубликатов не создаёт + assert again[0]["backup_id"] == found[0]["backup_id"] # повторный вызов (из кэша) ID не меняет и дубликатов не создаёт listing["items"] = [] - await backups.search() + await backups.search(refresh=True) # изменения бакета мимо кэша — принудительное обновление with session_scope() as s: b = s.query(Backup).one() assert b.deleted_at is not None and b.id == found[0]["backup_id"] @@ -635,3 +635,53 @@ def test_journal_page_filters_cursor_dialogs_and_settings(): assert c.get("/api/v1/events/settings", headers=h).json()["retention_days"] == 45 assert c.put("/api/v1/events/settings", headers=h, json={"retention_days": 0, "max_rows": 0}).json()["max_rows"] == 0 assert "ротация: без ограничений" in c.get("/events").text + + +@pytest.mark.asyncio +async def test_bucket_list_is_cached_between_reads(monkeypatch): + """Подряд идущие чтения бакета не обращаются к S3 повторно; refresh=True и invalidate() — обращаются.""" + calls = [] + + async def fake_list(device=None): + calls.append(1) + return [] + + monkeypatch.setattr(s3, "list_backups", fake_list) + await backups.search() + await backups.search() + assert len(calls) == 1 # второй вызов — из кэша + + await backups.search(refresh=True) + assert len(calls) == 2 # принудительное обновление, минуя кэш + + backups.invalidate() + await backups.search() + assert len(calls) == 3 # invalidate() сбрасывает кэш + + # invalidate() пришедший, пока чтение бакета уже шло (гонка: run_backup/delete_many завершились + # во время открытой страницы «Бэкапы») не должен теряться — следующий search() обязан перечитать бакет + async def fake_list_race(device=None): + calls.append(1) + if len(calls) == 4: # ровно один раз, при первом чтении в этом сценарии + backups.invalidate() + return [] + + monkeypatch.setattr(s3, "list_backups", fake_list_race) + backups.invalidate() # кэш после предыдущего шага свежий — начать сценарий с чтения бакета + await backups.search() + assert len(calls) == 4 + await backups.search() + assert len(calls) == 5 # кэш не помечен свежим из-за invalidate() во время предыдущего чтения + + +def test_process_lock_blocks_second_process(tmp_path): + """Второй процесс на той же файловой БД не стартует; после освобождения блокировки — снова можно.""" + db_url = f"sqlite:///{tmp_path}/lock.db" + process_lock.acquire(db_url) + try: + with pytest.raises(RuntimeError, match="уже используется"): + process_lock.acquire(db_url) + finally: + process_lock.release() + process_lock.acquire(db_url) # после освобождения — успешно + process_lock.release()