Глобально уникальные ID (docs/changes/016): - у устройств, групп, задач, резервных копий и записей журнала ID вида <префикс>_<uuid7> (dev_, grp_, job_, bkp_, evt_): типы не пересекаются, внутри типа ID не повторяются и сортируются по времени; - миграция БД v1 заменяет числовые ID с пересчётом ссылок (копия файла БД перед миграцией, сверка числа строк, одна транзакция); - резервная копия = пара файлов с одним bkp_ ID, файлы в S3 получают метаданные backup-id/device-id; синхронизация метаданных с бакетом; - журнал событий в БД (events): создание/изменение/удаление, задачи, бэкапы, смена online/offline, вход в UI; API чтения GET /api/v1/events. Журнал в UI, ротация и очистка (docs/changes/017): - страница «Журнал»: фильтры, подгрузка «Показать ещё», окно записи; - настройки ротации (срок и максимум записей) хранятся в БД (схема v2), ротация при старте, раз в час и после сохранения настроек; - очистка журнала только через окно с паролем пользователя, блокировка после 5 неверных попыток, остаётся запись о факте очистки; через API очистки нет. Тесты: 20 из 20. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
94 lines
3.8 KiB
Python
94 lines
3.8 KiB
Python
"""Фоновые задачи: запись в таблицу jobs, параллелизм ограничен MAX_CONCURRENCY."""
|
|
import asyncio
|
|
|
|
from app import ids
|
|
from app.config import get_settings
|
|
from app.db import session_scope
|
|
from app.models import Device, Job, now
|
|
from app.services import events, ops
|
|
|
|
JOB_TYPES = {
|
|
"backup": ops.run_backup,
|
|
"ros_update": ops.run_ros_update,
|
|
"fw_update": ops.run_fw_update,
|
|
}
|
|
|
|
_tasks: set[asyncio.Task] = set()
|
|
_sem: asyncio.Semaphore | None = None
|
|
|
|
|
|
def _semaphore() -> asyncio.Semaphore:
|
|
global _sem
|
|
if _sem is None:
|
|
_sem = asyncio.Semaphore(get_settings().max_concurrency)
|
|
return _sem
|
|
|
|
|
|
def _finish(job_id: str, status: str, message: str) -> None:
|
|
with session_scope() as s:
|
|
j = s.get(Job, job_id)
|
|
j.status, j.message = status, message
|
|
if status in ("done", "failed"):
|
|
j.finished_at = now()
|
|
etype = {"running": "job.started", "done": "job.done", "failed": "job.failed"}[status]
|
|
events.record(etype, "job", job_id, message or f"Задача {j.type}: {status}", device_id=j.device_id,
|
|
job_id=job_id, data={"type": j.type, "status": status}, s=s)
|
|
|
|
|
|
async def _run(job_id: str, job_type: str, device_id: str) -> None:
|
|
events.set_job(job_id) # события и бэкап, созданные внутри задачи, ссылаются на её ID
|
|
async with _semaphore():
|
|
_finish(job_id, "running", "")
|
|
try:
|
|
_finish(job_id, "done", await JOB_TYPES[job_type](device_id))
|
|
except Exception as e: # noqa: BLE001 — любой сбой фиксируем в задаче
|
|
_finish(job_id, "failed", str(e))
|
|
|
|
|
|
def start_jobs(job_type: str, device_ids: list[str]) -> list[str]:
|
|
"""Создаёт задачи для списка устройств и запускает их в фоне. Возвращает ID задач."""
|
|
if job_type not in JOB_TYPES:
|
|
raise ValueError(f"Неизвестный тип задачи: {job_type}")
|
|
for did in device_ids:
|
|
ids.check(did, "dev")
|
|
started = []
|
|
with session_scope() as s:
|
|
for did in device_ids:
|
|
d = s.get(Device, did)
|
|
if d is None:
|
|
raise LookupError(f"Устройство {did} не найдено")
|
|
j = Job(device_id=did, device_name=d.name, type=job_type)
|
|
s.add(j)
|
|
s.flush()
|
|
events.record("job.created", "job", j.id, f"Задача {job_type} для {d.name} создана", device_id=did,
|
|
job_id=j.id, data={"type": job_type}, s=s)
|
|
started.append((j.id, did))
|
|
for jid, did in started:
|
|
task = asyncio.create_task(_run(jid, job_type, did)) # наследует контекст (актор)
|
|
_tasks.add(task)
|
|
task.add_done_callback(_tasks.discard)
|
|
return [jid for jid, _ in started]
|
|
|
|
|
|
def list_jobs(limit: int = 50) -> list[Job]:
|
|
with session_scope() as s:
|
|
return list(s.query(Job).order_by(Job.id.desc()).limit(limit))
|
|
|
|
|
|
def get_job(job_id: str) -> Job:
|
|
ids.check(job_id, "job")
|
|
with session_scope() as s:
|
|
j = s.get(Job, job_id)
|
|
if j is None:
|
|
raise LookupError(f"Задача {job_id} не найдена")
|
|
return j
|
|
|
|
|
|
def fail_stale_jobs() -> None:
|
|
"""После перезапуска сервера незавершённые задачи помечаются как прерванные."""
|
|
with session_scope() as s:
|
|
for j in s.query(Job).filter(Job.status.in_(("pending", "running"))):
|
|
j.status, j.message, j.finished_at = "failed", "Прервано перезапуском сервера", now()
|
|
events.record("job.failed", "job", j.id, j.message, device_id=j.device_id, job_id=j.id,
|
|
data={"type": j.type, "status": "failed"}, s=s)
|