Кэш списка бакета, БД без блокировки event loop, один процесс на БД
Производительность и масштабирование, пункты 8–10 ревью (docs/changes/019): - кэш списка бакета и связей файлов с копиями (BACKUPS_CACHE_TTL, 60 с): sync_rows только при реальном чтении бакета; сброс после бэкапа и удаления, счётчик поколений против гонки; refresh=1 и «Обновить список»; - обработчики API/UI без await — обычные функции (пул потоков FastAPI), запуск задач остаётся async; запись статуса, задачи и sync_rows — через asyncio.to_thread; SQLite: WAL, synchronous=NORMAL, busy_timeout; - файловая блокировка <файл БД>.lock: второй процесс на той же БД не стартует; раздел «Ограничения» в README. Тесты: 24 из 24. Стенд (порт 8001): кэш 0,014 с против 0,138 с, второй uvicorn на боевой БД отклонён блокировкой, боевые данные не изменены. Ручная проверка UI пользователем на момент коммита не подтверждена. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
1 parent
123b5abdfc
commit
ba8ac7be71
17 files changed
+369
-84
No files matched your search
@@ -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
|
||||
|
||||
|
||||
@@ -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=<evt_…>`, `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`; файловая блокировка БД — один процесс на БД.
|
||||
+18
-17
@@ -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,20 +283,20 @@ 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,
|
||||
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);
|
||||
@@ -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)
|
||||
@@ -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 # секунд на установление соединения: так быстро определяется недоступность
|
||||
|
||||
@@ -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 (регистрация моделей)
|
||||
|
||||
|
||||
@@ -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,6 +18,8 @@ from app.ui import routes as ui
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
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
|
||||
@@ -28,6 +31,8 @@ async def lifespan(app: FastAPI):
|
||||
t.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await t
|
||||
finally:
|
||||
process_lock.release()
|
||||
|
||||
|
||||
def create_app() -> FastAPI:
|
||||
|
||||
@@ -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
|
||||
+51
-15
@@ -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)
|
||||
@@ -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]:
|
||||
|
||||
+18
-9
@@ -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"]
|
||||
|
||||
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()
|
||||
backup_id = b.id
|
||||
events.record("backup.created", "backup", backup_id, f"Резервная копия {conn.name} запрошена", device_id=device_id,
|
||||
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:
|
||||
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}"
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
+34
-31
@@ -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("Очистка журнала выполняется только через окно подтверждения")
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
<h1>Резервные копии</h1>
|
||||
<div class="sub">Бакет {{ bucket }} · {{ stats.total }} {{ stats.total|plural("файл", "файла", "файлов") }} · {{ stats.size|size }}</div>
|
||||
</div>
|
||||
<div class="actions"><a class="btn secondary" href="/backups?{{ request.url.query }}">{{ ui.icon("refresh") }}Обновить список</a></div>
|
||||
<div class="actions"><a class="btn secondary" href="{{ refresh_url }}">{{ ui.icon("refresh") }}Обновить список</a></div>
|
||||
</div>
|
||||
|
||||
{% if deleted >= 0 %}<div class="{{ 'form-err' if failed else 'form-ok' }}" role="status" style="margin-bottom:16px">
|
||||
|
||||
@@ -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 <override 8001> 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 — пользователь: страница «Бэкапы» быстро открывается при смене фильтров, «Обновить список» подтягивает изменения бакета.
|
||||
@@ -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 пользователем на момент коммита не подтверждена.
|
||||
@@ -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()
|
||||
+53
-3
@@ -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()
|
||||
Reference in new issue
Block a user