"""Задание наблюдения: полный жизненный цикл одного прохода."""
from datetime import UTC, datetime
from apscheduler.schedulers.base import BaseScheduler
from watchdog.observers import Observer
from core.configs import logger, settings
from core.configs.flowgraph import MODES
from core.exceptions import NoCompatibleRxDeviceError
from core.types import TLE
from soniks_client.antenna.rig import get_rig_controller
from soniks_client.antenna.rotator import get_rotator_controller
from soniks_client.api import modes_published
from soniks_client.health import observation_finished, observation_started
from soniks_client.jobs._log import log_next_observation
from soniks_client.jobs.file_surveillance import CreateFileHandler
from soniks_client.jobs.sending_data import (
send_data_after_observation,
send_data_during_observation,
)
from soniks_client.observation import Observation
from soniks_client.sat_data import baudrate_for, iq_mode_for, norad_of
[документация]
def execute_observation(
*,
job_id: str,
tle: TLE,
start: datetime,
end: datetime,
frequency: int,
mode: str,
baud: int,
norad_cat_id: int | None = None,
scheduler: BaseScheduler,
) -> None:
"""Провести один проход целиком.
Последовательность: подготовка, запуск наблюдателя файловой системы
(кадры уходят на портал во время прохода), ``satnogs-pre``, приём,
``satnogs-post``, постобработка и постановка задания на выгрузку.
Ошибки наблюдения логируются и гасятся — планировщик продолжает работу.
"""
# Неизвестный режим отказывается, только когда портал знает режимы станции:
# до публикации он считает, что станция умеет всё, и отказ снял бы проходы,
# которые она принимает. После — подмена на DEFAULT_MODE отдала бы скриптам
# частоту дискретизации чужого графа (канал C в docs/roadmap-network.md).
mode_requested = None
if mode not in MODES and modes_published():
# Спутник с satyaml декодирует gr-satellites, которому от satnogs-графа
# нужен только IQ: граф берётся по модуляции передатчика (Фаза 3.5).
norad = norad_of(tle, norad_cat_id)
fallback = iq_mode_for(norad)
if fallback is None:
logger.error(
"Режима %s нет в таблице клиента, а портал знает режимы станции:"
" проход не начинается [ID: %s]",
mode,
job_id,
)
return
logger.info(
"Режима %s нет в таблице клиента, спутник %s есть в satyaml:"
" IQ для gr-satellites даст граф режима %s [ID: %s]",
mode,
norad,
fallback,
job_id,
)
mode_requested, mode = mode, fallback
baud = baud or baudrate_for(norad)
# Координаты необязательны в .env и приходят с портала; без них расчёт
# положения спутника невозможен.
if settings.station.LATITUDE is None or settings.station.LONGITUDE is None:
logger.error(
"Координаты станции неизвестны: проход не начинается [ID: %s]", job_id
)
return
observation = Observation(get_rig_controller(), get_rotator_controller())
try:
observation.set_observation_parameters(
observation_id=job_id,
tle=tle,
observation_start=start,
observation_end=end,
frequency=frequency,
mode=mode,
baud=baud,
norad_cat_id=norad_cat_id,
mode_requested=mode_requested,
)
except NoCompatibleRxDeviceError:
return
file_observer = Observer()
file_handler = CreateFileHandler(
observation_id=job_id,
prefix=settings.observation.prefix.DATA,
send_callback=send_data_during_observation,
scheduler=scheduler,
batch_delay=settings.observation.BATCH_DELAY,
)
file_observer.schedule(
file_handler,
observation.observation_directory,
recursive=False,
)
file_observer.start()
# Признак для /healthz — не лок: наложившиеся проходы разрешены. Отметка
# именно здесь, а не в начале функции: выше есть ранний return по
# несовместимому приёмнику, и признак залипал бы до перезапуска процесса.
observation_started(job_id)
# Наблюдатель уже запущен, поэтому дальше — только гасимые ошибки: любой
# выход мимо finally оставил бы висеть поток наблюдателя, а данные прохода —
# в output/, который не обходит ни одно задание (повторы читают incomplete/).
try:
try:
observation.run_pre_script()
except Exception as e:
# Проход продолжается — так же, как при падении post-скрипта.
logger.error("Ошибка pre-скрипта наблюдения [ID: %s]: %s", job_id, e)
try:
observation.run()
except Exception as e:
logger.error("Ошибка во время наблюдения [ID: %s]: %s", job_id, e)
finally:
file_observer.stop()
file_observer.join()
# Пачка, накопленная к концу прохода: раньше её досылал __del__ во время
# сборки мусора. Досылка синхронная и до постановки задания выгрузки —
# поэтому один кадр не может уйти на портал дважды.
try:
file_handler.send_batch(force=True)
except Exception as e:
logger.error(
"Ошибка отправки последней пачки кадров [ID: %s]: %s", job_id, e
)
# Задание выгрузки должно быть поставлено даже если постобработка упала:
# иначе данные прохода не уйдут никогда.
try:
observation.run_post_script()
observation.post_processing()
except Exception as e:
logger.error("Ошибка постобработки наблюдения [ID: %s]: %s", job_id, e)
finally:
# Первой строкой: ниже в этом же finally стоит add_job, который на
# умирающем планировщике бросает, и признак остался бы поднятым.
observation_finished(job_id)
if file_observer.is_alive():
file_observer.stop()
file_observer.join()
scheduler.add_job(
func=send_data_after_observation,
trigger="date",
next_run_time=datetime.now(UTC),
id=f"send_{job_id}",
replace_existing=True,
args=[job_id],
)
logger.debug("Добавлено задание отправки данных наблюдения: %s", job_id)
log_next_observation(scheduler)