Редактировать rag_client.py

This commit is contained in:
Markov Andrey
2026-06-30 09:46:53 +00:00
parent 0dd4c1ed98
commit 3ed6bbfa26

187
rag/rag_client.py Normal file
View File

@@ -0,0 +1,187 @@
# -*- coding: utf-8 -*-
"""
HTTP-клиент для взаимодействия с RAG-сервером.
Используется ботами вместо локального RAGAPI, когда RAG-ядро вынесено
в отдельный HTTP-сервис.
Клиент отправляет асинхронные HTTP-запросы на эндпоинты /rag/query и /rag/index.
Все методы повторяют сигнатуру RAGAPI для бесшовной замены.
ИСТОРИЯ НЕ ПЕРЕДАЁТСЯ – сервер получает её из БД.
"""
import logging
from typing import Optional, Dict, List, Any
import aiohttp
import asyncio
logger = logging.getLogger(__name__)
class RAGClient:
"""
HTTP-клиент для RAG-сервера.
"""
def __init__(
self,
server_url: str,
timeout: int = 60,
max_retries: int = 3,
retry_delay: float = 1.0
):
"""
Инициализация клиента.
Аргументы:
server_url (str): базовый URL RAG-сервера (например, http://localhost:8080)
timeout (int): таймаут на запрос в секундах
max_retries (int): максимальное количество повторных попыток при ошибке
retry_delay (float): задержка между попытками (секунды)
"""
self.server_url = server_url.rstrip('/')
self.timeout = timeout
self.max_retries = max_retries
self.retry_delay = retry_delay
self._session: Optional[aiohttp.ClientSession] = None
logger.info(f"RAGClient инициализирован с сервером {server_url}")
async def _get_session(self) -> aiohttp.ClientSession:
"""
Возвращает сессию aiohttp (создаёт при первом вызове).
Сессия переиспользуется для всех запросов для эффективности.
"""
if self._session is None:
self._session = aiohttp.ClientSession(
timeout=aiohttp.ClientTimeout(total=self.timeout)
)
return self._session
async def close(self):
"""Закрывает HTTP-сессию."""
if self._session:
await self._session.close()
self._session = None
logger.debug("RAGClient сессия закрыта")
async def _post(self, endpoint: str, payload: Dict[str, Any]) -> Dict[str, Any]:
"""
Внутренний метод для отправки POST-запроса с повторными попытками.
Аргументы:
endpoint (str): путь эндпоинта (например, '/rag/query')
payload (Dict): тело запроса в формате JSON
Возвращает:
Dict: ответ сервера
Исключения:
aiohttp.ClientError: если все попытки не удались
"""
url = f"{self.server_url}{endpoint}"
session = await self._get_session()
for attempt in range(self.max_retries):
try:
async with session.post(url, json=payload) as resp:
if resp.status == 200:
return await resp.json()
else:
error_text = await resp.text()
logger.error(
f"RAG-сервер вернул {resp.status}: {error_text[:200]}"
)
raise aiohttp.ClientError(
f"HTTP {resp.status}: {error_text[:200]}"
)
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
logger.warning(
f"Попытка {attempt + 1}/{self.max_retries} к {endpoint} не удалась: {e}"
)
if attempt < self.max_retries - 1:
# Экспоненциальная задержка: 1, 2, 4, ... секунд
await asyncio.sleep(self.retry_delay * (2 ** attempt))
else:
logger.error(f"Все попытки к {endpoint} не удались")
raise
async def query(
self,
query: str,
user_jid: str,
room_jid: Optional[str],
prompts: Optional[Dict[str, str]] = None,
intent_override: Optional[str] = None,
last_file_path: Optional[str] = None,
last_file_text: Optional[str] = None,
) -> Dict[str, Any]:
"""
Выполняет RAG-запрос через HTTP-вызов /rag/query.
Сигнатура полностью совпадает с RAGAPI.query для бесшовной замены.
ИСТОРИЯ НЕ ПЕРЕДАЁТСЯ – сервер получает её из БД.
"""
payload = {
"query": query,
"user_jid": user_jid,
"room_jid": room_jid,
"prompts": prompts or {},
"intent_override": intent_override,
"last_file_path": last_file_path,
"last_file_text": last_file_text,
}
# Убираем None-значения, чтобы не засорять запрос
payload = {k: v for k, v in payload.items() if v is not None}
try:
result = await self._post("/rag/query", payload)
logger.debug(f"RAGClient.query успешно выполнен для {user_jid}")
return result
except Exception as e:
logger.error(f"Ошибка в RAGClient.query: {e}", exc_info=True)
return {
"answer": f"⚠️ Ошибка связи с RAG-сервером: {str(e)}",
"intent": "ERROR",
"context": "",
"sources": [],
"confidence": None,
"error": str(e)
}
async def index_document(
self,
file_name: str,
file_text: str,
user_jid: str,
room_jid: Optional[str],
is_global: bool = False,
title: Optional[str] = None,
metadata: Optional[Dict] = None,
file_hash: Optional[str] = None,
update_if_exists: bool = True
) -> Dict[str, Any]:
"""
Индексирует документ через HTTP-вызов /rag/index.
Сигнатура полностью совпадает с RAGAPI.index_document.
"""
payload = {
"file_name": file_name,
"file_text": file_text,
"user_jid": user_jid,
"room_jid": room_jid,
"is_global": is_global,
"title": title,
"metadata": metadata,
"file_hash": file_hash,
"update_if_exists": update_if_exists,
}
payload = {k: v for k, v in payload.items() if v is not None}
try:
result = await self._post("/rag/index", payload)
logger.info(f"Документ {file_name} проиндексирован через RAG-сервер")
return result
except Exception as e:
logger.error(f"Ошибка в RAGClient.index_document: {e}", exc_info=True)
return {"doc_id": None, "chunk_count": 0, "error": str(e)}