Пункт 18 ревью 2026-09-28 21:32 (docs/changes/027): после отката ROS
плата остаётся на более новой прошивке, а приложение сравнивало версии
на неравенство и показывало откат как обновление.
- devices.fw_state: unknown/update/downgrade/current; has_fw_update и
фильтры не считают откат обновлением; метка «в ROS: X».
- upgrade_firmware не понижает прошивку; общая запись — _flash_firmware.
- Задача fw_downgrade: проверка → обязательный бэкап → запись →
перезагрузка; API POST /devices/{id}/firmware/downgrade и
/batch/fw_downgrade (target_version обязателен).
- UI: «Откатить прошивку…» (строка и группа), общее окно отката с kind.
Тесты: 39 из 39. Стенд: временное недоступное устройство — 422/202,
задачи failed на проверке, бэкапов нет. Реальная запись прошивки при
откате и ручная проверка UI не выполнялись.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
106 lines
4.7 KiB
Python
106 lines
4.7 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,
|
|
"set_channel": ops.run_set_channel,
|
|
"ros_downgrade": ops.run_ros_downgrade,
|
|
"fw_downgrade": ops.run_fw_downgrade, # изменение 027
|
|
}
|
|
|
|
_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, params: dict) -> None:
|
|
events.set_job(job_id) # события и бэкап, созданные внутри задачи, ссылаются на её ID
|
|
async with _semaphore():
|
|
await asyncio.to_thread(_finish, job_id, "running", "")
|
|
try:
|
|
result = await JOB_TYPES[job_type](device_id, **params)
|
|
await asyncio.to_thread(_finish, job_id, "done", result)
|
|
except Exception as e: # noqa: BLE001 — любой сбой фиксируем в задаче
|
|
await asyncio.to_thread(_finish, job_id, "failed", str(e))
|
|
|
|
|
|
def _create_jobs(job_type: str, device_ids: list[str], params: dict) -> list[tuple[str, str]]:
|
|
"""Синхронная часть start_jobs: строки задач и события job.created — в потоке."""
|
|
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, **params}, s=s)
|
|
started.append((j.id, did))
|
|
return started
|
|
|
|
|
|
async def start_jobs(job_type: str, device_ids: list[str], params: dict | None = None) -> list[str]:
|
|
"""Создаёт задачи для списка устройств и запускает их в фоне. params передаются раннеру именованными
|
|
аргументами (device_id, **params) и попадают в data события job.created. Возвращает ID задач."""
|
|
params = params or {}
|
|
started = await asyncio.to_thread(_create_jobs, job_type, device_ids, params)
|
|
for jid, did in started:
|
|
task = asyncio.create_task(_run(jid, job_type, did, params)) # наследует контекст (актор)
|
|
_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)
|