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

"""Обмен с порталом СОНИКС: расписание проходов, метаданные и выгрузка файлов.

Ошибки 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_metadata( observation_id: str, data_to_send: DataToSend, file_path: Path, ) -> None: """Выгрузить метаданные наблюдения. Сигнатура общая с :func:`upload_observation_data` — метаданные проходят через ту же очередь. Отличие одно: содержимое уезжает полем формы, а не multipart-файлом, поэтому портал перезаписывает поля наблюдения и повтор безопасен без content-hash. Args: observation_id: Идентификатор наблюдения на портале. data_to_send: Пара «поле формы — содержимое». Словарь при отправке опустошается. file_path: Путь файла, нужен для сообщений в логе. Raises: ObservationNotFoundError: Наблюдение удалено на портале. FileNotUploadedError: Выгрузить не удалось, нужна повторная попытка. """ observation_url, headers = _get_url_and_headers_of_observation(observation_id) timeout = settings.api.TIMEOUT_IN_SECONDS_SEND_DATA field_name, value = data_to_send.popitem() data = { "client_version": __version__, field_name: value, } try: response = _session.put( observation_url, data=data, headers=headers, timeout=timeout, ) response.raise_for_status() logger.info("Успешно загружены метаданные [наблюдение: %s]", observation_id) except requests.Timeout as e: logger.error( "Не удалось загрузить метаданные [наблюдение: %s]: - превышено время ожидания", observation_id, ) raise FileNotUploadedError() from e except requests.HTTPError as e: if response.status_code == 404: logger.error( "Не удалось загрузить метаданные [наблюдение: %s], " "url: %s не существует (статус-код: 404). " "Возможно, наблюдение было удалено", observation_id, observation_url, ) raise ObservationNotFoundError() from e if response.status_code == 400: logger.error( "Портал отверг метаданные [наблюдение: %s] (статус-код: 400), " "ответ: %s", observation_id, response.text, ) raise FileRejectedError() from e logger.error( "Не удалось загрузить метаданные [наблюдение: %s], статус-код: %s, ответ: %s", observation_id, response.status_code, response.text, ) raise FileNotUploadedError() from e except requests.RequestException as e: logger.error( "Не удалось загрузить метаданные [наблюдение: %s], ошибка: %s", observation_id, e, ) raise FileNotUploadedError() from e except Exception as e: # Как и в upload_observation_data: без проброса непредвиденная ошибка # вернулась бы успехом, и вызывающий удалил бы невыгруженный файл. logger.exception( "Непредвиденная ошибка при загрузке метаданных [наблюдение: %s]: %s", observation_id, file_path, ) raise FileNotUploadedError() from e
[документация] 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