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

"""Запуск и остановка потокового графа GNU Radio отдельным процессом."""

import signal
import subprocess
from threading import Thread

from core._version import (
    flowgraphs_git_hash,
    gr_satellites_git_hash,
    gr_satnogs_git_hash,
)
from core.configs import logger, settings
from core.configs.flowgraph import MODES
from core.types import Metadata
from soniks_client.sat_data import applied_version
from soniks_client.subprocess_log import start_logged_process, stop_logging


[документация] class Flowgraph: """Потоковый граф GNU Radio, запускаемый отдельным процессом. Поля ``settings.flowgraph`` вместе с параметрами прохода превращаются в аргументы ``--kebab-case=значение``; значения ``None`` не передаются. Вывод процесса читается отдельным потоком и уходит в лог. """ def __init__( self, rx_device: str, frequency: int, mode: str, baud: int, payload_ogg_path: str, waterfall_raw_path: str, prefix_of_decoded_data_file: str, norad_cat_id: int | None = None, ) -> None: self.parameters = { "mode": mode, "soapy-rx-device": rx_device, "samp-rate-rx": settings.flowgraph.RX_SAMP_RATE, "rx-freq": frequency, "file-path": payload_ogg_path, "waterfall-file-path": waterfall_raw_path, "decoded-data-file-path": prefix_of_decoded_data_file, "doppler-correction-per-sec": settings.flowgraph.DOPPLER_CORR_PER_SEC, "lo-offset": settings.flowgraph.LO_OFFSET, "lo-transverter": settings.flowgraph.LO_TRANSVERTER, "ppm": settings.flowgraph.PPM_ERROR, "rigctl-host": settings.antenna.rig.IP, "rigctl-port": settings.antenna.rig.PORT, "gain-mode": settings.flowgraph.GAIN_MODE, "gain": settings.flowgraph.RF_GAIN, "antenna": settings.flowgraph.ANTENNA, "dev-args": settings.flowgraph.DEV_ARGS, "stream-args": settings.flowgraph.STREAM_ARGS, "tune-args": settings.flowgraph.TUNE_ARGS, "other-settings": settings.flowgraph.OTHER_SETTINGS, # Диспетчер объявляет оба bool как ``type=int``: ``int("True")`` # роняет его с кодом 2 до старта графа. "dc-removal": ( None if settings.flowgraph.DC_REMOVAL is None else int(settings.flowgraph.DC_REMOVAL) ), "bb-freq": settings.flowgraph.BB_FREQ, "bw": settings.flowgraph.RX_BANDWIDTH, "enable-iq-dump": int(settings.flowgraph.ENABLE_IQ_DUMP), "iq-file-path": settings.flowgraph.IQ_DUMP_FILENAME, "udp-dump-host": settings.flowgraph.UDP_DUMP_HOST, "udp-dump-port": settings.flowgraph.UDP_DUMP_PORT, } if baud: self.parameters["baud"] = baud if norad_cat_id is not None: self.parameters["norad-cat-id"] = norad_cat_id self.dispatcher_script = settings.flowgraph.FLOWGRAPH_DISPATCHER self.process: subprocess.Popen | None = None self.log_thread: Thread | None = None # Код возврата, если граф кончился сам, до нашей просьбы. ``None`` — # штатный путь: остановлен нами либо не запускался. Отдельное поле, а не # ``process.returncode``: на штатном пути там всегда -2 от SIGINT, и # отличить по нему упавший граф от успешного прохода нельзя. self.exit_code: int | None = None @property def is_running(self) -> bool: """Живёт ли процесс графа.""" return bool(self.process and self.process.poll() is None)
[документация] def start(self) -> None: """Запустить граф и поток чтения его вывода. Raises: OSError: Процесс не удалось запустить. Раньше ошибка только логировалась, ``self.process`` оставался ``None``, и вызывающий продолжал проход как ни в чём не бывало. """ if self.is_running: logger.info("Flowgraph уже запущен") return args = [self.dispatcher_script] for parameter, value in self.parameters.items(): if value is not None: args.append(f"--{parameter}={value}") try: self.process, self.log_thread = start_logged_process(args, _log_line) logger.info("Запуск Flowgraph") except OSError as e: logger.error("Ошибка при запуске Flowgraph: %s", e) raise
[документация] def stop(self) -> None: """Остановить граф. Сначала ``SIGINT`` с ожиданием 10 секунд — по нему GNU Radio успевает корректно закрыть выходные файлы. Только затем принудительное завершение. Граф, кончившийся сам до этого вызова, здесь же и опознаётся: отдельного состояния для этого не нужно, ``stop()`` вызывается из ``finally`` на любом пути прохода. Раньше такой отказ молчал, и проход без единого файла выглядел успешным наблюдением. """ if not self.is_running: if self.process is None: logger.info("Flowgraph не запущен") return self.exit_code = self.process.returncode # Сначала слить остаток вывода: причина смерти написана самим # диспетчером («unrecognized arguments: …»), и в логе она должна # стоять перед строкой с кодом, а не после неё. stop_logging(self.log_thread, "Flowgraph") # 2 — argparse диспетчера, то есть расхождение контракта клиент ↔ # диспетчер: граф не стартовал вовсе. Отрицательное значение или # 128+N — сигнал. Прочее — падение самого графа. logger.error( "Flowgraph завершился сам, код возврата: %s", self.exit_code, ) return try: self.process.send_signal(signal.SIGINT) self.process.wait(timeout=10) logger.info("Flowgraph остановлен") except subprocess.TimeoutExpired: logger.warning("Принудительное завершение Flowgraph") self.process.kill() self.process.wait() except Exception as e: logger.exception("Ошибка при остановке Flowgraph: %s", e) finally: # На пути «SIGINT не помог» поток логирования раньше не # join-ился, а stdout не закрывался ни на одном из путей. stop_logging(self.log_thread, "Flowgraph")
[документация] def get_metadata(self) -> Metadata: """Вернуть блок метаданных приёмного тракта для отправки на портал. Включает полный набор параметров графа, поэтому эти метаданные становятся публичными вместе с наблюдением. """ return { "radio": { "name": "gr-satnogs", "version": gr_satnogs_git_hash[:8], # Тег образа не говорит, какой код флоуграфов внутри: теги # переставляют. Коммит — говорит, и постфактум это единственный # способ узнать, чем принято конкретное наблюдение. "flowgraphs": flowgraphs_git_hash[:8], # Второй декодер на тот же сигнал — тем же приёмом. "gr_satellites": gr_satellites_git_hash[:8], # Тем же приёмом — определения спутников: они больше не едут в # образе намертво, а приезжают bundle'ом с портала. ``None`` # означает станцию на данных из образа. "sat_data": applied_version(), # Без этого упавший граф виден только в логе станции, до # которого никто не доходит. Здесь он доезжает до портала: # ``None`` — проход прошёл штатно, число — граф кончился сам. "exit_code": self.exit_code, # Тем же приёмом видимым делается расхождение таблиц режимов. # Сырой режим лежит ниже в ``parameters``, но судить по нему # нельзя, не зная версию таблицы клиента: ``False`` означает # ровно то, что скриптам прохода ушло имя графа чужого режима. "mode_known": self.parameters["mode"] in MODES, "parameters": self.parameters, } }
def _log_line(line: str) -> None: logger.log(settings.log.flowgraph_log_level, "[Flowgraph] %s", line)