Повторное ревью кодовой базы: 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>
191 lines
11 KiB
Python
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)
|