"""Конфигурация станции с портала: получение, проверка, применение.
Портал хранит желаемую конфигурацию станции (форма «Настройки станции» на
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]