/
YoungFreddy
/
NetologyPipeLineService
Обзор
Документация
Войти
/
YoungFreddy
/
NetologyPipeLineService
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
llm/client.py
279 строк
11 KB
Igor
Initial commit: LLM сервис с FastAPI, pipeline, ретраями, fallback, кэшем
18 июл 2026, 17:01
18 июл 2026, 17:01
5b26a48
Код
Авторство
О чём код?
"""Модуль LLM: LLMClient с ретраями/таймаутом/fallback и PromptBuilder.""" import hashlib import json import logging import time from typing import Any, Dict, List, Optional import httpx from config.loader import load_config logger = logging.getLogger(__name__) class LLMClient: """Клиент для вызова LLM API с ретраями, таймаутом и fallback-ответом.""" def __init__(self, config: Optional[Dict[str, Any]] = None): """Инициализирует клиент: загружает конфиг, устанавливает параметры модели.""" self.config = config or load_config() self.model = self.config.get("model", "gpt-3.5-turbo") self.temperature = self.config.get("temperature", 0.7) self.max_tokens = self.config.get("max_tokens", 1000) self.api_timeout = self.config.get("api_timeout", 30.0) self.max_retries = self.config.get("max_retries", 3) self.retry_delay_base = self.config.get("retry_delay_base", 1.0) self.fallback_response = self.config.get( "fallback_response", "Сервис временно недоступен, попробуйте позже" ) self.api_key = self.config.get("api_key", "") self.api_base = self.config.get("api_base", "https://api.openai.com/v1") def _build_cache_key(self, user_message: str, system_prompt: str) -> str: """Строит MD5-ключ кэша из сообщения, промпта, модели и температуры.""" key_data = { "user_message": user_message, "system_prompt": system_prompt, "model": self.model, "temperature": self.temperature, } key_str = json.dumps(key_data, sort_keys=True) return hashlib.md5(key_str.encode()).hexdigest() def _structured_log(self, event: str, **kwargs): """Логирует информационное событие в JSON.""" record = {"event": event, "timestamp": time.time(), **kwargs} logger.info(json.dumps(record, ensure_ascii=False)) def _structured_error(self, event: str, **kwargs): """Логирует ошибочное событие в JSON.""" record = {"event": event, "timestamp": time.time(), **kwargs} logger.error(json.dumps(record, ensure_ascii=False)) def _structured_warning(self, event: str, **kwargs): """Логирует предупреждение в JSON.""" record = {"event": event, "timestamp": time.time(), **kwargs} logger.warning(json.dumps(record, ensure_ascii=False)) def _make_api_call(self, messages: List[Dict[str, str]]) -> str: """Выполняет реальный HTTP-вызов к LLM API. Если API-ключ не задан — возвращает заглушку для тестов. Raises: httpx.TimeoutException: при превышении таймаута. httpx.HTTPStatusError: при HTTP-ошибке. httpx.RequestError: при сетевой ошибке. ValueError: при ошибке парсинга ответа. """ if not self.api_key: last = messages[-1]["content"] if messages else "" return f"Response to: {last}" url = f"{self.api_base.rstrip('/')}/chat/completions" headers = { "Authorization": f"Bearer {self.api_key}", "Content-Type": "application/json", } payload = { "model": self.model, "messages": messages, "temperature": self.temperature, "max_tokens": self.max_tokens, "stream": False, } with httpx.Client(timeout=self.api_timeout) as client: resp = client.post(url, json=payload, headers=headers) if resp.status_code == 200: data = resp.json() if "error" in data: err_msg = data["error"].get("message", "Unknown API error") self._structured_error( "llm_api_error", error_message=err_msg, raw_response=json.dumps(data, ensure_ascii=False)[:500], ) raise ValueError(f"Ошибка LLM API: {err_msg}") try: content = data["choices"][0]["message"]["content"] if content is None: content = "" return content except (KeyError, IndexError, TypeError) as e: self._structured_error( "llm_parse_error", error_type=type(e).__name__, error_message=str(e), raw_response=json.dumps(data, ensure_ascii=False)[:500], ) raise ValueError("Не удалось распарсить ответ LLM") resp.raise_for_status() return "" def _is_retryable(self, e: Exception) -> bool: """Определяет, стоит ли повторять запрос при этой ошибке.""" if isinstance(e, httpx.TimeoutException): return True if isinstance(e, httpx.RequestError): return True if isinstance(e, httpx.HTTPStatusError): code = e.response.status_code return code >= 500 or code == 429 return False def _call_with_retry(self, messages: List[Dict[str, str]], cache_key: str) -> str: """Вызывает API с экспоненциальными ретраями для временных ошибок. При исчерпании попыток — возвращает fallback-ответ. """ for attempt in range(1, self.max_retries + 1): try: return self._make_api_call(messages) except ValueError as e: self._structured_error( "llm_error", cache_key=cache_key, attempt=attempt, error_type="ValueError", error_message=str(e), status_code=None, ) return self._handle_failure(cache_key, str(e)) except Exception as e: status_code = getattr(e, "response", None) if status_code is not None: status_code = status_code.status_code self._structured_error( "llm_error", cache_key=cache_key, attempt=attempt, error_type=type(e).__name__, error_message=str(e), status_code=status_code, ) if attempt < self.max_retries and self._is_retryable(e): delay = self.retry_delay_base * (2 ** (attempt - 1)) self._structured_warning( "llm_retry", cache_key=cache_key, attempt=attempt, delay_seconds=delay, ) time.sleep(delay) else: return self._handle_failure(cache_key) return self._handle_failure(cache_key) def _handle_failure(self, cache_key: str, reason: str = "") -> str: """Возвращает fallback-ответ при недоступности модели.""" self._structured_warning( "llm_failure", cache_key=cache_key, fallback_used=True, reason=reason or "превышены ретраи", ) return self.fallback_response def _postprocess_response(self, response: str) -> str: """Очищает ответ: удаляет лишние пробелы, проверяет на пустоту.""" if not isinstance(response, str): return "Сервис вернул некорректный ответ" response = response.strip() if not response: return "Сервис вернул пустой ответ" return response def generate_with_cache( self, user_message: str, system_prompt: str, cache: Any, ) -> str: """Генерирует ответ, используя кэш. Args: user_message: Сообщение пользователя. system_prompt: Системный промпт. cache: Объект кэша (должен иметь методы get/set). Returns: Сгенерированный или закешированный ответ. """ start_time = time.time() cache_key = self._build_cache_key(user_message, system_prompt) self._structured_log( "llm_request_start", cache_key=cache_key, model=self.model, system_prompt=( system_prompt[:50] + "..." if len(system_prompt) > 50 else system_prompt ), user_message=( user_message[:100] + "..." if len(user_message) > 100 else user_message ), ) cached = cache.get(cache_key) if cached is not None: elapsed = time.time() - start_time self._structured_log( "llm_cache_hit", cache_key=cache_key, elapsed_seconds=round(elapsed, 3), ) return cached self._structured_log("llm_cache_miss", cache_key=cache_key) messages = [] if system_prompt: messages.append({"role": "system", "content": system_prompt}) messages.append({"role": "user", "content": user_message}) result = self._call_with_retry(messages, cache_key) result = self._postprocess_response(result) cache.set(cache_key, result) elapsed = time.time() - start_time self._structured_log( "llm_response", cache_key=cache_key, cache_hit=False, elapsed_seconds=round(elapsed, 3), response_length=len(result), ) return result class PromptBuilder: """Сборщик списка сообщений для LLM с поддержкой system/user/assistant.""" def __init__(self, system_prompt: str = ""): """Инициализирует сборщик с опциональным системным промптом.""" self.system_prompt = system_prompt self.messages: List[Dict[str, str]] = [] def add_user_message(self, content: str) -> "PromptBuilder": """Добавляет сообщение пользователя. Возвращает self для chaining.""" self.messages.append({"role": "user", "content": content}) return self def add_assistant_message(self, content: str) -> "PromptBuilder": """Добавляет сообщение ассистента. Возвращает self для chaining.""" self.messages.append({"role": "assistant", "content": content}) return self def build(self) -> List[Dict[str, str]]: """Собирает финальный список сообщений (system + history).""" if self.system_prompt: return [{"role": "system", "content": self.system_prompt}] + self.messages return self.messages