/
serohvostov
/
radio_proxy
Обзор
Документация
Войти
/
serohvostov
/
radio_proxy
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/protocol_handler.py
973 строки
51 KB
serohvostov
wip
22 фев 2026, 13:59
22 фев 2026, 13:59
7490624
Код
Авторство
О чём код?
""" Обработчик протокола для парсинга и обработки пакетов """ import struct import json import logging import asyncio import time from pathlib import Path from typing import Union, Tuple, Optional from src.models import ( PacketHeader, RegistrationRequest, RegistrationResponse, BitstreamPacket, MetadataPacket, ResultsPacket, PacketTypes, ProtocolError ) from src.logging_config import ErrorHandler, log_packet_event logger = logging.getLogger(__name__) error_handler = ErrorHandler(logger) class PacketParser: """Парсер для пакетов протокола""" def parse_header(self, data: bytes) -> PacketHeader: """Парсить заголовок пакета из байтов""" return PacketHeader.from_bytes(data) def parse_registration_request(self, data: bytes) -> RegistrationRequest: """Parse registration request packet data""" if len(data) < 8: raise ProtocolError("Недостаточно данных для запроса регистрации") client_number, protocol_version = struct.unpack('<II', data[:8]) try: return RegistrationRequest(client_number=client_number, protocol_version=protocol_version) except ValueError as e: # Преобразуем ошибки валидации модели в ProtocolError raise ProtocolError(f"Ошибка валидации запроса регистрации: {e}") def parse_registration_response(self, data: bytes) -> RegistrationResponse: """Parse registration response packet data""" if len(data) < 8: raise ProtocolError("Недостаточно данных для ответа регистрации") registration_result, protocol_version = struct.unpack('<II', data[:8]) try: return RegistrationResponse(registration_result=registration_result, protocol_version=protocol_version) except ValueError as e: # Преобразуем ошибки валидации модели в ProtocolError raise ProtocolError(f"Ошибка валидации ответа регистрации: {e}") def parse_bitstream_packet(self, data: bytes) -> BitstreamPacket: """Parse bitstream packet data""" if len(data) < 20: # 8 + 8 + 4 = 20 байт для заголовка raise ProtocolError("Недостаточно данных для пакета битового потока") timestamp, unique_id, content_type = struct.unpack('<QQI', data[:20]) binary_data = data[20:] try: return BitstreamPacket( timestamp=timestamp, unique_id=unique_id, content_type=content_type, binary_data=binary_data ) except ValueError as e: # Преобразуем ошибки валидации модели в ProtocolError raise ProtocolError(f"Ошибка валидации пакета битового потока: {e}") def parse_metadata_packet(self, data: bytes) -> MetadataPacket: """Parse metadata packet data""" try: # Удаляем нулевые байты в конце данных (padding) data = data.rstrip(b'\x00') json_str = data.decode('utf-8') # Удаляем нулевые символы из строки JSON (могут быть внутри) json_str = json_str.rstrip('\x00') metadata_dict = json.loads(json_str) except (UnicodeDecodeError, json.JSONDecodeError) as e: raise ProtocolError(f"Неверный JSON в пакете метаданных: {e}") # Проверяем, что JSON содержит объект (словарь) if not isinstance(metadata_dict, dict): raise ProtocolError(f"JSON должен содержать объект, получен: {type(metadata_dict).__name__}") # Extract required fields try: # Преобразуем session_id в строку, если пришло как число session_id_raw = metadata_dict['session_id'] session_id = str(session_id_raw) if not isinstance(session_id_raw, str) else session_id_raw time_start = metadata_dict['time'] time_end = metadata_dict['time_end'] # Преобразуем frequency в int, если пришло как строка frequency_raw = metadata_dict['frequency'] frequency = int(frequency_raw) if isinstance(frequency_raw, (str, int, float)) else frequency_raw signal_type = metadata_dict['signal_type'] radionet_name_raw = metadata_dict['radionet_name'] # Если radionet_name пустая строка, используем значение по умолчанию radionet_name = radionet_name_raw if radionet_name_raw and radionet_name_raw.strip() else "UNKNOWN" except KeyError as e: raise ProtocolError(f"Отсутствует обязательное поле в метаданных: {e}") except (TypeError, ValueError) as e: raise ProtocolError(f"Неверный тип данных в метаданных: {e}") # Extract optional fields radiodata = metadata_dict.get('radiodata', []) # Если radiodata пришел как строка, пытаемся распарсить как JSON if isinstance(radiodata, str): try: radiodata = json.loads(radiodata) # Если после парсинга получился не список, оборачиваем в список if not isinstance(radiodata, list): radiodata = [radiodata] if radiodata else [] except json.JSONDecodeError: # Если не удалось распарсить, сохраняем как есть в виде списка с одной строкой logger.info(f"radiodata не JSON, сохраняем как строку: {radiodata[:100]}") radiodata = [{"raw_data": radiodata}] files = metadata_dict.get('files', []) places = metadata_dict.get('places', []) # Извлекаем дополнительные поля для формата продакшена (если присутствуют) frequency_average = metadata_dict.get('Frequency_average') colour_code_s = metadata_dict.get('colour_code_s') is_bs = metadata_dict.get('is_bs') seance_type = metadata_dict.get('seance_type') tmo_dmo = metadata_dict.get('tmo_dmo') from_phone = metadata_dict.get('from') to_phone = metadata_dict.get('to') rpp_name = metadata_dict.get('rpp_name') # Новое поле для PostNumber try: return MetadataPacket( session_id=session_id, time_start=time_start, time_end=time_end, frequency=frequency, signal_type=signal_type, radionet_name=radionet_name, radiodata=radiodata, files=files, places=places, Frequency_average=frequency_average, colour_code_s=colour_code_s, is_bs=is_bs, seance_type=seance_type, tmo_dmo=tmo_dmo, from_phone=from_phone, to_phone=to_phone, rpp_name=rpp_name ) except ValueError as e: # Преобразуем ошибки валидации модели в ProtocolError raise ProtocolError(f"Ошибка валидации пакета метаданных: {e}") def parse_results_packet(self, data: bytes) -> ResultsPacket: """Parse results packet data""" try: json_str = data.decode('utf-8') results_dict = json.loads(json_str) except (UnicodeDecodeError, json.JSONDecodeError) as e: raise ProtocolError(f"Invalid JSON in results packet: {e}") try: request_id = results_dict['request_id'] except KeyError: raise ProtocolError("Missing request_id in results packet") text = results_dict.get('text') try: return ResultsPacket(request_id=request_id, text=text) except ValueError as e: # Преобразуем ошибки валидации модели в ProtocolError raise ProtocolError(f"Ошибка валидации пакета результатов: {e}") def serialize_registration_request(self, packet: RegistrationRequest) -> bytes: """Serialize registration request to bytes""" return struct.pack('<II', packet.client_number, packet.protocol_version) def serialize_registration_response(self, packet: RegistrationResponse) -> bytes: """Serialize registration response to bytes""" return struct.pack('<II', packet.registration_result, packet.protocol_version) def serialize_bitstream_packet(self, packet: BitstreamPacket) -> bytes: """Serialize bitstream packet to bytes""" header = struct.pack('<QQI', packet.timestamp, packet.unique_id, packet.content_type) return header + packet.binary_data def serialize_metadata_packet(self, packet: MetadataPacket) -> bytes: """Serialize metadata packet to bytes""" metadata_dict = { 'session_id': packet.session_id, 'time': packet.time_start, 'time_end': packet.time_end, 'frequency': packet.frequency, 'signal_type': packet.signal_type, 'radionet_name': packet.radionet_name, 'radiodata': packet.radiodata, 'files': packet.files, 'places': packet.places } # Добавляем дополнительные поля если они присутствуют if packet.Frequency_average is not None: metadata_dict['Frequency_average'] = packet.Frequency_average if packet.colour_code_s is not None: metadata_dict['colour_code_s'] = packet.colour_code_s if packet.is_bs is not None: metadata_dict['is_bs'] = packet.is_bs if packet.seance_type is not None: metadata_dict['seance_type'] = packet.seance_type if packet.tmo_dmo is not None: metadata_dict['tmo_dmo'] = packet.tmo_dmo if packet.from_phone is not None: metadata_dict['from'] = packet.from_phone if packet.to_phone is not None: metadata_dict['to'] = packet.to_phone json_str = json.dumps(metadata_dict, ensure_ascii=False) return json_str.encode('utf-8') def serialize_results_packet(self, packet: ResultsPacket) -> bytes: """Serialize results packet to bytes""" results_dict = {'request_id': packet.request_id} if packet.text is not None: results_dict['text'] = packet.text json_str = json.dumps(results_dict, ensure_ascii=False) return json_str.encode('utf-8') def parse_packet(self, data: bytes) -> Tuple[PacketHeader, Union[RegistrationRequest, RegistrationResponse, BitstreamPacket, MetadataPacket, ResultsPacket]]: """Parse complete packet with header and data""" if len(data) < 12: raise ProtocolError("Insufficient data for packet header") header = self.parse_header(data) packet_data = data[12:12 + header.data_length] if len(packet_data) != header.data_length: raise ProtocolError(f"Data length mismatch: expected {header.data_length}, got {len(packet_data)}") if header.packet_type == PacketTypes.REGISTRATION_REQUEST: return header, self.parse_registration_request(packet_data) elif header.packet_type == PacketTypes.REGISTRATION_RESPONSE: return header, self.parse_registration_response(packet_data) elif header.packet_type == PacketTypes.BITSTREAM: return header, self.parse_bitstream_packet(packet_data) elif header.packet_type == PacketTypes.METADATA: return header, self.parse_metadata_packet(packet_data) elif header.packet_type == PacketTypes.RESULTS: return header, self.parse_results_packet(packet_data) else: raise ProtocolError(f"Unknown packet type: {header.packet_type}") def serialize_packet(self, packet_type: int, packet_data: Union[RegistrationRequest, RegistrationResponse, BitstreamPacket, MetadataPacket, ResultsPacket]) -> bytes: """Serialize complete packet with header and data""" if packet_type == PacketTypes.REGISTRATION_REQUEST: data = self.serialize_registration_request(packet_data) elif packet_type == PacketTypes.REGISTRATION_RESPONSE: data = self.serialize_registration_response(packet_data) elif packet_type == PacketTypes.BITSTREAM: data = self.serialize_bitstream_packet(packet_data) elif packet_type == PacketTypes.METADATA: data = self.serialize_metadata_packet(packet_data) elif packet_type == PacketTypes.RESULTS: data = self.serialize_results_packet(packet_data) else: raise ProtocolError(f"Unknown packet type: {packet_type}") header = PacketHeader(packet_type=packet_type, data_length=len(data)) return header.to_bytes() + data class ProtocolHandler: """Handles protocol packet routing and client registration""" def __init__(self, connection_manager, supported_protocol_versions=None, mode="ts", grid_file_handler=None, config=None): """ Initialize protocol handler Args: connection_manager: ConnectionManager instance for managing connections supported_protocol_versions: List of supported protocol versions (default: [1]) mode: Режим работы ("ts" или "grid") grid_file_handler: Обработчик файлового режима для режима grid config: Конфигурация системы """ self.parser = PacketParser() self.connection_manager = connection_manager self.supported_protocol_versions = supported_protocol_versions or [1] self.max_supported_version = max(self.supported_protocol_versions) self.min_supported_version = min(self.supported_protocol_versions) self._process_audio_callback = None # Callback для обработки аудио # Режим работы и обработчик grid self.mode = mode self.grid_file_handler = grid_file_handler self.config = config self.results_monitor = None # Ссылка на ResultsMonitor для режима grid # Отслеживание соответствия request_id -> connection_id для режима grid self.grid_request_connections = {} # request_id -> connection_id # Отслеживание соответствия filename -> unique_id для режима grid self.grid_filename_mapping = {} # filename_without_extension -> unique_id logger.info(f"ProtocolHandler инициализирован в режиме: {self.mode}") async def _save_debug_raw_packet(self, packet: MetadataPacket, connection_id: str, raw_json_bytes: bytes = None) -> None: """ Сохранить исходные данные пакета 1211 для отладки Args: packet: MetadataPacket с распарсенными данными connection_id: Идентификатор соединения raw_json_bytes: Исходные байты JSON до парсинга (опционально) """ import asyncio try: debug_start_time = time.time() # Проверяем, включен ли режим отладки if not (self.config and hasattr(self.config, 'grid') and self.config.grid.debug_save_raw_packets): return # Создаем папку для отладочных файлов если её нет debug_dir = Path(self.config.grid.debug_raw_packets_directory) debug_dir.mkdir(parents=True, exist_ok=True) # Создаем имя файла с временной меткой и session_id from datetime import datetime timestamp = datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3] # миллисекунды filename = f"raw_packet_1211_{packet.session_id}_{timestamp}_{connection_id}.json" # Декодируем исходный JSON если он передан raw_json_string = None raw_json_parsed = None if raw_json_bytes: try: # Удаляем null-байты перед декодированием cleaned_bytes = raw_json_bytes.rstrip(b'\x00') raw_json_string = cleaned_bytes.decode('utf-8') # Удаляем null-символы из строки cleaned_json_string = raw_json_string.rstrip('\x00') raw_json_parsed = json.loads(cleaned_json_string) except Exception as e: logger.warning(f"Не удалось декодировать исходный JSON: {e}") raw_json_string = raw_json_bytes.decode('utf-8', errors='replace') # Собираем все данные пакета debug_data = { "timestamp": datetime.now().isoformat(), "connection_id": connection_id, "packet_type": 1211, "session_id": packet.session_id, "raw_json_string": raw_json_string, "raw_json_parsed": raw_json_parsed, "parsed_metadata": { "session_id": packet.session_id, "time_start": packet.time_start, "time_end": packet.time_end, "frequency": packet.frequency, "signal_type": packet.signal_type, "radionet_name": packet.radionet_name, "radiodata": packet.radiodata, "files": packet.files, "places": packet.places } } # Добавляем дополнительные поля если они присутствуют if packet.Frequency_average is not None: debug_data["parsed_metadata"]["Frequency_average"] = packet.Frequency_average if packet.colour_code_s is not None: debug_data["parsed_metadata"]["colour_code_s"] = packet.colour_code_s if packet.is_bs is not None: debug_data["parsed_metadata"]["is_bs"] = packet.is_bs if packet.seance_type is not None: debug_data["parsed_metadata"]["seance_type"] = packet.seance_type if packet.tmo_dmo is not None: debug_data["parsed_metadata"]["tmo_dmo"] = packet.tmo_dmo if packet.from_phone is not None: debug_data["parsed_metadata"]["from_phone"] = packet.from_phone if packet.to_phone is not None: debug_data["parsed_metadata"]["to_phone"] = packet.to_phone # Сохраняем в файл асинхронно (в отдельном потоке) debug_file_path = debug_dir / filename await asyncio.to_thread(self._write_debug_json, debug_file_path, debug_data) debug_duration = time.time() - debug_start_time logger.debug(f"Сохранены отладочные данные пакета 1211: {debug_file_path} за {debug_duration:.3f}с") except Exception as e: logger.warning(f"Ошибка при сохранении отладочных данных пакета 1211: {e}") def _write_debug_json(self, filepath, data): """Синхронная запись отладочного JSON файла (вызывается в отдельном потоке)""" with open(filepath, 'w', encoding='utf-8') as f: json.dump(data, f, ensure_ascii=False, indent=2) async def _save_debug_raw_packet_on_error(self, packet_data: bytes, packet_type: int, connection_id: str, error_message: str) -> None: """ Сохранить сырые данные пакета при ошибке парсинга для отладки Args: packet_data: Сырые данные пакета (только данные, без заголовка) packet_type: Тип пакета connection_id: Идентификатор соединения error_message: Сообщение об ошибке """ 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) # Создаем имя файла с временной меткой from datetime import datetime timestamp = datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3] # миллисекунды filename = f"error_packet_{packet_type}_{timestamp}_{connection_id}.json" # Пытаемся декодировать данные как JSON для пакетов метаданных raw_data_decoded = None if packet_type == PacketTypes.METADATA: try: json_str = packet_data.decode('utf-8') raw_data_decoded = json.loads(json_str) except Exception: # Если не удалось декодировать, сохраним как есть try: raw_data_decoded = packet_data.decode('utf-8', errors='replace') except Exception: raw_data_decoded = packet_data.hex() # Собираем все данные для отладки debug_data = { "timestamp": datetime.now().isoformat(), "connection_id": connection_id, "packet_type": packet_type, "error_message": error_message, "data_length": len(packet_data), "raw_data": raw_data_decoded if raw_data_decoded else packet_data.hex(), "raw_data_hex": packet_data[:200].hex() if len(packet_data) > 200 else packet_data.hex() } # Сохраняем в файл debug_file_path = debug_dir / filename with open(debug_file_path, 'w', encoding='utf-8') as f: json.dump(debug_data, f, ensure_ascii=False, indent=2) logger.info(f"Сохранены отладочные данные ошибочного пакета: {debug_file_path}") except Exception as e: logger.warning(f"Ошибка при сохранении отладочных данных ошибочного пакета: {e}") async def handle_packet(self, packet_data: bytes, connection_id: str) -> None: """ Handle incoming packet by routing to appropriate handler Args: packet_data: Raw packet bytes connection_id: Connection identifier """ start_time = time.time() header = None packet_type = None raw_metadata_bytes = None # Для сохранения исходных байтов метаданных try: # Добавляем отладочное логирование logger.debug(f"Получен пакет от {connection_id}, размер: {len(packet_data)} байт") # Сначала парсим заголовок, чтобы узнать тип пакета if len(packet_data) >= 12: header = self.parser.parse_header(packet_data[:12]) packet_type = header.packet_type # Сохраняем исходные байты для пакета метаданных (до парсинга) if packet_type == PacketTypes.METADATA and len(packet_data) > 12: raw_metadata_bytes = packet_data[12:12 + header.data_length] header, packet = self.parser.parse_packet(packet_data) connection = self.connection_manager.get_connection(connection_id) if not connection: raise ProtocolError(f"Соединение {connection_id} не найдено") log_packet_event(logger, "получен", connection_id, header.packet_type, data_length=header.data_length) # Route packet based on type if header.packet_type == PacketTypes.REGISTRATION_REQUEST: await self._handle_registration_request(packet, connection, connection_id) elif header.packet_type == PacketTypes.BITSTREAM: await self._handle_bitstream_packet(packet, connection, connection_id) elif header.packet_type == PacketTypes.METADATA: await self._handle_metadata_packet(packet, connection, connection_id, raw_metadata_bytes) else: # Ignore unknown packet types as per requirement 2.4 logger.warning(f"Игнорирование неизвестного типа пакета {header.packet_type} от {connection_id}") # duration = time.time() - start_time # log_performance(logger, f"обработка пакета типа {header.packet_type}", duration) except ProtocolError as e: # Сохраняем сырые данные пакета при ошибке парсинга if packet_type and len(packet_data) > 12: # Извлекаем данные пакета без заголовка packet_data_only = packet_data[12:] await self._save_debug_raw_packet_on_error( packet_data_only, packet_type, connection_id, str(e) ) error_handler.handle_protocol_error(e, connection_id) except Exception as e: # Сохраняем сырые данные пакета при любой ошибке if packet_type and len(packet_data) > 12: packet_data_only = packet_data[12:] await self._save_debug_raw_packet_on_error( packet_data_only, packet_type, connection_id, f"Неожиданная ошибка: {str(e)}" ) error_handler.handle_unexpected_error(e, f"обработка пакета от {connection_id}") async def _handle_registration_request(self, packet: RegistrationRequest, connection, connection_id: str) -> None: """ Handle client registration request Args: packet: RegistrationRequest packet connection: Connection state object connection_id: Connection identifier """ logger.info(f"Запрос регистрации от {connection_id}: клиент={packet.client_number}, версия={packet.protocol_version}") try: # Validate protocol version compatibility if packet.protocol_version < self.min_supported_version or packet.protocol_version > self.max_supported_version: # Protocol version not supported agreed_version = self.min_supported_version registration_result = 1 # failure logger.warning(f"Неподдерживаемая версия протокола {packet.protocol_version} от {connection_id}") else: # Use the minimum of requested and max supported version agreed_version = min(packet.protocol_version, self.max_supported_version) registration_result = 0 # success # Update connection state connection.client_number = packet.client_number connection.protocol_version = agreed_version connection.is_registered = (registration_result == 0) logger.info(f"Регистрация успешна для {connection_id}: клиент={packet.client_number}, согласованная_версия={agreed_version}") # Send registration response await self.send_registration_response(connection_id, registration_result, agreed_version) except Exception as e: error_handler.handle_protocol_error(e, f"регистрация клиента {connection_id}") async def _handle_bitstream_packet(self, packet: BitstreamPacket, connection, connection_id: str) -> None: """ Handle bitstream packet (audio data) Args: packet: BitstreamPacket connection: Connection state object connection_id: Connection identifier """ try: # Check if client is registered if not connection.is_registered: logger.warning(f"Отклонение пакета битового потока от незарегистрированного клиента {connection_id}") return # Логируем предупреждение для неизвестного типа контента if packet.content_type == 0: # ContentTypes.UNKNOWN logger.warning(f"Получен пакет битового потока с неизвестным типом контента от {connection_id}: id={packet.unique_id}, тип={packet.content_type}, размер={len(packet.binary_data)}") else: logger.info(f"Получен пакет битового потока от {connection_id}: id={packet.unique_id}, тип={packet.content_type}, размер={len(packet.binary_data)}") # Маршрутизация в зависимости от режима работы if self.mode == "grid" and self.grid_file_handler: # Режим grid - обрабатываем через GridFileHandler await self._handle_bitstream_packet_grid(packet, connection_id) else: # Режим ts - обрабатываем через callback await self._handle_bitstream_packet_ts(packet, connection, connection_id) except Exception as e: error_handler.handle_protocol_error(e, f"обработка битового потока от {connection_id}") async def _handle_bitstream_packet_ts(self, packet: BitstreamPacket, connection, connection_id: str) -> None: """ Обработка пакета битового потока в режиме ts Args: packet: BitstreamPacket connection: Connection state object connection_id: Connection identifier """ # Проверяем, есть ли уже соответствующие метаданные metadata_key = f"metadata_{packet.unique_id}" if metadata_key in connection.pending_requests: # Метаданные уже есть - сохраняем пакет для обработки парой logger.info(f"Найдены соответствующие метаданные для пакета {packet.unique_id}, обработка парой") bitstream_key = f"bitstream_{packet.unique_id}" connection.pending_requests[bitstream_key] = { 'type': 'bitstream', 'packet': packet, 'timestamp': packet.timestamp } else: # Метаданных нет - обрабатываем сразу с дефолтными метаданными logger.info(f"Метаданные для пакета {packet.unique_id} отсутствуют, обработка с дефолтными метаданными") if hasattr(self, '_process_audio_callback') and self._process_audio_callback: # Есть callback - обрабатываем сразу default_metadata = self._create_default_metadata(packet.unique_id) await self._process_audio_callback(packet, default_metadata, connection_id) else: # Нет callback - сохраняем в pending_requests для совместимости с тестами logger.info(f"Callback не установлен, сохраняем пакет {packet.unique_id} в pending_requests") bitstream_key = f"bitstream_{packet.unique_id}" connection.pending_requests[bitstream_key] = { 'type': 'bitstream', 'packet': packet, 'timestamp': packet.timestamp } async def _handle_bitstream_packet_grid(self, packet: BitstreamPacket, connection_id: str) -> None: """ Обработка пакета битового потока в режиме grid Args: packet: BitstreamPacket connection_id: Connection identifier """ try: logger.info(f"Обработка пакета битового потока в режиме grid: id={packet.unique_id}") # Сохраняем связь request_id -> connection_id для последующей отправки результатов self.grid_request_connections[packet.unique_id] = connection_id # Передаем пакет в GridFileHandler (теперь возвращает список) file_info_list = await self.grid_file_handler.process_bitstream_packet(packet, connection_id) if file_info_list: # Пара пакетов полная и обработана (обрабатываем все элементы списка) for file_info in file_info_list: logger.info(f"Файлы созданы в режиме grid: {file_info.audio_filename}, {file_info.json_filename}") # Регистрируем связь между именем файла и unique_id self.register_grid_file(file_info.audio_filename, file_info.request_id) # Регистрируем ожидающий запрос в ResultsMonitor для отслеживания таймаута if self.results_monitor: # Ожидаемое имя файла результатов (без расширения) # В режиме grid используем имя файла из FileNameGenerator expected_result_filename = file_info.json_filename self.results_monitor.register_pending_request(file_info.request_id, expected_result_filename) else: logger.warning(f"ResultsMonitor не установлен, таймаут для запроса {file_info.request_id} не будет отслеживаться") # В режиме grid мы не отправляем немедленный ответ клиенту # Ответ будет отправлен через ResultsMonitor когда появится файл результатов else: logger.debug(f"Пакет битового потока добавлен в grid, ожидание пакета метаданных") except Exception as e: logger.error(f"Ошибка обработки пакета битового потока в режиме grid: {e}") # В случае ошибки отправляем пустой ответ клиенту await self.send_results_packet(connection_id, packet.unique_id, text=None) # Удаляем связь при ошибке self.grid_request_connections.pop(packet.unique_id, None) def _create_default_metadata(self, unique_id: int) -> 'MetadataPacket': """ Create default metadata for bitstream packet without metadata Args: unique_id: Unique identifier from bitstream packet Returns: MetadataPacket: Default metadata packet """ from datetime import datetime current_time = datetime.now().isoformat() return MetadataPacket( session_id=unique_id, time_start=current_time, time_end=current_time, frequency=1, # Минимальная валидная частота signal_type="unknown", radionet_name="default", radiodata=[], files=[], places=[] ) def set_process_audio_callback(self, callback): """ Set callback for processing audio requests Args: callback: Async function to call for processing audio """ self._process_audio_callback = callback def set_grid_file_handler(self, grid_file_handler): """ Set grid file handler for grid mode Args: grid_file_handler: GridFileHandler instance """ self.grid_file_handler = grid_file_handler logger.info("GridFileHandler установлен для режима grid") def set_results_monitor(self, results_monitor): """ Установить монитор результатов для режима grid Args: results_monitor: ResultsMonitor instance """ self.results_monitor = results_monitor logger.info("ResultsMonitor установлен для режима grid") async def _handle_metadata_packet(self, packet: MetadataPacket, connection, connection_id: str, raw_json_bytes: bytes = None) -> None: """ Handle metadata packet Args: packet: MetadataPacket connection: Connection state object connection_id: Connection identifier raw_json_bytes: Исходные байты JSON до парсинга (опционально) """ try: # Check if client is registered if not connection.is_registered: logger.warning(f"Отклонение пакета метаданных от незарегистрированного клиента {connection_id}") return logger.info(f"Получен пакет метаданных от {connection_id}: сессия={packet.session_id}") # Сохраняем отладочные данные если включен режим отладки await self._save_debug_raw_packet(packet, connection_id, raw_json_bytes) # Маршрутизация в зависимости от режима работы if self.mode == "grid" and self.grid_file_handler: # Режим grid - обрабатываем через GridFileHandler await self._handle_metadata_packet_grid(packet, connection_id) else: # Режим ts - сохраняем в pending_requests await self._handle_metadata_packet_ts(packet, connection, connection_id) except Exception as e: error_handler.handle_protocol_error(e, f"обработка метаданных от {connection_id}") async def _handle_metadata_packet_ts(self, packet: MetadataPacket, connection, connection_id: str) -> None: """ Обработка пакета метаданных в режиме ts Args: packet: MetadataPacket connection: Connection state object connection_id: Connection identifier """ # Store metadata with a prefixed key to avoid collision with bitstream packets metadata_key = f"metadata_{packet.session_id}" connection.pending_requests[metadata_key] = { 'type': 'metadata', 'packet': packet } async def _handle_metadata_packet_grid(self, packet: MetadataPacket, connection_id: str) -> None: """ Обработка пакета метаданных в режиме grid Args: packet: MetadataPacket connection_id: Connection identifier """ try: logger.info(f"Обработка пакета метаданных в режиме grid: сессия={packet.session_id}") # Проверяем, что поле files не пустое if not packet.files or len(packet.files) == 0: logger.info(f"Пакет метаданных {packet.session_id} имеет пустое поле files, пропускаем обработку") return # Передаем пакет в GridFileHandler (теперь возвращает список) file_info_list = await self.grid_file_handler.process_metadata_packet(packet, connection_id) if file_info_list: # Пара пакетов полная и обработана (обрабатываем все элементы списка) for file_info in file_info_list: logger.info(f"Файлы созданы в режиме grid: {file_info.audio_filename}, {file_info.json_filename}") # Сохраняем связь request_id -> connection_id для последующей отправки результатов self.grid_request_connections[file_info.request_id] = connection_id # Регистрируем связь между именем файла и unique_id self.register_grid_file(file_info.audio_filename, file_info.request_id) # Регистрируем ожидающий запрос в ResultsMonitor для отслеживания таймаута if self.results_monitor: # Ожидаемое имя файла результатов (без расширения) # В режиме grid используем имя файла из FileNameGenerator expected_result_filename = file_info.json_filename self.results_monitor.register_pending_request(file_info.request_id, expected_result_filename) else: logger.warning(f"ResultsMonitor не установлен, таймаут для запроса {file_info.request_id} не будет отслеживаться") # В режиме grid мы не отправляем немедленный ответ клиенту # Ответ будет отправлен через ResultsMonitor когда появится файл результатов else: logger.debug(f"Пакет метаданных добавлен в grid, ожидание пакета битового потока") except Exception as e: logger.error(f"Ошибка обработки пакета метаданных в режиме grid: {e}") # В случае ошибки для всех файлов из массива files отправляем пустой ответ for file_id in packet.files: await self.send_results_packet(connection_id, file_id, text=None) self.grid_request_connections.pop(file_id, None) async def send_registration_response(self, connection_id: str, result: int, protocol_version: int) -> None: """ Send registration response to client Args: connection_id: Connection identifier result: Registration result (0=success, 1=failure) protocol_version: Agreed protocol version """ writer = self.connection_manager.get_writer(connection_id) if not writer: logger.error(f"Writer не найден для соединения {connection_id}") return try: response_packet = RegistrationResponse( registration_result=result, protocol_version=protocol_version ) packet_bytes = self.parser.serialize_packet(PacketTypes.REGISTRATION_RESPONSE, response_packet) writer.write(packet_bytes) await writer.drain() log_packet_event(logger, "отправлен ответ регистрации", connection_id, PacketTypes.REGISTRATION_RESPONSE, result=result, version=protocol_version) except Exception as e: error_handler.handle_network_error(e, f"отправка ответа регистрации для {connection_id}") async def send_results_packet(self, connection_id: str, request_id: int, text: str = None) -> None: """ Send results packet to client Args: connection_id: Connection identifier request_id: Original request ID text: Transcription text (None if recognition failed) """ writer = self.connection_manager.get_writer(connection_id) if not writer: logger.error(f"Writer не найден для соединения {connection_id}") return try: results_packet = ResultsPacket(request_id=request_id, text=text) packet_bytes = self.parser.serialize_packet(PacketTypes.RESULTS, results_packet) writer.write(packet_bytes) await writer.drain() log_packet_event(logger, "отправлен пакет результатов", connection_id, PacketTypes.RESULTS, request_id=request_id, has_text=(text is not None)) except Exception as e: error_handler.handle_network_error(e, f"отправка пакета результатов для {connection_id}") def register_grid_file(self, filename: str, unique_id: int) -> None: """ Зарегистрировать связь между именем файла и unique_id для режима grid Args: filename: Имя файла (с расширением) unique_id: Уникальный идентификатор пакета """ # Убираем расширение из имени файла from pathlib import Path filename_without_ext = Path(filename).stem self.grid_filename_mapping[filename_without_ext] = unique_id logger.debug(f"Зарегистрирована связь файла: {filename_without_ext} -> unique_id={unique_id}") def get_unique_id_by_filename(self, filename: str) -> Optional[int]: """ Получить unique_id по имени файла для режима grid Args: filename: Имя файла (с расширением или без) Returns: unique_id или None если не найден """ from pathlib import Path filename_without_ext = Path(filename).stem unique_id = self.grid_filename_mapping.get(filename_without_ext) if unique_id: logger.debug(f"Найден unique_id для файла {filename_without_ext}: {unique_id}") else: logger.debug(f"unique_id не найден для файла: {filename_without_ext}") return unique_id async def send_grid_results(self, request_id: int, transcription_text: str) -> None: """ Отправить результаты распознавания клиенту в режиме grid Args: request_id: ID запроса (unique_id пакета) transcription_text: Текст транскрипции (может быть None при ошибке) """ try: # Найти connection_id по request_id connection_id = self.grid_request_connections.get(request_id) if not connection_id: logger.error(f"Не найден connection_id для request_id {request_id} в режиме grid") return # Отправить пакет результатов клиенту await self.send_results_packet(connection_id, request_id, transcription_text) # Удалить связь после отправки результата self.grid_request_connections.pop(request_id, None) logger.info(f"Результаты grid отправлены клиенту {connection_id}: request_id={request_id}, " f"текст={'есть' if transcription_text else 'отсутствует'}") except Exception as e: logger.error(f"Ошибка отправки результатов grid для request_id {request_id}: {e}") # Удаляем связь даже при ошибке, чтобы избежать утечек памяти self.grid_request_connections.pop(request_id, None)