"""Оркестрация операций: 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, 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 refresh_status(device_id: str) -> 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: str, 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[str] | None = None) -> None: targets = 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: 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 = 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, 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() backup_id = b.id events.record("backup.created", "backup", backup_id, f"Резервная копия {conn.name} запрошена", device_id=device_id, data={"key_binary": key_bin, "key_rsc": key_rsc}, s=s) 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: 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) return f"Бэкап загружен в S3: {key_bin}, {key_rsc}" async def set_channel(device_id: str, 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_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(devices.get_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(devices.get_conn(device_id)) as c: return await ros.upgrade_firmware(c)