После смены канала на long-term колонка показывала «актуально», хотя
версия канала (7.23.7) старше установленной (7.24.4) — docs/changes/023:
- состояния колонки: «↑ X», «актуально», «канал: X» (версия канала старше
установленной), «проверка не удалась» (ошибка check-for-updates больше
не маскируется под «актуально»), «—»;
- откат ROS до версии канала как отдельная задача ros_downgrade: проверка
версии на устройстве → обязательный бэкап (сбой прерывает) → install;
- UI: «Откатить ROS…» в меню устройства и групповой пункт в «Обновление»,
окно с вводом целевой версии, список не затрагиваемых устройств;
- API: POST /devices/{id}/update/downgrade и /batch/ros_downgrade с
обязательным target_version. Штатное обновление откат не выполняет.
Тесты: 33 из 33. Стенд: 422 без версии, задачи отката на недоступном
устройстве завершаются на проверке без бэкапа. Реальный откат и ручная
проверка UI пользователем на момент коммита не выполнены.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
170 lines
8.5 KiB
Python
170 lines
8.5 KiB
Python
"""Оркестрация операций: RouterOS + S3 + БД."""
|
|
import asyncio
|
|
import json
|
|
import tempfile
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
|
|
from app import s3
|
|
from app.config import get_settings
|
|
from app.db import session_scope
|
|
from app.models import Backup, Device, now
|
|
from app.ros import operations as ros
|
|
from app.ros.client import RosError
|
|
from app.services import backups, devices, events
|
|
|
|
|
|
def _save_status(device_id: str, fields: dict) -> None:
|
|
with session_scope() as s:
|
|
d = s.get(Device, device_id)
|
|
if d is not None: # устройство могли удалить во время опроса
|
|
was = d.online
|
|
for k, v in fields.items():
|
|
setattr(d, k, v)
|
|
d.status_at = now()
|
|
if "online" in fields and fields["online"] is not was: # в журнал — только смены состояния, не каждый опрос
|
|
up = fields["online"]
|
|
events.record("device.online" if up else "device.offline", "device", d.id,
|
|
f"Устройство {d.name} {'доступно' if up else 'недоступно'}", device_id=d.id,
|
|
data={"error": (fields.get("last_error") or "")[:300]} if not up else None, s=s)
|
|
|
|
|
|
async def _conn(device_id: str) -> devices.Conn:
|
|
return await asyncio.to_thread(devices.get_conn, device_id)
|
|
|
|
|
|
async def refresh_status(device_id: str) -> None:
|
|
"""Полный опрос: статус, версии и проверка обновлений. Недоступность — не исключение."""
|
|
conn = await _conn(device_id)
|
|
try:
|
|
async with devices.open_client(conn) as c:
|
|
status = await ros.get_status(c)
|
|
fields = dict(online=True, status_json=json.dumps(status), last_error=None)
|
|
except RosError as e:
|
|
fields = dict(online=False, last_error=str(e))
|
|
await asyncio.to_thread(_save_status, device_id, fields)
|
|
|
|
|
|
async def poll_device(device_id: str, full: bool = False) -> None:
|
|
"""Фоновый опрос. Лёгкий (один запрос system/resource) обновляет online, uptime и версию ROS;
|
|
полный — если нужна проверка обновлений (full) либо устройство было недоступно/ещё не опрошено
|
|
(после возвращения версии могли измениться)."""
|
|
d = await asyncio.to_thread(devices.get_device, device_id)
|
|
if full or d.online is not True or not d.status:
|
|
return await refresh_status(device_id)
|
|
try:
|
|
conn = await _conn(device_id)
|
|
async with devices.open_client(conn) as c:
|
|
res = await c.get("system/resource")
|
|
except RosError as e:
|
|
return await asyncio.to_thread(_save_status, device_id, dict(online=False, last_error=str(e)))
|
|
status = {**d.status, "uptime": res.get("uptime"), "ros_version": res.get("version")}
|
|
await asyncio.to_thread(_save_status, device_id, dict(online=True, status_json=json.dumps(status), last_error=None))
|
|
|
|
|
|
async def refresh_many(device_ids: list[str] | None = None) -> None:
|
|
targets = device_ids if device_ids is not None else [d.id for d in await asyncio.to_thread(devices.list_devices)]
|
|
sem = asyncio.Semaphore(get_settings().max_concurrency)
|
|
|
|
async def one(i: str):
|
|
async with sem:
|
|
await refresh_status(i)
|
|
|
|
await asyncio.gather(*(one(i) for i in targets))
|
|
|
|
|
|
async def run_backup(device_id: str) -> str:
|
|
"""Бэкап .backup + .rsc: создать на устройстве → скачать через REST во временную папку →
|
|
загрузить в S3 от имени сервера → удалить файлы с устройства."""
|
|
conn = await _conn(device_id)
|
|
stamp = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S")
|
|
base = f"ros_control-{stamp}"
|
|
files = { # имя на устройстве -> ключ в S3
|
|
f"{base}.backup": f"{s3.device_prefix(conn.name)}{stamp}.backup",
|
|
f"{base}.rsc": f"{s3.device_prefix(conn.name)}{stamp}.rsc",
|
|
}
|
|
key_bin, key_rsc = files[f"{base}.backup"], files[f"{base}.rsc"]
|
|
|
|
def _create_row() -> str:
|
|
with session_scope() as s:
|
|
b = Backup(device_id=device_id, device_name=conn.name, job_id=events.current_job(),
|
|
key_binary=key_bin, key_rsc=key_rsc)
|
|
s.add(b)
|
|
s.get(Device, device_id).last_backup_requested_at = b.requested_at = now()
|
|
s.flush()
|
|
events.record("backup.created", "backup", b.id, f"Резервная копия {conn.name} запрошена", device_id=device_id,
|
|
data={"key_binary": key_bin, "key_rsc": key_rsc}, s=s)
|
|
return b.id
|
|
|
|
backup_id = await asyncio.to_thread(_create_row)
|
|
meta = {"backup-id": backup_id, "device-id": device_id} # связь файлов в бакете с метаданными
|
|
|
|
status, error = "failed", "прервано"
|
|
try:
|
|
with tempfile.TemporaryDirectory(prefix="ros_control-") as tmp:
|
|
async with devices.open_client(conn) as c:
|
|
try:
|
|
await ros.create_backup_files(c, base)
|
|
for name, key in files.items():
|
|
local = Path(tmp) / name
|
|
await ros.download_file(c, name, local)
|
|
await s3.upload_file(local, key, meta)
|
|
finally:
|
|
for name in files: # на устройстве файлы не остаются ни при успехе, ни при сбое
|
|
try:
|
|
await ros.remove_file(c, name)
|
|
except RosError:
|
|
pass
|
|
status, error = "done", None
|
|
except Exception as e:
|
|
error = str(e)
|
|
raise
|
|
finally:
|
|
backups.invalidate() # часть файлов могла успеть загрузиться в бакет и при неудаче
|
|
|
|
def _finish_row() -> None:
|
|
with session_scope() as s:
|
|
b = s.get(Backup, backup_id)
|
|
b.status, b.error = status, error
|
|
events.record("backup.done" if status == "done" else "backup.failed", "backup", backup_id,
|
|
"Резервная копия загружена в S3" if status == "done" else f"Резервная копия не создана: {error}",
|
|
device_id=device_id, data={"key_binary": key_bin, "key_rsc": key_rsc} if status == "done" else {"error": error}, s=s)
|
|
|
|
await asyncio.to_thread(_finish_row)
|
|
return f"Бэкап загружен в S3: {key_bin}, {key_rsc}"
|
|
|
|
|
|
async def set_channel(device_id: str, channel: str) -> None:
|
|
async with devices.open_client(await _conn(device_id)) as c:
|
|
await ros.set_channel(c, channel)
|
|
await refresh_status(device_id)
|
|
|
|
|
|
async def run_set_channel(device_id: str, channel: str) -> str:
|
|
"""Раннер задачи set_channel: та же смена канала, что и для одного устройства."""
|
|
await set_channel(device_id, channel)
|
|
return f"Канал {channel} установлен"
|
|
|
|
|
|
async def run_ros_update(device_id: str) -> str:
|
|
async with devices.open_client(await _conn(device_id)) as c:
|
|
return await ros.install_ros_update(c)
|
|
|
|
|
|
async def run_fw_update(device_id: str) -> str:
|
|
async with devices.open_client(await _conn(device_id)) as c:
|
|
return await ros.upgrade_firmware(c)
|
|
|
|
|
|
async def run_ros_downgrade(device_id: str, target_version: str) -> str:
|
|
"""Осознанный откат ROS до версии канала: 1) проверка, что версия канала совпадает с подтверждённой
|
|
и действительно старше установленной; 2) обязательный бэкап — сбой прерывает задачу, откат не
|
|
запускается (бэкап привязывается к этой же задаче через job_id, ContextVar events.current_job);
|
|
3) повторная проверка и запуск отката на новом подключении."""
|
|
async with devices.open_client(await _conn(device_id)) as c:
|
|
await ros.check_downgrade(c, target_version)
|
|
backup_msg = await run_backup(device_id)
|
|
async with devices.open_client(await _conn(device_id)) as c:
|
|
install_msg = await ros.downgrade_ros(c, target_version)
|
|
return f"{install_msg}. {backup_msg}"
|