"""Наблюдение одного прохода: подготовка, приём, постобработка."""
import json
import os
from datetime import UTC, datetime
from pathlib import Path
from time import sleep
from core.configs import logger, settings
from core.types import TLE, Metadata, StationLocation
from soniks_client.antenna.communication_session import (
RigRXDopplerCorrectedSession,
RotatorTrackingSession,
)
from soniks_client.antenna.rig import RigController
from soniks_client.antenna.rotator import RotatorController
from soniks_client.flowgraph import Flowgraph
from soniks_client.observation_files import create_observation_files
from soniks_client.observation_scripts import build_script_argv, run_script
from soniks_client.post_processing import build_waterfall
from soniks_client.rx_device import select_rx_device_by_frequency
from soniks_client.sat_data import decoder_for
[документация]
class Observation:
"""Наблюдение одного прохода спутника.
Выбирает приёмник по частоте, строит пути файлов и потоковый граф,
ведёт сессии слежения и доплеровской коррекции, а после прохода строит
водопад и отправляет метаданные.
"""
def __init__(
self,
rig_controller: RigController,
rotator_controller: RotatorController | None = None,
) -> None:
self.station_location: StationLocation = {
"latitude": settings.station.LATITUDE,
"longitude": settings.station.LONGITUDE,
"elevation": settings.station.ELEVATION,
}
self.tle: TLE | None = None
self.timestamp: str | None = None
self.observation_start: datetime | None = None
self.observation_end: datetime | None = None
self.frequency: int | None = None
self.rx_device: str | None = None
self.mode: str | None = None
self.baud: int | None = None
self.norad_cat_id: int | None = None
self.flowgraph: Flowgraph | None = None
self.observation_dir: Path | None = None
self.payload_ogg_path: str | None = None
self.waterfall_raw_path: str | None = None
self.waterfall_png_path: str | None = None
self.prefix_of_decoded_data_file: str | None = None
self.metadata_json_path: str | None = None
self.rig_session = RigRXDopplerCorrectedSession(
update_interval=settings.observation.RIG_UPDATE_INTERVAL,
rig=rig_controller,
)
self.rotator_session: RotatorTrackingSession | None = None
if rotator_controller:
self.rotator_session = RotatorTrackingSession(
update_interval=settings.observation.ROTATOR_UPDATE_INTERVAL,
rotator=rotator_controller,
)
@property
def observation_directory(self) -> str:
"""Директория файлов текущего наблюдения."""
return str(self.observation_dir)
[документация]
def set_observation_parameters(
self,
observation_id: str,
tle: TLE,
observation_start: datetime,
observation_end: datetime,
frequency: int,
mode: str,
baud: int,
norad_cat_id: int | None = None,
mode_requested: str | None = None,
) -> None:
"""Подготовить наблюдение к проходу.
Выбирает приёмник по частоте, строит пути выходных файлов и
конструирует потоковый граф.
Raises:
NoCompatibleRxDeviceError: Нет приёмника под частоту прохода.
"""
self.frequency = frequency
self.rx_device = select_rx_device_by_frequency(
settings.observation.SOAPY_RX_DEVICE,
self.frequency,
)
self.observation_id = observation_id
self.tle = tle
self.observation_start = observation_start
self.observation_end = observation_end
self.mode = mode
self.baud = baud
self.norad_cat_id = norad_cat_id
# Режим портала, подменённый на IQ-граф для gr-satellites (Фаза 3.5);
# ``None`` — граф выбран по режиму портала как есть.
self.mode_requested = mode_requested
# Один выбор декодера на проход (Фаза 3.2). NORAD от портала — опт-ин,
# старый портал его не отдаёт; в TLE он есть всегда.
norad = norad_cat_id if norad_cat_id is not None else int(tle["tle2"].split()[1])
self.decoder = decoder_for(norad)
logger.info(
"Декодер прохода %s: %s (NORAD %s)", observation_id, self.decoder, norad
)
self.timestamp = datetime.now(tz=UTC).strftime(
settings.observation.TIME_FORMAT_IN_FILENAME
)
self._create_observation_file_paths()
self.flowgraph = Flowgraph(
self.rx_device,
self.frequency,
self.mode,
self.baud,
self.payload_ogg_path,
self.waterfall_raw_path,
self.prefix_of_decoded_data_file,
self.norad_cat_id,
)
[документация]
def run_pre_script(self) -> None:
"""Выполнить скрипт, запускаемый перед проходом (``satnogs-pre``)."""
if settings.observation.RUN_PRE_OBSERVATION_SCRIPT:
self._around_observation(settings.observation.PRE_OBSERVATION_SCRIPT)
[документация]
def run(self) -> None:
"""Провести проход: слежение, доплеровская коррекция, потоковый граф.
Завершается по окончании окна наблюдения или при смерти графа.
Оборудование останавливается в любом случае: без ``finally`` любое
исключение оставляло GNU Radio жить дальше, а потоки rig и rotator —
крутить антенну в следующий проход.
"""
try:
self._start_observation()
sleep(1)
while datetime.now(tz=UTC) <= self.observation_end:
if not self.flowgraph.is_running:
break
sleep(1)
finally:
self._stop_observation()
[документация]
def run_post_script(self) -> None:
"""Выполнить скрипт, запускаемый после прохода (``satnogs-post``)."""
if settings.observation.RUN_POST_OBSERVATION_SCRIPT:
self._around_observation(settings.observation.POST_OBSERVATION_SCRIPT)
[документация]
def post_processing(self) -> None:
"""Построить водопад, собрать метаданные и положить их в очередь выгрузки.
Метаданные не отправляются отсюда: файл забирает та же выгрузка, что и
аудио с кадрами. Раньше это был единственный «выстрелил и забыл» путь —
при ошибке сети блок ``signal`` (SNR, девиация, ppm) исчезал навсегда,
очереди для него не было.
"""
signal_metadata = build_waterfall(
self.waterfall_raw_path,
self.waterfall_png_path,
self.observation_id,
)
metadata = self._get_metadata()
if signal_metadata:
metadata["signal"] = signal_metadata
Path(self.metadata_json_path).write_text(
json.dumps(metadata), encoding="utf-8"
)
def _start_observation(self) -> None:
logger.info("Старт наблюдения: %s", self.observation_id)
if self.rotator_session:
self.rotator_session.set_session_parameters(
self.station_location,
self.tle,
self.observation_start,
self.observation_end,
)
self.rotator_session.start_session()
self.rig_session.set_session_parameters(
self.station_location,
self.tle,
self.frequency,
)
self.rig_session.start_session()
self.flowgraph.start()
def _stop_observation(self) -> None:
self.flowgraph.stop()
self.rig_session.stop_session()
if self.rotator_session:
self.rotator_session.stop_session()
logger.info("Завершено наблюдение: %s", self.observation_id)
def _around_observation(self, script_type: str) -> None:
argv = build_script_argv(
script_type,
self.observation_id,
self.frequency,
self.tle,
self.timestamp,
self.baud,
self.mode,
)
# Скрипты пакет не импортируют, а позиционный контракт хуков жёсткий:
# выбор декодера едет к ``grsat.py`` переменной окружения.
run_script(
argv,
settings.observation.SCRIPT_TIMEOUT_IN_SECONDS,
{**os.environ, "SONIKS_DECODER": self.decoder},
)
def _create_observation_file_paths(self) -> None:
files = create_observation_files(
self.observation_id,
self.timestamp,
self.mode,
)
self.observation_dir = files.directory
self.payload_ogg_path = files.payload_ogg
self.waterfall_raw_path = files.waterfall_raw
self.waterfall_png_path = files.waterfall_png
self.prefix_of_decoded_data_file = files.decoded_data_prefix
self.metadata_json_path = files.metadata_json
def _get_metadata(self) -> Metadata:
metadata = self.flowgraph.get_metadata()
# Чем декодировали — постфактум по наблюдению должно быть видно. Не в
# ``Flowgraph.parameters``: те целиком уезжают в argv диспетчера.
metadata["radio"]["decoder"] = self.decoder
if self.mode_requested is not None:
metadata["radio"]["mode_requested"] = self.mode_requested
metadata.update(self.station_location)
metadata["frequency"] = self.frequency
return metadata