"""Оркестрация операций: 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 devices def _save_status(device_id: int, fields: dict) -> None: with session_scope() as s: d = s.get(Device, device_id) if d is not None: # устройство могли удалить во время опроса for k, v in fields.items(): setattr(d, k, v) d.status_at = now() async def refresh_status(device_id: int) -> None: """Полный опрос: статус, версии и проверка обновлений. Недоступность — не исключение.""" conn = devices.get_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)) _save_status(device_id, fields) async def poll_device(device_id: int, full: bool = False) -> None: """Фоновый опрос. Лёгкий (один запрос system/resource) обновляет online, uptime и версию ROS; полный — если нужна проверка обновлений (full) либо устройство было недоступно/ещё не опрошено (после возвращения версии могли измениться).""" d = devices.get_device(device_id) if full or d.online is not True or not d.status: return await refresh_status(device_id) try: async with devices.open_client(devices.get_conn(device_id)) as c: res = await c.get("system/resource") except RosError as e: return _save_status(device_id, dict(online=False, last_error=str(e))) status = {**d.status, "uptime": res.get("uptime"), "ros_version": res.get("version")} _save_status(device_id, dict(online=True, status_json=json.dumps(status), last_error=None)) async def refresh_many(device_ids: list[int] | None = None) -> None: ids = device_ids if device_ids is not None else [d.id for d in devices.list_devices()] sem = asyncio.Semaphore(get_settings().max_concurrency) async def one(i: int): async with sem: await refresh_status(i) await asyncio.gather(*(one(i) for i in ids)) async def run_backup(device_id: int) -> str: """Бэкап .backup + .rsc: создать на устройстве → скачать через REST во временную папку → загрузить в S3 от имени сервера → удалить файлы с устройства.""" conn = devices.get_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"] with session_scope() as s: b = Backup(device_id=device_id, 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() backup_id = b.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) 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: with session_scope() as s: b = s.get(Backup, backup_id) b.status, b.error = status, error return f"Бэкап загружен в S3: {key_bin}, {key_rsc}" async def set_channel(device_id: int, channel: str) -> None: async with devices.open_client(devices.get_conn(device_id)) as c: await ros.set_channel(c, channel) await refresh_status(device_id) async def run_ros_update(device_id: int) -> str: async with devices.open_client(devices.get_conn(device_id)) as c: return await ros.install_ros_update(c) async def run_fw_update(device_id: int) -> str: async with devices.open_client(devices.get_conn(device_id)) as c: return await ros.upgrade_firmware(c)