/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/services/rag_client.py
173 строки
5 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""HTTP клиент для RAG сервиса.""" from __future__ import annotations from dataclasses import dataclass from typing import Any import httpx import structlog from src.config import get_settings logger = structlog.get_logger() @dataclass class RAGSource: """Источник из RAG поиска.""" chunk_id: str document_id: str content: str score: float rerank_score: float | None metadata: dict[str, Any] def to_dict(self) -> dict[str, Any]: return { "chunk_id": self.chunk_id, "document_id": self.document_id, "content": self.content, "score": self.score, "rerank_score": self.rerank_score, "document_title": self.metadata.get("title", ""), "source": self.metadata.get("source", ""), } @dataclass class RAGSearchResult: """Результат RAG поиска.""" query: str sources: list[RAGSource] context: str results_count: int def format_context_for_llm(self, max_sources: int = 5) -> str: """Форматировать контекст для вставки в LLM prompt.""" if not self.sources: return "" parts = [] for i, src in enumerate(self.sources[:max_sources], 1): title = src.metadata.get("title", "Unknown") parts.append(f"[Source {i}: {title}]\n{src.content}") return "\n\n".join(parts) class RAGClient: """HTTP клиент для RAG сервиса.""" def __init__(self, base_url: str | None = None) -> None: settings = get_settings() self.base_url = ( base_url or getattr(settings, "rag_service_url", None) or "http://rag:8000" ).rstrip("/") # ✅ Детализированные таймауты: connect 10s, остальные 30s self._client = httpx.AsyncClient( timeout=httpx.Timeout(30.0, connect=10.0), verify=False, # Cloud.ru corporate proxy ) async def search( self, query: str, workspace_id: str, top_k: int = 5, score_threshold: float | None = None, ) -> RAGSearchResult: """Поиск в RAG (graceful degradation при ошибке).""" url = f"{self.base_url}/search" payload: dict[str, Any] = { "query": query, "workspace_id": workspace_id, "top_k": top_k, } if score_threshold is not None: payload["score_threshold"] = score_threshold logger.info( "rag.search.request", url=url, query=query[:50], workspace_id=workspace_id, ) try: response = await self._client.post(url, json=payload) response.raise_for_status() data = response.json() sources = [ RAGSource( chunk_id=r.get("chunk_id", ""), document_id=r.get("document_id", ""), content=r.get("content", ""), score=r.get("score", 0.0), rerank_score=r.get("rerank_score"), metadata=r.get("metadata", {}), ) for r in data.get("results", []) ] result = RAGSearchResult( query=query, sources=sources, context=data.get("context", ""), results_count=len(sources), ) logger.info( "rag.search.success", results_count=result.results_count, top_score=sources[0].score if sources else 0.0, ) return result except httpx.HTTPStatusError as e: logger.error( "rag.search.http_error", status_code=e.response.status_code, response=e.response.text[:200], ) return RAGSearchResult(query=query, sources=[], context="", results_count=0) except Exception as e: logger.error("rag.search.error", error=str(e)) return RAGSearchResult(query=query, sources=[], context="", results_count=0) async def health_check(self) -> bool: """Проверить доступность RAG.""" try: response = await self._client.get(f"{self.base_url}/health") return response.status_code == 200 except Exception: return False async def aclose(self) -> None: """Закрыть HTTP клиент (вызывать при shutdown).""" await self._client.aclose() _rag_client: RAGClient | None = None def get_rag_client() -> RAGClient: """Получить singleton RAG клиент.""" global _rag_client if _rag_client is None: _rag_client = RAGClient() return _rag_client async def close_rag_client() -> None: """Закрыть singleton RAG клиент (для lifespan shutdown).""" global _rag_client if _rag_client is not None: await _rag_client.aclose() _rag_client = None