Files
ros_control/app/services/jobs.py
T
ayurishchevandClaude Opus 5.5 5f7bd6588a Состояния «Upgrade FW» и осознанный откат прошивки RouterBOARD
Пункт 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>
2026-09-29 00:04:51 +03:00

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)