"""Запуск и остановка потокового графа 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 _log_line(line: str) -> None:
logger.log(settings.log.flowgraph_log_level, "[Flowgraph] %s", line)