Files
ros_control/app/services/backups.py
T
ayurishchevandClaude Opus 5.5 0bda0038d0 Остаток п. 9 ревью: синхронная БД вне event loop, повторное ревью
Повторное ревью кодовой базы: docs/reviews/2026-09-28-1243-codebase-review.md
(статус 12 замечаний, новые замечания 13–16).

Остаток п. 9 (docs/changes/022):
- 14 async-функций API, UI и сервисов больше не обращаются к SQLite
  напрямую — через asyncio.to_thread; jobs.start_jobs стал async (БД в
  потоке, create_task в event loop); ops._conn для подключения к устройству;
- тест-линтер по AST: в async def нет прямых вызовов функций с
  session_scope — защита от регресса.

Тесты: 29 из 29. Стенд: задачи и актор событий в порядке, параллельные
запросы не ждут медленного устройства, боевые данные не изменены.
Ручная проверка UI пользователем на момент коммита не подтверждена.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-28 17:12:44 +03:00

191 lines
11 KiB
Python

"""Список бэкапов из бакета с фильтрами и связью с метаданными (backups.id = bkp_…)."""
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
from app.db import session_scope
from app.models import Backup, Device, now
from app.services import devices, events, groups
log = logging.getLogger("ros_control.backups")
_sync_lock = threading.Lock() # один процесс: не создаём дубликаты «сирот» при параллельных запросах
# кэш списка бакета: 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"
def sync_rows(objects: list[dict]) -> dict[str, dict]:
"""Сопоставляет файлы бакета со строками backups (по ключу) и возвращает key -> {backup_id, device_id}.
- файлы без строки (загруженные до появления ID или вручную) получают строку bkp_… по паре <устройство>/<метка>;
- копия, у которой не осталось файлов в бакете, помечается удалённой (строка и ID остаются в истории);
- вернувшиеся файлы снимают пометку."""
keys = {o["key"]: o for o in objects}
with _sync_lock, session_scope() as s:
by_key: dict[str, dict] = {}
for b in s.query(Backup).all():
present = [k for k in (b.key_binary, b.key_rsc) if k and k in keys]
for k in present:
by_key[k] = {"backup_id": b.id, "device_id": b.device_id}
if b.status == "done" and b.deleted_at is None and not present:
b.deleted_at = now()
events.record("backup.deleted", "backup", b.id, f"Файлы резервной копии {b.device_name} удалены из бакета",
device_id=b.device_id, data={"key_binary": b.key_binary, "key_rsc": b.key_rsc}, s=s)
elif present and b.deleted_at is not None:
b.deleted_at = None
orphans: dict[tuple[str, str], dict] = {}
for key, o in keys.items():
parts = key.split("/")
if key in by_key or len(parts) < 3 or _kind(key) == "other":
continue
stem = parts[-1].rsplit(".", 1)[0]
orphans.setdefault((parts[1], stem), {})[_kind(key)] = (key, o["last_modified"])
dev_by_name = {d.name: d.id for d in s.query(Device)}
for (name, _stem), files in orphans.items():
at = max(t for _, t in files.values())
b = Backup(id=ids.new_id("bkp", at), device_id=dev_by_name.get(name), device_name=name, status="done",
requested_at=at, key_binary=files.get("backup", (None,))[0], key_rsc=files.get("rsc", (None,))[0])
s.add(b)
s.flush()
events.record("backup.imported", "backup", b.id, f"Файлы бакета связаны с резервной копией {name}",
device_id=b.device_id, data={"key_binary": b.key_binary, "key_rsc": b.key_rsc}, s=s)
for k in (b.key_binary, b.key_rsc):
if k:
by_key[k] = {"backup_id": b.id, "device_id": b.device_id}
return by_key
async def 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:
events.set_actor("system")
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 = "", refresh: bool = False) -> tuple[list[dict], dict]:
"""Отфильтрованные бэкапы и итоги по всему бакету ({"total": N, "size": байты}).
Ключи вида backups/<устройство>/<файл>: устройство берётся из ключа, группа — по имени устройства.
group: "" — все, "none" — без группы (в т.ч. бэкапы удалённых устройств), иначе ID группы (grp_…).
Даты — по времени изменения объекта (UTC), включительно. Каждому файлу сопоставлен backup_id.
refresh — принудительно перечитать бакет, минуя кэш."""
def _group_maps():
names = {g["id"]: g["name"] for g in groups.list_groups()}
return names, {d.name: names.get(d.group_id) for d in devices.list_devices()}
group_names, device_group = await asyncio.to_thread(_group_maps)
want_group = group_names.get(group) if ids.is_id(group, "grp") else None
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 = []
for item in everything:
parts = item["key"].split("/")
name = parts[1] if len(parts) >= 3 else ""
link = links.get(item["key"], {})
item = {**item, "device": name, "group": device_group.get(name), "kind": _kind(item["key"]),
"backup_id": link.get("backup_id"), "device_id": link.get("device_id")}
if device and name != device:
continue
if group == "none" and item["group"] is not None:
continue
if ids.is_id(group, "grp") and (want_group is None or item["group"] != want_group):
continue
if kind and item["kind"] != kind:
continue
day = item["last_modified"].date()
if (date_from and day < date_from) or (date_to and day > date_to):
continue
if q and q not in item["key"].lower():
continue
out.append(item)
return out, stats
async def list_backups(device: str = "", group: str = "", kind: str = "", date_from: date | None = None,
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
async def delete_many(keys: list[str]) -> tuple[int, int]:
"""Групповое удаление файлов из бакета. Все ключи проверяются до начала удаления
(чужой ключ — ValueError, ничего не удаляется). Возвращает (удалено, не удалось)."""
keys = list(dict.fromkeys(keys))
if not keys:
raise ValueError("Не выбрано ни одного файла")
if len(keys) > MAX_DELETE:
raise ValueError(f"За один раз можно удалить не более {MAX_DELETE} файлов")
if not all(s3.key_allowed(k) for k in keys):
raise ValueError("Недопустимый ключ")
sem = asyncio.Semaphore(8)
async def one(key: str) -> bool:
async with sem:
try:
await s3.delete_object(key)
return True
except Exception: # noqa: BLE001 — считаем неудачей одного файла, остальные удаляем
return False
results = await asyncio.gather(*(one(k) for k in keys))
invalidate()
try: # перечитать бакет и метаданные: копии, у которых не осталось файлов, помечаются удалёнными (с событием)
await bucket_objects(refresh=True)
except Exception: # noqa: BLE001
log.exception("bucket_objects после удаления не выполнен")
return sum(results), len(results) - sum(results)