Исходный код soniks_client.jobs.sending_data

"""Выгрузка файлов наблюдения и раскладка их по директориям по итогу отправки."""

import threading
from collections.abc import Callable
from pathlib import Path

from core.configs import logger, settings
from core.exceptions import (
    FileNotUploadedError,
    FileRejectedError,
    ObservationNotFoundError,
)
from core.types import DataToSend
from soniks_client.api import upload_observation_data, upload_observation_metadata
from soniks_client.jobs.file_manager.file_paths import (
    get_paths_observation_files,
)
from soniks_client.jobs.file_manager.move_file import (
    delete_data_file,
    delete_observation_directory,
    move_data_file_to_complete_directory,
    move_file_to_incomplete_directory,
    move_observation_directory_to_complete_directory,
)
from soniks_client.jobs.file_manager.read_files import (
    add_demoddata_to_sending_data,
    add_metadata_to_sending_data,
    add_payload_to_sending_data,
    add_waterfall_to_sending_data,
)

# Наблюдения, повторная отправка которых идёт прямо сейчас. Задание повтора
# ставится раз в RESENDING_INTERVAL_IN_MINUTES с replace_existing=True, но это
# переставляет запланированное задание, а не отменяет уже запущенное: медленный
# и свежий повтор обходили одну директорию одновременно.
_resending_observations: set[str] = set()
_resending_lock = threading.Lock()


def _handlers_by_prefix() -> dict[str, Callable[[Path], DataToSend]]:
    """Соответствие «префикс имени файла — читатель для выгрузки».

    Собирается на каждый вызов, а не в константе модуля: и префиксы, и сами
    читатели подменяются в тестах.
    """

    return {
        settings.observation.prefix.METADATA: add_metadata_to_sending_data,
        settings.observation.prefix.PAYLOAD: add_payload_to_sending_data,
        settings.observation.prefix.WATERFALL: add_waterfall_to_sending_data,
        settings.observation.prefix.DATA: add_demoddata_to_sending_data,
    }


def _uploader(prefix: str) -> Callable[[str, DataToSend, Path], None]:
    """Выбрать отправителя по префиксу файла.

    Метаданные уходят полями формы, всё остальное — multipart-файлом.
    """

    if prefix == settings.observation.prefix.METADATA:
        return upload_observation_metadata

    return upload_observation_data


def _close_observation_directory(observation_dir_path: Path) -> None:
    """Убрать директорию наблюдения, все файлы которой пристроены."""

    if settings.observation.REMOVE_OBSERVATION_DATA:
        delete_observation_directory(observation_dir_path)
    else:
        move_observation_directory_to_complete_directory(observation_dir_path)


def _close_data_file(file_path: Path) -> None:
    """Убрать файл, который больше не нужно отправлять."""

    if settings.observation.REMOVE_OBSERVATION_DATA:
        delete_data_file(file_path)
    else:
        move_data_file_to_complete_directory(file_path)


def _defer_observation_to_incomplete(observation_dir_path: Path) -> None:
    """Отложить оставшиеся файлы наблюдения в ``incomplete``.

    Вызывается, когда портал ответил 404. Данные не удаляются с первой
    попытки: решение об удалении принимает повторная отправка, до которой у
    портала есть ещё как минимум ``RESENDING_INTERVAL_IN_MINUTES``.
    """

    # Список целиком до первого перемещения: обход директории, из которой в это
    # же время уезжают файлы, зависит от порядка readdir.
    for file_path in list(observation_dir_path.iterdir()):
        if file_path.is_file():
            move_file_to_incomplete_directory(file_path)

    remaining = list(observation_dir_path.iterdir())

    if remaining:
        logger.error(
            "Часть файлов не удалось перенести в incomplete. Директория "
            "оставлена в output для ручного разбора: %s",
            observation_dir_path,
        )

        return

    delete_observation_directory(observation_dir_path)


[документация] def send_data_during_observation(file_paths: list[Path], observation_id: str) -> None: """Выгрузить пачку декодированных кадров, появившихся во время прохода. По итогу отправки файл удаляется, переезжает в ``complete`` либо в ``incomplete`` для повторной попытки. """ for index, file_path in enumerate(file_paths): try: data_to_send = add_demoddata_to_sending_data(file_path) except Exception: # Падение на одном кадре не должно отменять остальные: раньше # чтение стояло вне ``try`` и первый же нечитаемый файл обрывал # всю пачку. Сам кадр остаётся на диске — его подберёт выгрузка # после прохода. В send_data_after_observation это чинили ещё # в фазе 2, а соседний вызывающий остался без правки. logger.exception("Не удалось прочитать кадр для отправки: %s", file_path) continue try: upload_observation_data( observation_id, data_to_send, file_path, ) except ObservationNotFoundError: # 404 не удаляет данные с первой попытки: кадры уходят в # incomplete, а удалить их может только повторная отправка. # Начало пачки не трогаем — оно уже удалено или уехало в complete. logger.warning( "Наблюдение %s не найдено на портале (404), кадры отложены " "в incomplete до повторной попытки", observation_id, ) for pending_path in file_paths[index:]: move_file_to_incomplete_directory(pending_path) return except FileNotUploadedError: move_file_to_incomplete_directory(file_path) continue _close_data_file(file_path)
[документация] def send_data_after_observation(observation_id: str) -> None: """Выгрузить аудио, водопад и оставшиеся кадры после прохода.""" handlers = _handlers_by_prefix() observation_dir_path, file_paths = get_paths_observation_files(observation_id) # Директорию можно закрывать, только если каждый файл либо ушёл на портал, # либо лёг в incomplete. Иначе единственная копия данных будет удалена. observation_dir_closable = True for type_of_data, paths in file_paths.items(): handler = handlers[type_of_data] for file_path in paths: try: data_to_send = handler(file_path) except FileNotFoundError: logger.error( "Файл не найден, возможно он уже был отправлен: %s", file_path ) continue except Exception: # Падение на одном файле не должно отменять остальные: без # этого PermissionError или битый кадр обрывали задание # целиком, и наблюдение навсегда оставалось в output. logger.exception( "Не удалось прочитать файл для отправки: %s", file_path ) observation_dir_closable = False continue try: _uploader(type_of_data)( observation_id, data_to_send, file_path, ) except ObservationNotFoundError: # Первая попытка выгрузки — не повод удалять единственную копию # данных. Файлы откладываются, удалит их повторная отправка, # если портал ответит 404 и через несколько минут. logger.warning( "Наблюдение %s не найдено на портале (404), файлы отложены " "в incomplete до повторной попытки", observation_id, ) _defer_observation_to_incomplete(observation_dir_path) return except FileNotUploadedError: if not move_file_to_incomplete_directory(file_path): observation_dir_closable = False if not observation_dir_closable: logger.error( "Часть файлов наблюдения %s не удалось ни выгрузить, ни перенести " "в incomplete. Директория оставлена в output для ручного разбора: %s", observation_id, observation_dir_path, ) return _close_observation_directory(observation_dir_path)
[документация] def send_observation_data_again(observation_dir_path: Path) -> None: """Повторно выгрузить файлы наблюдения из директории ``incomplete``. Директория закрывается, только если ушли все файлы. Одну и ту же директорию одновременно обходит не больше одного задания. """ observation_id = observation_dir_path.name with _resending_lock: if observation_id in _resending_observations: logger.debug( "Повторная отправка наблюдения %s уже идёт, задание пропущено", observation_id, ) return _resending_observations.add(observation_id) try: _resend_observation_directory(observation_dir_path, observation_id) finally: with _resending_lock: _resending_observations.discard(observation_id)
def _resend_observation_directory( observation_dir_path: Path, observation_id: str, ) -> None: all_files_sent = True handlers = _handlers_by_prefix() # Список целиком до обхода: отвергнутый порталом файл убирается прямо в # цикле, а обход директории, из которой уезжают файлы, зависит от readdir. for file_path in list(observation_dir_path.iterdir()): # Ищется префикс, а не сам читатель: по префиксу выбирается ещё и # отправитель — метаданные уходят полями формы, остальное файлом. prefix = next( ( prefix for prefix in handlers if file_path.name.startswith(prefix) ), None, ) if prefix is None: # Без этой ветки файл с неизвестным префиксом давал NameError либо # переиспользовал словарь прошлой итерации и уходил не в то поле # формы. Отправляемым он не станет, поэтому директорию не держим. logger.warning( "Неизвестный префикс файла, отправка пропущена [наблюдение: %s]: %s", observation_id, file_path.name, ) continue handler = handlers[prefix] try: data_to_send = handler(file_path) except OSError as e: logger.warning("Не удалось прочитать файл %s: %s", file_path, e) # Директория с нечитаемым файлом не выгружена целиком: без этой # строки она закрывалась как успешная и файл удалялся, ни разу # не уйдя на портал. all_files_sent = False continue try: _uploader(prefix)( observation_id, data_to_send, file_path, ) except ObservationNotFoundError: # Здесь удалять можно: директория лежит в incomplete, то есть это # уже не первая попытка — портал отвечает 404 повторно, спустя # RESENDING_INTERVAL_IN_MINUTES после предыдущей. logger.warning( "Наблюдение %s не найдено на портале (404) при повторной " "отправке, данные удалены", observation_id, ) delete_observation_directory(observation_dir_path) return except FileRejectedError: # Портал отверг файл повторно, как и 404 — не раньше чем через # RESENDING_INTERVAL_IN_MINUTES после прошлой попытки. Те же байты # получат тот же ответ, а файл держал бы в incomplete всё # наблюдение. Закрывается как выгруженный: при # REMOVE_OBSERVATION_DATA=False он остаётся в complete/. logger.error( "Портал отверг файл повторно, отправка прекращена " "[наблюдение: %s]: %s", observation_id, file_path.name, ) _close_data_file(file_path) except FileNotUploadedError: all_files_sent = False if not all_files_sent: logger.warning( "Во время повторной оправки файлов наблюдения: %s " "часть файлов отправить не удалось. " "Повтор через %d мин.", observation_id, settings.scheduler.RESENDING_INTERVAL_IN_MINUTES, ) return _close_observation_directory(observation_dir_path)