Исходный код soniks_client.remote_config

"""Конфигурация станции с портала: получение, проверка, применение.

Портал хранит желаемую конфигурацию станции (форма «Настройки станции» на
sonik.space) и отдаёт её документом ``GET /api/v2/stations/<id>/state/``:
``generation`` — номер сохранения, ``config`` — плоский словарь, ключи которого
— имена переменных окружения (``FLOWGRAPH__RF_GAIN``), та же лексика, что в
docs/station/environment_variables.md; ``location`` — координаты станции;
``release`` — канал обновлений. Тот же документ приходит по MQTT
(``soniks_client.mqtt``) в момент сохранения формы.

Применение — в процессе, без перезапуска контейнера: разделы ``settings``
пересобираются через pydantic поверх значений из ``.env``, значения зеркалятся
в ``os.environ`` для скриптов из ``scripts/``, кэш контроллеров ротатора и рига
сбрасывается. Пока владелец ничего не сохранил на портале (``generation`` 0),
станция работает по ``.env`` как раньше.

Последний полученный документ лежит на диске (``PATHS__STATE_FILE``): станция
без портала стартует на нём.
"""

import json
import logging
import os
import re
import threading
from dataclasses import dataclass
from datetime import UTC, datetime
from enum import Enum
from pathlib import Path
from urllib.parse import urljoin

import requests
from pydantic import BaseModel, ValidationError

from core.configs import logger, settings
from core.configs.antenna import RigSettings, RotatorSettings
from core.configs.flowgraph import FlowgraphSettings
from core.configs.log import LoggingSettings
from core.configs.observation import ObservationSettings
from soniks_client import api, health, sdr_check, sdr_survey
from soniks_client.antenna.rig import get_rig_controller
from soniks_client.antenna.rotator import get_rotator_controller

# Что портал вправе менять (решение 48 в docs/roadmap-network.md): тракт
# приёма, поворотное устройство, риг и уровни логов. Всё остальное — адреса,
# пути, планировщик, таймауты — остаётся в ``.env``. Тот же список лежит в
# ``STATION_CONFIG_FIELDS`` портала; ключ вне списка — отказ всего документа.
_SECTIONS: dict[str, tuple[type[BaseModel], tuple[str, ...], tuple[str, ...]]] = {
    "FLOWGRAPH": (
        FlowgraphSettings,
        ("flowgraph",),
        (
            "RX_SAMP_RATE", "RX_BANDWIDTH", "DOPPLER_CORR_PER_SEC", "LO_OFFSET",
            "LO_TRANSVERTER", "PPM_ERROR", "GAIN_MODE", "RF_GAIN", "ANTENNA",
            "DEV_ARGS", "STREAM_ARGS", "TUNE_ARGS", "OTHER_SETTINGS", "DC_REMOVAL",
            "BB_FREQ", "ENABLE_IQ_DUMP", "IQ_DUMP_FILENAME",
        ),
    ),
    "OBSERVATION": (
        ObservationSettings,
        ("observation",),
        (
            "SOAPY_RX_DEVICE", "REMOVE_WATERFALL_RAW_FILES", "REMOVE_OBSERVATION_DATA",
            "RUN_PRE_OBSERVATION_SCRIPT", "RUN_POST_OBSERVATION_SCRIPT",
        ),
    ),
    "ANTENNA__ROTATOR": (
        RotatorSettings,
        ("antenna", "rotator"),
        ("ENABLED", "MODEL", "BAUD", "PORT", "THRESHOLD", "MODE", "MIN_ELEVATION"),
    ),
    "ANTENNA__RIG": (RigSettings, ("antenna", "rig"), ("IP", "PORT")),
    "LOG": (LoggingSettings, ("log",), ("LEVEL", "SCRIPT_LEVEL", "FLOWGRAPH_LEVEL")),
}

MANAGED_FIELDS: frozenset[str] = frozenset(
    f"{prefix}__{name}" for prefix, (_, _, names) in _SECTIONS.items() for name in names
)

# Смена любого из них означает другое устройство или другой режим его работы —
# перед применением приёмник открывается проверкой.
_SDR_FIELDS = frozenset({
    "OBSERVATION__SOAPY_RX_DEVICE", "FLOWGRAPH__RX_SAMP_RATE", "FLOWGRAPH__ANTENNA",
    "FLOWGRAPH__GAIN_MODE", "FLOWGRAPH__RF_GAIN", "FLOWGRAPH__DEV_ARGS",
})

_LOG_LEVELS = ("DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL")


[документация] @dataclass class ApplyResult: """Итог применения документа состояния.""" generation: int ok: bool error: str = "" applied_at: str | None = None
def _section(path: tuple[str, ...]) -> BaseModel: model = settings for name in path: model = getattr(model, name) return model def _set_section(path: tuple[str, ...], model: BaseModel) -> None: setattr(_section(path[:-1]) if len(path) > 1 else settings, path[-1], model) # Снимки того, с чем станция стартовала: разделы настроек из ``.env`` и их # окружение. Слияние всегда идёт поверх них, а не поверх текущих значений — # иначе ключ, который владелец убрал с портала, оставался бы от прошлого # поколения. _BASELINE: dict[str, dict] = { prefix: _section(path).model_dump() for prefix, (_, path, _) in _SECTIONS.items() } _ENV_BASELINE: dict[str, str | None] = { field: os.environ.get(field) for field in MANAGED_FIELDS } _lock = threading.Lock() _applied = ApplyResult(generation=0, ok=True) _applied_config: dict | None = None _pending: dict | None = None _etag: str | None = None
[документация] def actual_config() -> dict: """Фактическая конфигурация станции — плоский словарь по ``MANAGED_FIELDS``.""" config = {} for prefix, (_, path, names) in _SECTIONS.items(): values = _section(path).model_dump() for name in names: value = values[name] config[f"{prefix}__{name}"] = value.value if isinstance(value, Enum) else value return config
[документация] def report() -> dict: """Что станция сообщает порталу о конфигурации — ключ ``config`` статуса. ``actual`` уходит всегда, и до первого сохранения на портале тоже: форма предзаполняется им, так что станция, настроенная через ``.env``, переезжает на портал без перепечатывания значений. """ return { "generation": _applied.generation, "ok": _applied.ok, "error": _applied.error, "applied_at": _applied.applied_at, "actual": actual_config(), }
[документация] def load_state_file() -> dict | None: """Прочитать последний сохранённый документ состояния.""" global _etag path = Path(settings.paths.STATE_FILE) try: stored = json.loads(path.read_text()) doc = stored["state"] except FileNotFoundError: return None except (OSError, ValueError, KeyError, TypeError) as e: logger.warning("Файл состояния %s не прочитан: %s", path, e) return None if not isinstance(doc, dict): return None _etag = stored.get("etag") return doc
[документация] def save_state_file(doc: dict) -> None: """Записать документ состояния на диск — атомарно, через соседний файл.""" path = Path(settings.paths.STATE_FILE) tmp = path.with_suffix(".tmp") try: path.parent.mkdir(parents=True, exist_ok=True) tmp.write_text(json.dumps({"etag": _etag, "state": doc}, ensure_ascii=False)) tmp.replace(path) except OSError as e: logger.warning("Файл состояния %s не записан: %s", path, e)
[документация] def fetch_state() -> dict | None: """Забрать документ состояния с портала. Returns: Документ, либо ``None``: не изменился (304), портал недоступен или эндпоинта ещё не знает (404 — не ошибка станции). """ global _etag url = urljoin(settings.url.stations_url, f"{settings.station.ID}/state/") headers = {"Authorization": f"Token {settings.station.TOKEN}"} if _etag: headers["If-None-Match"] = _etag try: response = api._session.get( url, headers=headers, timeout=settings.api.TIMEOUT_IN_SECONDS_REQUEST_DATA ) if response.status_code == 304: return None response.raise_for_status() doc = response.json() except requests.HTTPError as e: log = logger.debug if e.response.status_code == 404 else logger.warning log("Состояние станции не получено, статус-код: %s", e.response.status_code) return None except (requests.RequestException, ValueError) as e: logger.warning("Состояние станции не получено: %s", e) return None if not isinstance(doc, dict): logger.warning("Портал вернул состояние не словарём: %s", type(doc).__name__) return None _etag = response.headers.get("ETag") return doc
[документация] def sync() -> None: """Шаг сверки: забрать новое состояние либо доприменить отложенное.""" doc = fetch_state() or _pending if doc is not None: apply_state(doc)
[документация] def bootstrap() -> None: """Старт: применить сохранённое состояние, затем свежее с портала. Порядок важен: портал может лежать, а станция обязана подняться на том, что применяла в прошлый раз. """ stored = load_state_file() if stored is not None: apply_state(stored) fresh = fetch_state() if fresh is not None: apply_state(fresh) if settings.station.LATITUDE is None: logger.error( "Координаты станции не заданы ни в .env, ни на портале: проходы не " "проводятся, пока они не появятся" )
[документация] def apply_state(doc: dict) -> ApplyResult: """Применить документ состояния. Отказ (неизвестный ключ, невалидное значение, не открылся приёмник) оставляет прежнюю конфигурацию и уезжает на портал причиной в ``report()``. Идущий проход откладывает применение до следующей сверки. """ with _lock: return _apply_locked(doc)
def _apply_locked(doc: dict) -> ApplyResult: global _applied, _applied_config, _pending _fill_location(doc.get("location")) # Команды портала (пересканировать приёмник, откалибровать усиление) не # зависят от судьбы конфигурации: их выполнит ``sdr_survey.run_due()``. sdr_survey.request(doc.get("commands")) try: generation = int(doc.get("generation") or 0) except (TypeError, ValueError): generation = 0 config = doc.get("config") or {} if generation == 0 or not isinstance(config, dict) or not config: # Владелец ничего не сохранял — станция на .env. Документ всё равно # пишется: в нём координаты, а без них станция без портала не поднимется. _pending = None save_state_file(doc) return _applied if _applied.generation == generation and _applied_config == config: return _applied if health.observation_running(): if _pending is None: logger.info("Конфигурация %d получена, применится после прохода", generation) _pending = doc return ApplyResult(generation=generation, ok=False, error="deferred") _pending = None error, sections = _build_sections(config) if error is None: error = _check_sdr(sections) if error is not None: # Железо — не документ: приёмник могли ещё не воткнуть. Следующая # сверка повторит проверку, поэтому конфиг не запоминается как # отклонённый. _applied_config = None _applied = ApplyResult(generation=generation, ok=False, error=error) logger.error("Конфигурация %d отклонена: %s", generation, error) return _applied if error is not None: _applied_config = config _applied = ApplyResult(generation=generation, ok=False, error=error) logger.error("Конфигурация %d отклонена: %s", generation, error) return _applied for prefix, model in sections.items(): _set_section(_SECTIONS[prefix][1], model) _mirror_environment(config) logging.getLogger("app").setLevel(settings.log.LEVEL) logging.getLogger("row_logs").setLevel(settings.log.SCRIPT_LEVEL) get_rotator_controller.cache_clear() get_rig_controller.cache_clear() _applied_config = config _applied = ApplyResult( generation=generation, ok=True, applied_at=datetime.now(UTC).isoformat() ) health.set_config_generation(generation) save_state_file(doc) logger.info("Применена конфигурация %d с портала: %s", generation, sorted(config)) return _applied def _fill_location(location) -> None: """Взять координаты с портала, если в ``.env`` их нет.""" if settings.station.LATITUDE is not None or not isinstance(location, dict): return try: lat, lng, alt = location["lat"], location["lng"], location["alt"] if lat is None or lng is None: return settings.station.LATITUDE = float(lat) settings.station.LONGITUDE = float(lng) settings.station.ELEVATION = int(alt or 0) except (KeyError, TypeError, ValueError) as e: logger.warning("Координаты станции с портала не разобраны: %s", e) return logger.info( "Координаты станции взяты с портала: %s, %s, %s м", settings.station.LATITUDE, settings.station.LONGITUDE, settings.station.ELEVATION, ) def _build_sections(config: dict) -> tuple[str | None, dict[str, BaseModel]]: """Собрать разделы настроек поверх ``.env`` — либо вернуть причину отказа.""" unknown = sorted(set(config) - MANAGED_FIELDS) if unknown: return f"неизвестные параметры: {', '.join(unknown)}", {} sections: dict[str, BaseModel] = {} for prefix, (model_class, _, names) in _SECTIONS.items(): overrides = { name: config[f"{prefix}__{name}"] for name in names if f"{prefix}__{name}" in config } try: sections[prefix] = model_class.model_validate({**_BASELINE[prefix], **overrides}) except ValidationError as e: first = e.errors()[0] field = "__".join(str(part) for part in first["loc"]) return f"{prefix}__{field}: {first['msg']}", {} log = sections["LOG"] for name in ("LEVEL", "SCRIPT_LEVEL", "FLOWGRAPH_LEVEL"): if getattr(log, name) not in _LOG_LEVELS: return f"LOG__{name}: допустимы {', '.join(_LOG_LEVELS)}", {} return None, sections def _check_sdr(sections: dict[str, BaseModel]) -> str | None: """Открыть приёмник с новыми параметрами, если они изменились.""" current = actual_config() new_values = { f"{prefix}__{name}": getattr(model, name) for prefix, model in sections.items() for name in _SECTIONS[prefix][2] } if all(new_values[field] == current[field] for field in _SDR_FIELDS): return None flowgraph = sections["FLOWGRAPH"] gain = flowgraph.RF_GAIN if flowgraph.GAIN_MODE == "Overall" else None try: sample_rate = float(flowgraph.RX_SAMP_RATE) if flowgraph.RX_SAMP_RATE else None except ValueError: return f"FLOWGRAPH__RX_SAMP_RATE: не число: {flowgraph.RX_SAMP_RATE!r}" for device in _devices(sections["OBSERVATION"].SOAPY_RX_DEVICE): error = sdr_check.check(device, sample_rate, flowgraph.ANTENNA, gain) if error is not None: return error return None def _devices(spec: str) -> list[str]: """Строки драйверов из ``SOAPY_RX_DEVICE`` — одиночной или списка диапазонов.""" return [re.sub(r"^\d+(\.\d+)?-\d+(\.\d+)?:", "", token) for token in spec.split()] def _mirror_environment(config: dict) -> None: """Отразить применённое в ``os.environ``: скрипты читают переменные сами.""" for field in MANAGED_FIELDS: if field in config: value = config[field] if value is None: os.environ.pop(field, None) else: os.environ[field] = str(value) elif _ENV_BASELINE[field] is None: os.environ.pop(field, None) else: os.environ[field] = _ENV_BASELINE[field]