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

"""Обследование приёмника: какие устройства подключены и какое усиление ставить.

Две операции, обе — python-биндингом SoapySDR в подпроцессе с таймаутом, как
``sdr_check``, и только когда приёмник свободен: во время прохода он занят
графом, а Airspy и RTL-SDR второй раз не открываются.

* **Поиск устройств** — ``SoapySDR.Device.enumerate()`` плюс возможности
  каждого найденного (входы, диапазоны усиления, частоты дискретизации). На
  старте, раз в ``SDR__SCAN_INTERVAL_IN_MINUTES`` в простое и по команде с
  портала. Портал показывает найденное в форме настроек и подставляет строку
  устройства, вход и границы усиления одной кнопкой.
* **Калибровка усиления** — по команде с портала: свип по общему усилению на
  заданных частотах, средняя мощность на каждом шаге. Шумовая полка растёт с
  усилением, пока внешний шум не перекроет собственный; рекомендуется первое
  усиление, на котором полка поднялась на ``SDR__NOISE_LIFT_DB`` над минимумом
  (Фаза 4, шаг 1 в docs/roadmap-network.md).

Команды приходят блоком ``commands`` документа состояния
(``remote_config``): ``{"rescan_sdr": {"at": ...}, "calibrate_gain": {"at":
..., "frequencies": [...]}}``. Команда выполняется, если её отметка ``at``
отличается от последней выполненной; отметки и последние результаты лежат
рядом с файлом состояния. Итоги уходят ключами ``sdr`` и ``calibration``
документа статуса (``api.status_body()``).
"""

import json
import math
import subprocess
import sys
import threading
import time
from datetime import UTC, datetime
from pathlib import Path

NO_BINDING = 3
SCAN_TIMEOUT_SECONDS = 60
CALIBRATION_TIMEOUT_SECONDS = 180
# Пауза после смены усиления, чтобы АРУ-цепочки и фильтры драйвера устоялись.
SETTLE_SECONDS = 0.05

_lock = threading.Lock()
_worker: threading.Thread | None = None
_requested: dict = {}
_state: dict | None = None
_scanned_this_run = False


[документация] def recommend(points: list[list[float]], lift_db: float) -> float | None: """Выбрать усиление по подъёму шумовой полки. Args: points: Пары ``[усиление, мощность_дБ]`` по возрастанию усиления. lift_db: На сколько полка должна подняться над минимальной. Returns: Первое усиление, на котором подъём достигнут; если не достигнут нигде — усиление с максимальной мощностью (приёмник упёрся раньше); ``None`` без точек. """ if not points: return None floor = min(power for _, power in points) for gain, power in points: if power >= floor + lift_db: return gain return max(points, key=lambda point: point[1])[0]
[документация] def report() -> dict: """Ключи ``sdr`` и ``calibration`` для документа статуса — только известные.""" state = _load() return {key: state[key] for key in ("sdr", "calibration") if state.get(key)}
[документация] def request(commands) -> None: """Запомнить команды из документа состояния; выполнит ``run_due()``.""" global _requested if isinstance(commands, dict): _requested = commands
[документация] def run_due() -> None: """Запустить назревшее обследование в фоновом потоке, если приёмник свободен. Зовётся после каждой сверки и после каждого документа по MQTT. Один поток на всё: два обследования разом не откроют один приёмник. """ global _worker with _lock: if _worker is not None and _worker.is_alive(): return window = _idle_window() if window is not None and window < 0: return tasks = _due_tasks() if not tasks: return _worker = threading.Thread(target=_work, args=(tasks,), name="sdr-survey", daemon=True) _worker.start()
[документация] def scan() -> dict: """Перечислить приёмники и их возможности; итог — ключ ``sdr`` статуса.""" from core.configs import logger global _scanned_this_run result, error = _run({"op": "scan"}, SCAN_TIMEOUT_SECONDS) doc = { "scanned_at": _now(), "devices": result.get("devices", []) if result else [], "error": error, } _scanned_this_run = True _save({"sdr": doc}) if error: logger.warning("Поиск приёмников не удался: %s", error) else: logger.info( "Найдено приёмников: %d (%s)", len(doc["devices"]), ", ".join(device["args"] for device in doc["devices"]) or "ни одного", ) return doc
[документация] def calibrate(frequencies: list) -> dict: """Свип по усилению на частотах; итог — ключ ``calibration`` статуса. Приёмник для каждой частоты — тот, что взял бы проход (``OBSERVATION__SOAPY_RX_DEVICE`` может быть списком диапазонов), поэтому частоты группируются по устройству и каждое обследуется одним подпроцессом. """ from core.configs import logger, settings from core.exceptions import NoCompatibleRxDeviceError from soniks_client.rx_device import select_rx_device_by_frequency by_device: dict[str, list[float]] = {} results: list[dict] = [] errors: list[str] = [] for frequency in frequencies: try: frequency = float(frequency) device = select_rx_device_by_frequency( settings.observation.SOAPY_RX_DEVICE, int(frequency) ) except (TypeError, ValueError, NoCompatibleRxDeviceError): errors.append(f"нет приёмника для {frequency}") continue by_device.setdefault(device, []).append(frequency) flowgraph = settings.flowgraph try: sample_rate = float(flowgraph.RX_SAMP_RATE) if flowgraph.RX_SAMP_RATE else None except ValueError: sample_rate = None for device, device_frequencies in by_device.items(): result, error = _run( { "op": "calibrate", "device": device, "sample_rate": sample_rate, "antenna": flowgraph.ANTENNA or None, "frequencies": device_frequencies, "step": settings.sdr.GAIN_STEP_DB, "lift": settings.sdr.NOISE_LIFT_DB, "samples": settings.sdr.CALIBRATION_SAMPLES, }, CALIBRATION_TIMEOUT_SECONDS, ) if error: errors.append(f"{device}: {error}") continue for entry in result.get("results", []): results.append({"device": device, **entry}) recommended = [entry["recommended"] for entry in results if entry["recommended"] is not None] doc = { "calibrated_at": _now(), "results": results, # Одно усиление на все диапазоны — наибольшее из рекомендованных: тогда # подъём полки достигнут на каждом из них. "recommended": max(recommended) if recommended else None, "lift_db": settings.sdr.NOISE_LIFT_DB, "error": "; ".join(errors) or None, } _save({"calibration": doc}) if errors: logger.warning("Калибровка усиления: %s", doc["error"]) if results: logger.info( "Калибровка усиления: рекомендовано %s дБ по %d частотам", doc["recommended"], len(results), ) return doc
def _due_tasks() -> list: from core.configs import settings state = _load() done = state["done"] tasks = [] rescan = _requested.get("rescan_sdr") calibration = _requested.get("calibrate_gain") scan_due = not _scanned_this_run or _marker(rescan) != done.get("rescan_sdr") last_scan = (state.get("sdr") or {}).get("scanned_at") if not scan_due and last_scan: age = (datetime.now(UTC) - datetime.fromisoformat(last_scan)).total_seconds() scan_due = age >= settings.sdr.SCAN_INTERVAL_IN_MINUTES * 60 if scan_due: tasks.append(("rescan_sdr", _marker(rescan), scan)) if isinstance(calibration, dict) and _marker(calibration) != done.get("calibrate_gain"): frequencies = calibration.get("frequencies") or [] tasks.append(("calibrate_gain", _marker(calibration), lambda: calibrate(frequencies))) return tasks def _work(tasks: list) -> None: from core.configs import logger for name, marker, task in tasks: try: task() except Exception as e: logger.exception("Обследование приёмника (%s) не удалось: %s", name, e) _save({"done": {**_load()["done"], name: marker}}) # Ленивый импорт: mqtt и remote_config импортируют этот модуль. from soniks_client import api, mqtt api.publish_station_status() mqtt.publish_status() def _idle_window() -> float | None: """Секунды до ближайшего прохода за вычетом запаса; отрицательно — нельзя.""" from core.configs import settings from soniks_client import health seconds = health.seconds_to_next_observation() if seconds is None: return None return seconds - settings.sdr.IDLE_WINDOW_IN_MINUTES * 60 def _marker(command) -> str | None: if isinstance(command, dict): return str(command.get("at") or "") or None return None def _now() -> str: return datetime.now(UTC).isoformat() def _path() -> Path: from core.configs import settings return Path(settings.paths.STATE_FILE).with_name("sdr-survey.json") def _load() -> dict: global _state if _state is None: _state = {"done": {}, "sdr": None, "calibration": None} try: stored = json.loads(_path().read_text()) if isinstance(stored, dict): _state.update(stored) _state["done"] = dict(_state.get("done") or {}) except FileNotFoundError: pass except (OSError, ValueError) as e: from core.configs import logger logger.warning("Файл обследования %s не прочитан: %s", _path(), e) return _state def _save(update: dict) -> None: from core.configs import logger state = _load() state.update(update) path = _path() tmp = path.with_suffix(".tmp") try: path.parent.mkdir(parents=True, exist_ok=True) tmp.write_text(json.dumps(state, ensure_ascii=False)) tmp.replace(path) except OSError as e: logger.warning("Файл обследования %s не записан: %s", path, e) def _run(args: dict, timeout: float) -> tuple[dict | None, str | None]: """Выполнить операцию в дочернем процессе; ``(результат, ошибка)``.""" try: proc = subprocess.run( [sys.executable, __file__, json.dumps(args)], capture_output=True, text=True, timeout=timeout, ) except subprocess.TimeoutExpired: return None, f"приёмник не ответил за {timeout:g} с" except OSError as e: return None, f"обследование не запустилось: {e}" if proc.returncode == NO_BINDING: return None, "биндинг SoapySDR недоступен" if proc.returncode != 0: lines = proc.stderr.strip().splitlines() return None, lines[-1] if lines else f"код {proc.returncode}" try: return json.loads(proc.stdout), None except ValueError: return None, "ответ обследования не разобран" # --- Дочерний процесс: отсюда и ниже — только SoapySDR, без настроек станции. def _device_args(kwargs: dict) -> str: """Строка для ``OBSERVATION__SOAPY_RX_DEVICE``: драйвер и серийный номер.""" args = f"driver={kwargs.get('driver', '')}" if kwargs.get("serial"): args += f",serial={kwargs['serial']}" return args def _range(value) -> list[float]: return [value.minimum(), value.maximum()] def _scan_devices(SoapySDR) -> dict: rx = SoapySDR.SOAPY_SDR_RX devices = [] for found in SoapySDR.Device.enumerate(): kwargs = dict(found) entry = { "args": _device_args(kwargs), "driver": kwargs.get("driver", ""), "label": kwargs.get("label", ""), "serial": kwargs.get("serial", ""), } try: device = SoapySDR.Device(kwargs) ranges = device.getFrequencyRange(rx, 0) entry.update( antennas=list(device.listAntennas(rx, 0)), gains={ "overall": _range(device.getGainRange(rx, 0)), "stages": { name: _range(device.getGainRange(rx, 0, name)) for name in device.listGains(rx, 0) }, }, sample_rates=[float(rate) for rate in device.listSampleRates(rx, 0)], frequency_range=[ min(r.minimum() for r in ranges), max(r.maximum() for r in ranges) ] if ranges else None, dc_offset_mode=bool(device.hasDCOffsetMode(rx, 0)), ) SoapySDR.Device.unmake(device) except Exception as e: entry["error"] = str(e) devices.append(entry) return {"devices": devices} def _read_power_db(SoapySDR, device, stream, buffer) -> float: import numpy # Первое чтение после смены усиления — отбросить: в нём хвост старого. for _ in range(2): result = device.readStream(stream, [buffer], len(buffer), timeoutUs=500_000) if result.ret < 0: raise RuntimeError(SoapySDR.errToStr(result.ret)) samples = buffer[: result.ret] return 10 * math.log10(float(numpy.mean(numpy.abs(samples) ** 2)) + 1e-20) def _calibrate_device(SoapySDR, args: dict) -> dict: import numpy rx = SoapySDR.SOAPY_SDR_RX device = SoapySDR.Device(args["device"]) if args.get("sample_rate"): device.setSampleRate(rx, 0, float(args["sample_rate"])) if args.get("antenna"): device.setAntenna(rx, 0, args["antenna"]) low, high = _range(device.getGainRange(rx, 0)) step = float(args["step"]) gains = [low + i * step for i in range(int((high - low) / step) + 1)] if gains[-1] < high: gains.append(high) buffer = numpy.empty(int(args["samples"]), numpy.complex64) stream = device.setupStream(rx, SoapySDR.SOAPY_SDR_CF32) device.activateStream(stream) results = [] try: for frequency in args["frequencies"]: device.setFrequency(rx, 0, float(frequency)) points = [] for gain in gains: device.setGain(rx, 0, gain) time.sleep(SETTLE_SECONDS) points.append([gain, _read_power_db(SoapySDR, device, stream, buffer)]) results.append({ "frequency": float(frequency), "points": points, "recommended": recommend(points, float(args["lift"])), }) finally: device.deactivateStream(stream) device.closeStream(stream) SoapySDR.Device.unmake(device) return {"results": results} def _main(args: dict) -> None: try: import SoapySDR except ImportError: sys.exit(NO_BINDING) try: if args["op"] == "scan": result = _scan_devices(SoapySDR) else: result = _calibrate_device(SoapySDR, args) except Exception as e: print(e, file=sys.stderr) sys.exit(1) print(json.dumps(result)) if __name__ == "__main__": _main(json.loads(sys.argv[1]))