/
serohvostov
/
radio_proxy
Обзор
Документация
Войти
/
serohvostov
/
radio_proxy
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/grid_file_handler.py
611 строк
31 KB
serohvostov
wip
09 фев 2026, 10:36
09 фев 2026, 10:36
80d625e
Код
Авторство
О чём код?
""" Обработчик файлового режима для режима grid """ import os import logging from pathlib import Path from typing import Dict, Any, Optional, Tuple import time from .models import ( BitstreamPacket, MetadataPacket, PacketPair, GridFileInfo, ContentTypes, GridModeError ) from .file_name_generator import FileNameGenerator from .json_mapper import JSONMapper class PacketPairTracker: """Класс для отслеживания пакетов и их сопоставления через поле files""" def __init__(self): """Инициализация трекера пар пакетов""" # Хранилище пакетов 1210 по их unique_id self.bitstream_packets: Dict[int, BitstreamPacket] = {} # Хранилище пакетов 1211 по session_id (UUID) self.metadata_packets: Dict[str, MetadataPacket] = {} self.logger = logging.getLogger(__name__) def add_bitstream_packet(self, packet: BitstreamPacket) -> None: """ Добавить пакет битового потока Args: packet: Пакет битового потока """ self.bitstream_packets[packet.unique_id] = packet self.logger.debug(f'Добавлен пакет битового потока с unique_id: {packet.unique_id}') def add_metadata_packet(self, packet: MetadataPacket) -> None: """ Добавить пакет метаданных Args: packet: Пакет метаданных """ self.metadata_packets[packet.session_id] = packet self.logger.debug(f'Добавлен пакет метаданных с session_id: {packet.session_id}, ' f'files: {packet.files}') def find_matching_bitstream_packets(self, metadata_packet: MetadataPacket) -> list[BitstreamPacket]: """ Найти все пакеты битового потока, соответствующие пакету метаданных Args: metadata_packet: Пакет метаданных Returns: Список пакетов битового потока, ID которых присутствуют в поле files """ matching_packets = [] for file_id in metadata_packet.files: if file_id in self.bitstream_packets: matching_packets.append(self.bitstream_packets[file_id]) self.logger.debug(f'Найден соответствующий пакет битового потока: {file_id}') else: self.logger.debug(f'Пакет битового потока {file_id} еще не получен') return matching_packets def find_matching_metadata_packet(self, bitstream_packet: BitstreamPacket) -> Optional[MetadataPacket]: """ Найти пакет метаданных, который ссылается на данный пакет битового потока Args: bitstream_packet: Пакет битового потока Returns: Пакет метаданных, если найден, иначе None """ for metadata_packet in self.metadata_packets.values(): if bitstream_packet.unique_id in metadata_packet.files: self.logger.debug(f'Найден соответствующий пакет метаданных с session_id: {metadata_packet.session_id}') return metadata_packet self.logger.debug(f'Пакет метаданных для bitstream {bitstream_packet.unique_id} еще не получен') return None def get_complete_pairs(self, metadata_packet: MetadataPacket) -> list[Tuple[BitstreamPacket, MetadataPacket]]: """ Получить все полные пары для данного пакета метаданных Args: metadata_packet: Пакет метаданных Returns: Список кортежей (BitstreamPacket, MetadataPacket) """ matching_bitstreams = self.find_matching_bitstream_packets(metadata_packet) # Проверить, что все файлы из списка files получены if len(matching_bitstreams) == len(metadata_packet.files): return [(bs, metadata_packet) for bs in matching_bitstreams] return [] def is_metadata_complete(self, metadata_packet: MetadataPacket) -> bool: """ Проверить, получены ли все пакеты битового потока для данного пакета метаданных Args: metadata_packet: Пакет метаданных Returns: True если все пакеты битового потока получены """ if not metadata_packet.files: return False for file_id in metadata_packet.files: if file_id not in self.bitstream_packets: return False return True def remove_processed_packets(self, metadata_packet: MetadataPacket) -> None: """ Удалить обработанные пакеты Args: metadata_packet: Пакет метаданных """ # Удалить пакеты битового потока for file_id in metadata_packet.files: if file_id in self.bitstream_packets: del self.bitstream_packets[file_id] self.logger.debug(f'Удален пакет битового потока: {file_id}') # Удалить пакет метаданных if metadata_packet.session_id in self.metadata_packets: del self.metadata_packets[metadata_packet.session_id] self.logger.debug(f'Удален пакет метаданных: {metadata_packet.session_id}') def cleanup_old_pairs(self, max_age_seconds: int = 3600) -> None: """ Очистить старые пакеты (не реализовано в текущей версии) Args: max_age_seconds: Максимальный возраст пакета в секундах """ # TODO: Добавить временные метки к пакетам для очистки старых pass class GridFileHandler: """Обработчик файлового режима для режима grid""" def __init__(self, files_directory: str, jinja_template_path: str, config=None): """ Инициализация обработчика файлового режима Args: files_directory: Директория для сохранения файлов jinja_template_path: Путь к Jinja2 шаблону config: Конфигурация системы (для отладочных настроек) """ self.files_directory = Path(files_directory) self.logger = logging.getLogger(__name__) self.config = config # Создать директорию если не существует self.files_directory.mkdir(parents=True, exist_ok=True) # Инициализировать компоненты try: self.filename_generator = FileNameGenerator() self.json_mapper = JSONMapper(jinja_template_path) self.packet_tracker = PacketPairTracker() self.logger.info(f'GridFileHandler инициализирован для директории: {self.files_directory}') except Exception as e: raise GridModeError(f'Ошибка инициализации GridFileHandler: {e}') async def process_bitstream_packet(self, packet: BitstreamPacket, connection_id: str) -> Optional[list[GridFileInfo]]: """ Обработать пакет битового потока Args: packet: Пакет битового потока connection_id: Идентификатор соединения Returns: Список GridFileInfo если найден соответствующий пакет метаданных, иначе None """ try: # Добавить пакет в трекер self.packet_tracker.add_bitstream_packet(packet) # Проверить, есть ли соответствующий пакет метаданных metadata_packet = self.packet_tracker.find_matching_metadata_packet(packet) if metadata_packet and self.packet_tracker.is_metadata_complete(metadata_packet): # Обработать все пары для этого пакета метаданных return await self._process_complete_pairs(metadata_packet, connection_id) self.logger.debug(f'Пакет битового потока добавлен, ожидание пакета метаданных для unique_id: {packet.unique_id}') return None except Exception as e: self.logger.error(f'Ошибка обработки пакета битового потока: {e}') # Сохраняем ошибочный пакет для отладки await self._save_error_packet(packet, None, connection_id, str(e)) raise GridModeError(f'Ошибка обработки пакета битового потока: {e}') async def process_metadata_packet(self, packet: MetadataPacket, connection_id: str) -> Optional[list[GridFileInfo]]: """ Обработать пакет метаданных Args: packet: Пакет метаданных connection_id: Идентификатор соединения Returns: Список GridFileInfo если все соответствующие пакеты битового потока получены, иначе None """ try: # Добавить пакет в трекер self.packet_tracker.add_metadata_packet(packet) # Проверить, получены ли все пакеты битового потока if self.packet_tracker.is_metadata_complete(packet): return await self._process_complete_pairs(packet, connection_id) self.logger.debug(f'Пакет метаданных добавлен, ожидание пакетов битового потока. ' f'session_id: {packet.session_id}, files: {packet.files}') return None except Exception as e: self.logger.error(f'Ошибка обработки пакета метаданных: {e}') # Сохраняем ошибочный пакет для отладки await self._save_error_packet(None, packet, connection_id, str(e)) raise GridModeError(f'Ошибка обработки пакета метаданных: {e}') async def _process_complete_pairs(self, metadata_packet: MetadataPacket, connection_id: str) -> list[GridFileInfo]: """ Обработать все полные пары для данного пакета метаданных Args: metadata_packet: Пакет метаданных connection_id: Идентификатор соединения Returns: Список GridFileInfo с информацией о созданных файлах """ try: # Получить все полные пары pairs = self.packet_tracker.get_complete_pairs(metadata_packet) if not pairs: error_msg = f'Не найдены полные пары для session_id: {metadata_packet.session_id}' await self._save_error_packet(None, metadata_packet, connection_id, error_msg) raise GridModeError(error_msg) file_infos = [] # Обработать каждую пару for bitstream_packet, meta_packet in pairs: # Валидировать пару пакетов if not self.validate_packet_pair(bitstream_packet, meta_packet): self.logger.warning(f'Валидация пары пакетов не прошла для unique_id: {bitstream_packet.unique_id}') await self._save_error_packet(bitstream_packet, meta_packet, connection_id, 'Валидация пары пакетов не прошла') continue try: # Сохранить файлы, передавая bitstream_packet для извлечения request_id audio_filename, json_filename = await self.save_files( bitstream_packet, meta_packet ) # Создать информацию о файлах file_info = GridFileInfo( audio_filename=audio_filename, json_filename=json_filename, request_id=bitstream_packet.unique_id, connection_id=connection_id ) file_infos.append(file_info) self.logger.info(f'Пара пакетов обработана: bitstream_id={bitstream_packet.unique_id}, ' f'session_id={meta_packet.session_id}, файлы: {audio_filename}, {json_filename}') except Exception as e: # Сохраняем ошибочную пару пакетов self.logger.error(f'Ошибка сохранения файлов для пары пакетов: {e}') await self._save_error_packet(bitstream_packet, meta_packet, connection_id, str(e)) # Продолжаем обработку других пар continue # Удалить обработанные пакеты из трекера self.packet_tracker.remove_processed_packets(metadata_packet) return file_infos except Exception as e: self.logger.error(f'Ошибка обработки полных пар пакетов: {e}') await self._save_error_packet(None, metadata_packet, connection_id, str(e)) raise GridModeError(f'Ошибка обработки полных пар пакетов: {e}') async def save_files(self, bitstream_packet: BitstreamPacket, metadata_packet: MetadataPacket) -> Tuple[str, str]: """ Сохранить аудиофайл и JSON файл с метаданными Args: bitstream_packet: Пакет битового потока metadata_packet: Пакет метаданных Returns: Кортеж (имя_аудиофайла, имя_json_файла) """ import time import json import asyncio try: start_time = time.time() # Извлечь request_id из bitstream_packet.unique_id request_id = bitstream_packet.unique_id content_type = bitstream_packet.content_type audio_data = bitstream_packet.binary_data self.logger.debug(f'[TIMING] Начало сохранения файлов для request_id={request_id}') # Сгенерировать имя аудиофайла используя request_id и content_type audio_filename = self.filename_generator.generate_filename(request_id, content_type) audio_filepath = self.files_directory / audio_filename # Сгенерировать имя JSON файла используя request_id json_filename = self.filename_generator.generate_json_filename(request_id) json_filepath = self.files_directory / json_filename # Сохранить аудиофайл асинхронно (в отдельном потоке, чтобы не блокировать event loop) audio_save_start = time.time() await asyncio.to_thread(self._write_audio_file, audio_filepath, audio_data) audio_save_duration = time.time() - audio_save_start self.logger.info(f'Аудиофайл сохранен: {audio_filepath} ({len(audio_data)} байт) за {audio_save_duration:.3f}с') # Создать JSON с метаданными, передавая unique_id из bitstream_packet json_create_start = time.time() json_data = self.json_mapper.create_json_from_metadata_packet(metadata_packet, unique_id=request_id) json_create_duration = time.time() - json_create_start self.logger.debug(f'[TIMING] JSON создан за {json_create_duration:.3f}с') # Сохранить JSON файл асинхронно json_save_start = time.time() await asyncio.to_thread(self._write_json_file, json_filepath, json_data) json_save_duration = time.time() - json_save_start self.logger.info(f'JSON файл с метаданными сохранен: {json_filepath} за {json_save_duration:.3f}с') total_duration = time.time() - start_time self.logger.debug(f'[TIMING] Общее время сохранения файлов: {total_duration:.3f}с') return audio_filename, json_filename except Exception as e: self.logger.error(f'Ошибка сохранения файлов: {e}') raise GridModeError(f'Ошибка сохранения файлов: {e}') def _write_audio_file(self, filepath, data): """Синхронная запись аудиофайла (вызывается в отдельном потоке)""" with open(filepath, 'wb') as f: f.write(data) def _write_json_file(self, filepath, data): """Синхронная запись JSON файла (вызывается в отдельном потоке)""" import json with open(filepath, 'w', encoding='utf-8') as f: json.dump(data, f, ensure_ascii=False, indent=2) def validate_packet_pair(self, bitstream_packet: BitstreamPacket, metadata_packet: MetadataPacket) -> bool: """ Валидировать пару пакетов Args: bitstream_packet: Пакет битового потока metadata_packet: Пакет метаданных Returns: True если пара валидна """ try: # Проверить, что пакеты не None if not bitstream_packet or not metadata_packet: self.logger.error('Один из пакетов равен None') return False # Проверить, что unique_id пакета битового потока присутствует в массиве files if bitstream_packet.unique_id not in metadata_packet.files: self.logger.error(f'unique_id пакета битового потока {bitstream_packet.unique_id} ' f'не найден в массиве files пакета метаданных: {metadata_packet.files}') return False # Проверить, что есть аудиоданные if not bitstream_packet.binary_data: self.logger.error('Отсутствуют аудиоданные в пакете битового потока') return False # Проверить поддерживаемые типы контента для режима grid supported_types = [ContentTypes.WAV, ContentTypes.JPEG] if bitstream_packet.content_type not in supported_types: self.logger.warning(f'Неподдерживаемый тип контента в режиме grid: {bitstream_packet.content_type}') return False # Проверить обязательные поля метаданных # time_end и frequency не являются обязательными - могут быть пустыми строками required_fields = ['session_id', 'time_start', 'signal_type', 'radionet_name'] for field in required_fields: if not hasattr(metadata_packet, field) or not getattr(metadata_packet, field): self.logger.error(f'Отсутствует обязательное поле метаданных: {field}') return False # Проверить, что time_end и frequency существуют (но могут быть пустыми) if not hasattr(metadata_packet, 'time_end'): self.logger.error('Отсутствует поле time_end в метаданных') return False if not hasattr(metadata_packet, 'frequency'): self.logger.error('Отсутствует поле frequency в метаданных') return False # Проверить, что session_id является строкой (UUID) if not isinstance(metadata_packet.session_id, str): self.logger.error(f'session_id должен быть строкой (UUID), получен: {type(metadata_packet.session_id)}') return False return True except Exception as e: self.logger.error(f'Ошибка валидации пары пакетов: {e}') return False async def cleanup_old_pairs(self) -> None: """Очистить старые неполные пары пакетов""" try: self.packet_tracker.cleanup_old_pairs() except Exception as e: self.logger.error(f'Ошибка очистки старых пар пакетов: {e}') def get_pending_pairs_count(self) -> int: """ Получить количество ожидающих пакетов Returns: Общее количество пакетов в трекере (битовые потоки + метаданные) """ return len(self.packet_tracker.bitstream_packets) + len(self.packet_tracker.metadata_packets) def get_pending_pairs_details(self) -> Dict[str, Any]: """ Получить детальную информацию об ожидающих пакетах Returns: Словарь с детальной информацией о битовых потоках и метаданных """ details = { 'bitstream_packets': [], 'metadata_packets': [] } # Собираем информацию о пакетах битового потока for unique_id, packet in self.packet_tracker.bitstream_packets.items(): details['bitstream_packets'].append({ 'unique_id': unique_id, 'content_type': packet.content_type, 'data_size': len(packet.binary_data) }) # Собираем информацию о пакетах метаданных for session_id, packet in self.packet_tracker.metadata_packets.items(): # Определяем какие файлы уже получены, а какие ожидаются received_files = [fid for fid in packet.files if fid in self.packet_tracker.bitstream_packets] missing_files = [fid for fid in packet.files if fid not in self.packet_tracker.bitstream_packets] details['metadata_packets'].append({ 'session_id': session_id, 'files_expected': packet.files, 'files_received': received_files, 'files_missing': missing_files, 'is_complete': len(missing_files) == 0 }) return details def get_files_directory(self) -> Path: """ Получить путь к директории файлов Returns: Путь к директории файлов """ return self.files_directory async def _save_error_packet(self, bitstream_packet: Optional[BitstreamPacket], metadata_packet: Optional[MetadataPacket], connection_id: str, error_message: str) -> None: """ Сохранить ошибочные пакеты для отладки Args: bitstream_packet: Пакет битового потока (может быть None) metadata_packet: Пакет метаданных (может быть None) connection_id: Идентификатор соединения error_message: Сообщение об ошибке """ import asyncio import json from datetime import datetime try: # Проверяем, включен ли режим отладки ошибочных пакетов if not (self.config and hasattr(self.config, 'grid') and self.config.grid.debug_save_error_packets): return # Создаем папку для отладочных файлов если её нет debug_dir = Path(self.config.grid.debug_raw_packets_directory) debug_dir.mkdir(parents=True, exist_ok=True) # Создаем имя файла с временной меткой timestamp = datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3] # миллисекунды # Определяем идентификатор для имени файла packet_id = "unknown" if metadata_packet: packet_id = metadata_packet.session_id elif bitstream_packet: packet_id = str(bitstream_packet.unique_id) filename = f"error_grid_pair_{packet_id}_{timestamp}_{connection_id}.json" # Собираем данные для отладки debug_data = { "timestamp": datetime.now().isoformat(), "connection_id": connection_id, "error_message": error_message, "bitstream_packet": None, "metadata_packet": None } # Добавляем данные пакета битового потока если есть if bitstream_packet: debug_data["bitstream_packet"] = { "unique_id": bitstream_packet.unique_id, "timestamp": bitstream_packet.timestamp, "content_type": bitstream_packet.content_type, "binary_data_length": len(bitstream_packet.binary_data), "binary_data_hex_preview": bitstream_packet.binary_data[:100].hex() if len(bitstream_packet.binary_data) > 100 else bitstream_packet.binary_data.hex() } # Добавляем данные пакета метаданных если есть if metadata_packet: debug_data["metadata_packet"] = { "session_id": metadata_packet.session_id, "time_start": metadata_packet.time_start, "time_end": metadata_packet.time_end, "frequency": metadata_packet.frequency, "signal_type": metadata_packet.signal_type, "radionet_name": metadata_packet.radionet_name, "radiodata": metadata_packet.radiodata, "files": metadata_packet.files, "places": metadata_packet.places } # Добавляем дополнительные поля если они присутствуют if metadata_packet.Frequency_average is not None: debug_data["metadata_packet"]["Frequency_average"] = metadata_packet.Frequency_average if metadata_packet.colour_code_s is not None: debug_data["metadata_packet"]["colour_code_s"] = metadata_packet.colour_code_s if metadata_packet.is_bs is not None: debug_data["metadata_packet"]["is_bs"] = metadata_packet.is_bs if metadata_packet.seance_type is not None: debug_data["metadata_packet"]["seance_type"] = metadata_packet.seance_type if metadata_packet.tmo_dmo is not None: debug_data["metadata_packet"]["tmo_dmo"] = metadata_packet.tmo_dmo # Сохраняем в файл асинхронно debug_file_path = debug_dir / filename await asyncio.to_thread(self._write_debug_json, debug_file_path, debug_data) self.logger.info(f"Сохранены отладочные данные ошибочной пары пакетов: {debug_file_path}") except Exception as e: self.logger.warning(f"Ошибка при сохранении отладочных данных ошибочной пары пакетов: {e}") def _write_debug_json(self, filepath, data): """Синхронная запись отладочного JSON файла (вызывается в отдельном потоке)""" import json with open(filepath, 'w', encoding='utf-8') as f: json.dump(data, f, ensure_ascii=False, indent=2)