/
serohvostov
/
radio_proxy
Обзор
Документация
Войти
/
serohvostov
/
radio_proxy
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/callback_server.py
457 строк
20 KB
serohvostov
version 2025
12 янв 2026, 10:30
12 янв 2026, 10:30
cd6cb9b
Код
Авторство
О чём код?
""" HTTP-сервер для получения callback-запросов от внешнего сервиса распознавания """ import asyncio import json import logging from typing import Dict, Any, Optional, Callable from aiohttp import web, ClientError from aiohttp.web import Request, Response from src.models import ProxyServiceError, RecognitionServiceError class CallbackServer: """ HTTP-сервер для получения результатов callback от внешнего сервиса распознавания Обрабатывает входящие HTTP POST запросы с результатами распознавания и обрабатывает их согласно требованиям 7.1, 7.2, 7.3, 7.4, 7.5 """ def __init__(self, host: str = "0.0.0.0", port: int = 8080): """ Инициализация callback-сервера Args: host: Хост для привязки сервера port: Порт для привязки сервера """ self.host = host self.port = port self.app = web.Application() self.runner: Optional[web.AppRunner] = None self.site: Optional[web.TCPSite] = None self.logger = logging.getLogger(__name__) self.results_processor = ResultsProcessor() self.callback_handler: Optional[Callable[[Dict[str, Any]], None]] = None # Настраиваем маршруты self.app.router.add_post('/callback', self.handle_callback) self.app.router.add_get('/health', self.health_check) async def start_server(self) -> None: """ Запустить HTTP callback-сервер Требование 7.1: Прослушивать HTTP callback от внешнего сервиса """ try: self.runner = web.AppRunner(self.app) await self.runner.setup() self.site = web.TCPSite(self.runner, self.host, self.port) await self.site.start() self.logger.info(f"Callback-сервер запущен на {self.host}:{self.port}") except Exception as e: self.logger.error(f"Не удалось запустить callback-сервер: {e}") raise RecognitionServiceError(f"Не удалось запустить callback-сервер: {e}") async def stop_server(self) -> None: """Остановить HTTP callback-сервер""" try: if self.site: await self.site.stop() self.site = None if self.runner: await self.runner.cleanup() self.runner = None self.logger.info("Callback-сервер остановлен") except Exception as e: self.logger.error(f"Ошибка остановки callback-сервера: {e}") def set_callback_handler(self, handler: Callable[[Dict[str, Any]], None]) -> None: """ Установить функцию обработчика callback Args: handler: Функция для вызова при получении результатов callback """ self.callback_handler = handler async def handle_callback(self, request: Request) -> Response: """ Handle incoming callback requests from external recognition service Requirements 7.2, 7.3, 7.4, 7.5: - Parse JSON response from external service - Extract transcribed text from successful responses - Handle error responses from external service - Match callback results with original client requests Args: request: HTTP request from external service Returns: HTTP response confirming receipt """ try: # Парсим JSON данные из тела запроса if not request.can_read_body: self.logger.warning("Callback-запрос не содержит тела") return web.Response(status=400, text="Отсутствует тело запроса") try: callback_data = await request.json() except json.JSONDecodeError as e: self.logger.error(f"Неверный JSON в callback-запросе: {e}") return web.Response(status=400, text="Неверный формат JSON") # Валидируем формат callback if not self.results_processor.validate_callback_format(callback_data): self.logger.warning(f"Неверный формат callback: {callback_data}") return web.Response(status=400, text="Неверный формат callback") # Обрабатываем результаты распознавания await self.process_recognition_results(callback_data) # Возвращаем успешный ответ return web.Response(status=200, text="OK") except Exception as e: self.logger.error(f"Ошибка обработки callback: {e}") return web.Response(status=500, text="Внутренняя ошибка сервера") async def process_recognition_results(self, results: Dict[str, Any]) -> None: """ Обработать результаты распознавания из callback Args: results: Данные callback от внешнего сервиса """ try: # Извлекаем job_id из header.id if 'header' not in results or not isinstance(results['header'], dict): self.logger.warning("Callback без обязательного поля header") return header = results['header'] job_id = header.get('id') if not job_id: self.logger.warning("Callback без job_id в header.id") return self.logger.info(f"Обработка результатов callback для задания {job_id}") # Обрабатываем результаты через ResultsProcessor processed_results = self.results_processor.process_results(results) # Вызываем обработчик callback если установлен if self.callback_handler: await asyncio.create_task( self._call_handler_safely(processed_results) ) else: self.logger.warning("Обработчик callback не установлен, результаты не переданы") except Exception as e: self.logger.error(f"Ошибка обработки результатов распознавания: {e}") raise RecognitionServiceError(f"Ошибка обработки результатов: {e}") async def _call_handler_safely(self, results: Dict[str, Any]) -> None: """ Безопасно вызвать обработчик callback Args: results: Обработанные результаты для передачи обработчику """ try: if asyncio.iscoroutinefunction(self.callback_handler): await self.callback_handler(results) else: self.callback_handler(results) except Exception as e: self.logger.error(f"Ошибка в обработчике callback: {e}") async def health_check(self, request: Request) -> Response: """Эндпоинт проверки здоровья""" return web.Response(status=200, text="OK") class ResultsProcessor: """ Процессор для извлечения и валидации результатов распознавания Обрабатывает извлечение текста транскрипции из результатов callback согласно требованию 7.3 """ def __init__(self): self.logger = logging.getLogger(__name__) def validate_callback_format(self, data: Dict[str, Any]) -> bool: """ Validate callback data format Args: data: Callback data to validate Returns: True if format is valid, False otherwise """ try: if 'header' not in data: self.logger.warning("Missing required field: header") return False header = data['header'] if not isinstance(header, dict): self.logger.warning("Header должен быть объектом") return False # Проверяем обязательные поля в header if 'id' not in header: self.logger.warning("Missing required field: header.id") return False if 'status' not in header: self.logger.warning("Missing required field: header.status") return False # Валидируем job_id в header job_id = header.get('id') if not isinstance(job_id, str) or not job_id.strip(): self.logger.warning(f"Invalid job_id format in header: {job_id}") return False # Валидируем status в header status = header.get('status') if not isinstance(status, str): self.logger.warning(f"Invalid status format in header: {status}") return False return True except Exception as e: self.logger.error(f"Error validating callback format: {e}") return False def extract_transcription_from_result(self, results: Dict[str, Any]) -> Optional[str]: """ Извлечь текст транскрипции из результатов в формате result.find Args: results: Результаты распознавания от внешнего сервиса Returns: Текст транскрипции в формате "Канал 1: текст. Канал 2: текст" если доступен, None иначе """ try: # Проверяем обязательную структуру result.find if 'result' not in results: self.logger.warning("Отсутствует поле 'result' в результатах") return None result = results['result'] if not isinstance(result, dict) or 'find' not in result: self.logger.warning("Отсутствует поле 'find' в result") return None find_results = result['find'] if not isinstance(find_results, list): self.logger.warning("Поле 'find' не является списком") return None # Собираем текст из всех каналов channel_texts = [] for channel_index, channel_result in enumerate(find_results, 1): if not isinstance(channel_result, dict) or 'words' not in channel_result: continue words = channel_result['words'] if not isinstance(words, list): continue # Извлекаем текст из слов канала channel_text_parts = [] for word in words: if isinstance(word, dict) and 'text' in word: text = word['text'] if isinstance(text, str) and text.strip(): channel_text_parts.append(text.strip()) if channel_text_parts: channel_text = ' '.join(channel_text_parts) channel_texts.append((channel_index, channel_text)) if channel_texts: # Если канал только один, не показываем номер канала if len(channel_texts) == 1: full_text = channel_texts[0][1] # Просто текст без "Канал 1:" else: # Если каналов несколько, показываем с номерами formatted_channels = [f"Канал {idx}: {text}" for idx, text in channel_texts] full_text = '. '.join(formatted_channels) self.logger.info(f"Извлечен текст транскрипции: {full_text[:100]}...") return full_text self.logger.info("Не найден текст транскрипции в результатах") return None except Exception as e: self.logger.error(f"Ошибка извлечения текста транскрипции: {e}") return None def handle_error_response(self, results: Dict[str, Any]) -> Dict[str, Any]: """ Handle error responses from external service Requirement 7.4: Handle error responses from external service Args: results: Error response from external service Returns: Processed error information """ try: # Извлекаем job_id и status из header if 'header' not in results or not isinstance(results['header'], dict): self.logger.warning("Error response без обязательного поля header") return { 'job_id': None, 'status': 'error', 'error': True, 'error_message': 'Missing header in error response', 'error_code': 'invalid_format' } header = results['header'] job_id = header.get('id') status = header.get('status', 'unknown') error_info = { 'job_id': job_id, 'status': status.lower(), # Приводим к нижнему регистру 'error': True, 'error_message': None, 'error_code': None } # Извлекаем информацию об ошибке из result (основной источник для вашего сервиса) if 'result' in results and isinstance(results['result'], dict): result = results['result'] # Код ошибки из result.code if 'code' in result: error_info['error_code'] = result['code'] # Сообщение об ошибке из result.description if 'description' in result: error_info['error_message'] = str(result['description']) # Дополнительно проверяем корневые поля (для совместимости) if not error_info['error_message']: error_fields = ['error', 'error_message', 'message', 'description'] for field in error_fields: if field in results: error_info['error_message'] = str(results[field]) break if not error_info['error_code']: code_fields = ['error_code', 'code', 'error_type'] for field in code_fields: if field in results: error_info['error_code'] = results[field] break # Проверяем поля в header (дополнительно) if not error_info['error_message'] and 'error_message' in header: error_info['error_message'] = str(header['error_message']) if not error_info['error_code'] and 'error_code' in header: error_info['error_code'] = header['error_code'] self.logger.warning(f"Recognition error for job {error_info['job_id']}: " f"{error_info['error_message']} (code: {error_info['error_code']})") return error_info except Exception as e: self.logger.error(f"Error handling error response: {e}") return { 'job_id': None, 'status': 'error', 'error': True, 'error_message': f"Failed to process error response: {e}", 'error_code': 'processing_error' } def process_results(self, results: Dict[str, Any]) -> Dict[str, Any]: """ Process callback results and extract relevant information Args: results: Raw callback results Returns: Processed results with extracted information """ try: # Извлекаем job_id и status из header if 'header' not in results or not isinstance(results['header'], dict): self.logger.warning("Callback results без обязательного поля header") return { 'job_id': None, 'status': 'error', 'error': True, 'error_message': 'Missing header in callback results', 'raw_results': results } header = results['header'] job_id = header.get('id') status = header.get('status', '').lower() processed = { 'job_id': job_id, 'status': status, 'error': False, 'text': None, 'error_message': None, 'raw_results': results } # Проверяем статус на ошибку (учитываем разные варианты написания) error_statuses = ['error', 'failed', 'failure'] if status.lower() in error_statuses: error_info = self.handle_error_response(results) processed.update(error_info) else: # Извлекаем текст транскрипции из result.find (единственный поддерживаемый формат) text = self.extract_transcription_from_result(results) processed['text'] = text if text is None: self.logger.info(f"Текст транскрипции недоступен для задания {job_id}") return processed except Exception as e: self.logger.error(f"Error processing results: {e}") return { 'job_id': None, 'status': 'error', 'error': True, 'error_message': f"Failed to process results: {e}", 'raw_results': results }