"""Обмен с порталом СОНИКС: расписание проходов, метаданные и выгрузка файлов.
Ошибки HTTP логируются здесь же. Наружу поднимаются только
:class:`~core.exceptions.FileNotUploadedError` и
:class:`~core.exceptions.ObservationNotFoundError` — они определяют, что
делать с файлом дальше."""
import mimetypes
from datetime import UTC, datetime
from email.utils import parsedate_to_datetime
from pathlib import Path
from urllib.parse import urljoin
import requests
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
from core._version import __version__
from core.configs import logger, settings
from core.configs.flowgraph import MODES
from core.exceptions import (
FileNotUploadedError,
FileRejectedError,
ObservationNotFoundError,
)
from core.types import DataToSend
from soniks_client.health import set_clock_skew
from soniks_client.models import JobData
# Что этот клиент понимает в ответе портала. Уезжает параметром запроса
# расписания, и портал отдаёт новое поле только попросившему.
#
# Парк не обновится целиком никогда, поэтому портал не имеет права менять то,
# что отдаёт непопросившей станции: она получает ответ байт-в-байт как раньше и
# сломаться не может. Обновлённая получает больше — решение 21 в
# docs/roadmap-network.md.
#
# Заявление честно по построению: клиент и flowgraph_dispatcher едут в одном
# образе, значит объявленное клиентом умеет и диспетчер. ``norad_cat_id``
# опасен именно этим — станция с июльским образом шлёт --norad-cat-id
# диспетчеру, который его не знает, и проход умирает молча.
CAPABILITIES = ("norad_cat_id",)
# Портал отвергает кадр, метка которого дальше пяти минут от окна наблюдения
# (400 frame_outside_window), а метку ставит диспетчер по часам станции.
# Предупреждать надо заметно раньше, чем кадры начнут отвергаться.
CLOCK_SKEW_WARNING_SECONDS = 60
_retry = Retry(
# Четыре попытки, паузы 0/2/4 с. Больше незачем: невыгруженный файл не
# теряется, он ложится в incomplete/ и повторяется раз в
# SCHEDULER__RESENDING_INTERVAL_IN_MINUTES — вот это и есть долгий backoff.
total=3,
backoff_factor=1,
status_forcelist=(429, 500, 502, 503, 504),
# Только PUT. Расписание и так опрашивается раз в минуту, а четыре попытки
# по 45 с съели бы MISFIRE_GRACE_TIME и раздули last_sync_age в /healthz.
# 404 и 403 в forcelist отсутствуют намеренно — это осмысленные ответы.
allowed_methods=frozenset({"PUT"}),
# Иначе исчерпанный ретрай приходит как RetryError (это RequestException,
# не HTTPError), и ветки ниже теряют статус-код и тело ответа.
raise_on_status=False,
# urllib3 спит по Retry-After без верхней границы, а троттлинг портала
# отдаёт там всё окно лимита — поток пула завис бы на тысячи секунд.
respect_retry_after_header=False,
)
# Одна сессия на процесс: до неё на каждый файл открывалось новое TLS-соединение.
# После создания не меняется, заголовки идут в каждый запрос — обращений из
# потоков пула APScheduler это выдерживает.
_session = requests.Session()
_session.mount("https://", HTTPAdapter(max_retries=_retry))
_session.mount("http://", HTTPAdapter(max_retries=_retry))
[документация]
def get_observation_jobs() -> list[JobData] | None:
"""Получить расписание проходов с портала.
Ошибки сети логируются и гасятся: недоступность портала не должна
ломать цикл синхронизации.
Returns:
Список заданий либо ``None``, если запрос не удался.
"""
url = settings.url.jobs_url
timeout = settings.api.TIMEOUT_IN_SECONDS_REQUEST_DATA
params = {
"ground_station": settings.station.ID,
"lat": settings.station.LATITUDE,
"lon": settings.station.LONGITUDE,
"alt": settings.station.ELEVATION,
"capabilities": ",".join(CAPABILITIES),
}
# Координаты в .env необязательны: без них портал по-прежнему пишет
# last_seen, а положение станции берётся с него же (remote_config).
params = {key: value for key, value in params.items() if value is not None}
headers = {"Authorization": f"Token {settings.station.TOKEN}"}
logger.debug("Получение заданий из сети")
try:
response = _session.get(
url,
params=params,
headers=headers,
timeout=timeout,
)
response.raise_for_status()
except requests.Timeout:
logger.error("Не удалось получить задания из сети - превышено время ожидания")
return
except requests.HTTPError as e:
logger.error(
"Не удалось получить задания из сети, статус-код: %s, ошибка: %s",
e.response.status_code,
e,
)
return
except requests.RequestException as e:
logger.error("Не удалось получить задания из сети, ошибка: %s", e)
return
_check_clock(response)
# Разбор ответа тоже гасится: битый JSON давал ValueError, изменившаяся
# схема — KeyError/TypeError, и падал весь цикл синхронизации.
try:
jobs_data = response.json()
except ValueError as e:
logger.error("Не удалось разобрать ответ портала с заданиями: %s", e)
return
# Не список — схема ответа изменилась. Пустой список здесь означал бы
# «портал снял все проходы» и стёр бы расписание (см. sync.py), поэтому
# именно None.
if not isinstance(jobs_data, list):
logger.error(
"Портал вернул задания не списком: %s", type(jobs_data).__name__
)
return
# Разбор по одному заданию: одна битая запись пропускается, остальные
# планируются. Раньше список собирался внутри одного try, и битая запись
# обнуляла всё расписание.
jobs = []
for job_data in jobs_data:
try:
jobs.append(JobData.from_dict(job_data))
except (KeyError, TypeError, AttributeError) as e:
logger.error("Задание пропущено, разбор не удался: %s (%s)", e, job_data)
logger.debug("Получено заданий: %d", len(jobs))
return jobs
_status_published = False
_published_body: dict | None = None
[документация]
def status_body() -> dict:
"""Документ статуса станции: режимы, версия, конфигурация.
Один и тот же для ``POST status/`` и для MQTT. Портал сливает его со своей
копией по ключам верхнего уровня, поэтому ключи здесь — собственность
клиента: ``modes``, ``client_version``, ``config``, ``sdr``,
``calibration``. Два последних есть только после обследования приёмника —
отсутствующий ключ портал не трогает.
"""
# Ленивый импорт: remote_config и sat_data пользуются сессией этого модуля.
from soniks_client import remote_config, sat_data, sdr_survey
return {
"modes": sorted(MODES),
# NORAD спутников с satyaml: их станция декодирует gr-satellites
# независимо от режима передатчика (Фаза 3.5). Портал, который ключа
# не знает, отбрасывает его — планирование остаётся по ``modes``.
"satellites": sorted(sat_data.known_norads()),
"client_version": __version__,
"config": remote_config.report(),
**sdr_survey.report(),
}
[документация]
def publish_station_status() -> None:
"""Сообщить порталу режимы станции, её версию и состояние конфигурации.
Портал не планирует станции режим, которого нет в её списке, — так
закрывается подмена неизвестного режима на FM (канал C в
docs/roadmap-network.md). Станция, которая ничего не объявила, для портала
умеет всё, поэтому неудача здесь ничего не ломает: остаётся прежнее
поведение.
Публикуется при каждом изменении документа (применённая конфигурация — тоже
изменение) и повторяется на каждой сверке расписания, пока не удастся.
404 значит, что портал эндпоинта ещё не знает: это не ошибка станции.
"""
global _status_published, _published_body
body = status_body()
if body == _published_body:
return
url = urljoin(settings.url.stations_url, f"{settings.station.ID}/status/")
headers = {"Authorization": f"Token {settings.station.TOKEN}"}
try:
response = _session.post(
url,
json=body,
headers=headers,
timeout=settings.api.TIMEOUT_IN_SECONDS_REQUEST_DATA,
)
response.raise_for_status()
except requests.HTTPError as e:
log = logger.debug if e.response.status_code == 404 else logger.warning
log(
"Режимы станции не опубликованы, статус-код: %s",
e.response.status_code,
)
return
except requests.RequestException as e:
logger.warning("Режимы станции не опубликованы: %s", e)
return
_status_published = True
_published_body = body
logger.info("Статус станции опубликован на портале: режимов %d", len(body["modes"]))
[документация]
def modes_published() -> bool:
"""Портал принял режимы станции в этом процессе — и планирует по ним.
Принять их может только портал, который фильтрует планирование: эндпоинт
и фильтр пришли одним изменением. Поэтому ``True`` значит, что режим вне
``MODES`` станции планироваться не должен был.
"""
return _status_published
def _check_clock(response: requests.Response) -> None:
"""Сверить часы станции с заголовком ``Date`` ответа портала.
NTP на станции никто не проверяет, а уход часов виден только порталу —
отказом на каждый кадр. Запрос расписания и так идёт раз в минуту, поэтому
сверка не стоит ни одного лишнего обращения. Точность — секунды, для
допуска в пять минут этого достаточно.
"""
try:
portal_time = parsedate_to_datetime(response.headers["Date"])
except (KeyError, TypeError, ValueError):
set_clock_skew(None)
return
# Плюс — часы станции спешат, минус — отстают.
skew = int((datetime.now(UTC) - portal_time).total_seconds())
set_clock_skew(skew)
if abs(skew) > CLOCK_SKEW_WARNING_SECONDS:
logger.warning(
"Часы станции расходятся с порталом на %d с. Кадры с меткой дальше "
"5 минут от окна наблюдения портал отвергает — проверьте NTP",
skew,
)
[документация]
def upload_observation_data(
observation_id: str,
data_to_send: DataToSend,
file_path: Path,
) -> None:
"""Выгрузить один файл наблюдения.
Args:
observation_id: Идентификатор наблюдения на портале.
data_to_send: Пара «поле формы — содержимое». Словарь при отправке
опустошается.
file_path: Путь файла, нужен для имени и MIME-типа.
Raises:
ObservationNotFoundError: Наблюдение удалено на портале.
FileNotUploadedError: Выгрузить не удалось, нужна повторная попытка.
"""
observation_url, headers = _get_url_and_headers_of_observation(observation_id)
logger.debug("Загрузка данных [наблюдение: %s]: %s", observation_id, file_path)
timeout = settings.api.TIMEOUT_IN_SECONDS_SEND_DATA
field_name, value = data_to_send.popitem()
mime_type, _ = mimetypes.guess_type(file_path.name)
if mime_type is None:
mime_type = "application/octet-stream"
files = {field_name: (file_path.name, value, mime_type)}
try:
response = _session.put(
observation_url,
headers=headers,
files=files,
timeout=timeout,
)
response.raise_for_status()
logger.info(
"Успешно загружены данные [наблюдение: %s]: %s",
observation_id,
file_path,
)
except requests.Timeout as e:
logger.error(
"Не удалось загрузить данные [наблюдение: %s]: %s - превышено время ожидания",
observation_id,
file_path,
)
raise FileNotUploadedError() from e
except requests.HTTPError as e:
if response.status_code == 404:
logger.error(
"Не удалось загрузить данные [наблюдение: %s]: %s , "
"url: %s не существует (статус-код: 404). "
"Возможно, наблюдение было удалено",
observation_id,
file_path,
observation_url,
)
raise ObservationNotFoundError() from e
elif response.status_code == 400:
# Постоянный отказ: имя кадра не по контракту
# (malformed_filename) или метка дальше допуска от окна
# наблюдения (frame_outside_window). Код разбирать не нужно —
# любой 400 на те же байты повторится.
logger.error(
"Портал отверг данные [наблюдение: %s]: %s (статус-код: 400), "
"ответ: %s",
observation_id,
file_path,
response.text,
)
raise FileRejectedError() from e
elif (
response.status_code == 403 and "has already been uploaded" in response.text
):
# Единственный случай, когда ошибка портала означает успех: файл уже
# лежит на той стороне, повторять отправку незачем.
logger.warning(
"Данные уже были загружены ранее [наблюдение: %s]: %s",
observation_id,
file_path,
)
else:
logger.error(
"Не удалось загрузить данные [наблюдение: %s]: %s, статус-код: %s, ответ: %s",
observation_id,
file_path,
response.status_code,
response.text,
)
raise FileNotUploadedError() from e
except requests.RequestException as e:
logger.error(
"Не удалось загрузить данные [наблюдение: %s]: %s - ошибка клиента: %s",
observation_id,
file_path,
e,
)
raise FileNotUploadedError() from e
except Exception as e:
# Без этого проброса непредвиденная ошибка возвращалась как успех, и
# вызывающий удалял невыгруженный файл.
logger.exception(
"Непредвиденная ошибка при загрузке данных [наблюдение: %s]: %s",
observation_id,
file_path,
)
raise FileNotUploadedError() from e
def _get_url_and_headers_of_observation(
observation_id: str,
) -> tuple[str, dict[str, str]]:
base_observation_url = settings.url.observations_url
observation_url = urljoin(base_observation_url, observation_id)
if not observation_url.endswith("/"):
observation_url += "/"
headers = {"Authorization": f"Token {settings.station.TOKEN}"}
return observation_url, headers