/
dictator
/
ragflow_sync
Обзор
Документация
Войти
/
dictator
/
ragflow_sync
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
ragflow_uploader/client.py
509 строк
14 KB
Developer
fix: use body-based DELETE /api/v1/datasets instead of path-based
14 июл 2026, 13:37
14 июл 2026, 13:37
d34f967
Код
Авторство
О чём код?
""" RAGFlow HTTP API client wrapper. Provides a Pythonic interface to RAGFlow's REST API for dataset and document management. """ import logging from typing import Any, Optional import httpx from .config import RAGFlowConfig logger = logging.getLogger(__name__) class RAGFlowError(Exception): """Exception raised for RAGFlow API errors.""" def __init__(self, message: str, status_code: Optional[int] = None): self.message = message self.status_code = status_code super().__init__(self.message) class RAGFlowClient: """Client for RAGFlow HTTP API.""" def __init__(self, config: RAGFlowConfig, verify_ssl: bool = True): """ Initialize RAGFlow client. Args: config: RAGFlow configuration with API key and base URL. verify_ssl: Whether to verify SSL certificates (default: True). Set to False for self-signed certificates. """ self.base_url = config.base_url.rstrip("/") self.api_key = config.api_key self.timeout = config.timeout self.verify_ssl = verify_ssl self._client = httpx.Client( base_url=self.base_url, timeout=self.timeout, headers={ "Authorization": f"Bearer {self.api_key}", }, verify=verify_ssl, ) def _request( self, method: str, path: str, **kwargs: Any, ) -> dict[str, Any]: """ Make HTTP request to RAGFlow API. Args: method: HTTP method (GET, POST, PUT, DELETE) path: API path (e.g., "/api/v1/datasets") **kwargs: Additional arguments passed to httpx Returns: Parsed JSON response Raises: RAGFlowError: If API request fails """ url = f"/api/v1{path}" if not path.startswith("/api") else path # Set Content-Type to application/json for JSON requests (not for file uploads) if "files" not in kwargs and "json" in kwargs: kwargs["headers"] = kwargs.get("headers", {}) kwargs["headers"]["Content-Type"] = "application/json" try: response = self._client.request(method, url, **kwargs) response.raise_for_status() # Handle empty responses if not response.text: return {} data = response.json() if data.get("code", 0) != 0: raise RAGFlowError( message=data.get("message", "Unknown error"), status_code=response.status_code, ) return data.get("data", {}) except httpx.HTTPStatusError as e: logger.error(f"HTTP error: {e.response.status_code} - {e.response.text}") try: error_data = e.response.json() raise RAGFlowError( message=error_data.get("message", str(e)), status_code=e.response.status_code, ) except (ValueError, Exception): raise RAGFlowError( message=str(e), status_code=e.response.status_code, ) except Exception as e: logger.error(f"Request error: {e}") raise RAGFlowError(message=str(e)) # Dataset operations def create_dataset( self, name: str, description: str = "", embedding_model: str = "BAAI/bge-large-zh-v1.5@BAAI", permission: str = "me", ) -> dict[str, Any]: """ Create a new dataset. Args: name: Dataset name description: Dataset description embedding_model: Embedding model to use permission: Permission level ("me" or "team") Returns: Dataset information including ID """ payload = { "name": name, "description": description, "embedding_model": embedding_model, "permission": permission, } logger.info(f"Creating dataset: {name}") return self._request("POST", "/datasets", json=payload) def list_datasets( self, page: int = 1, page_size: int = 100, name: Optional[str] = None, ) -> list[dict[str, Any]]: """ List all datasets. Args: page: Page number (1-indexed) page_size: Number of items per page name: Filter by dataset name (partial match) Returns: List of dataset information """ params = {"page": page, "page_size": page_size} if name: params["name"] = name logger.debug(f"Listing datasets (page={page}, page_size={page_size})") data = self._request("GET", "/datasets", params=params) return data.get("datasets", []) if isinstance(data, dict) else data def get_dataset(self, dataset_id: str) -> dict[str, Any]: """ Get dataset information. Args: dataset_id: Dataset ID Returns: Dataset information """ logger.debug(f"Getting dataset: {dataset_id}") return self._request("GET", f"/datasets/{dataset_id}") def delete_dataset(self, dataset_id: str) -> bool: """ Delete a dataset. Uses body-based DELETE /api/v1/datasets with {"ids": [...]} (required by RAGFlow v0.26.1+ — path-based DELETE returns 405). Args: dataset_id: Dataset ID Returns: True if successful """ logger.info(f"Deleting dataset: {dataset_id}") self._request("DELETE", "/datasets", json={"ids": [dataset_id]}) return True # Document operations def upload_document( self, dataset_id: str, file_name: str, content: bytes, ) -> dict[str, Any]: """ Upload a document to a dataset. Args: dataset_id: Dataset ID file_name: Name of the file content: File content as bytes Returns: Document information including ID """ files = {"file": (file_name, content, "text/markdown")} logger.info(f"Uploading document: {file_name} to dataset {dataset_id}") return self._request( "POST", f"/datasets/{dataset_id}/documents", files=files, ) def upload_document_text( self, dataset_id: str, file_name: str, text: str, ) -> dict[str, Any]: """ Upload a text document to a dataset. Args: dataset_id: Dataset ID file_name: Name of the file text: Document content as text Returns: Document information including ID """ return self.upload_document( dataset_id=dataset_id, file_name=file_name, content=text.encode("utf-8"), ) def list_documents( self, dataset_id: str, page: int = 1, page_size: int = 100, keywords: Optional[str] = None, ) -> list[dict[str, Any]]: """ List documents in a dataset. Args: dataset_id: Dataset ID page: Page number (1-indexed) page_size: Number of items per page keywords: Filter by keywords Returns: List of document information """ params = {"page": page, "page_size": page_size} if keywords: params["keywords"] = keywords logger.debug( f"Listing documents in dataset {dataset_id} (page={page}, page_size={page_size})" ) data = self._request( "GET", f"/datasets/{dataset_id}/documents", params=params, ) return ( data.get("documents", []) or data.get("docs", []) if isinstance(data, dict) else data ) def get_document(self, dataset_id: str, document_id: str) -> dict[str, Any]: """ Get document information. Args: dataset_id: Dataset ID document_id: Document ID Returns: Document information """ logger.debug(f"Getting document: {document_id}") try: return self._request("GET", f"/datasets/{dataset_id}/documents/{document_id}") except RAGFlowError: return {} def delete_document(self, dataset_id: str, document_id: str) -> bool: """ Delete a document from a dataset. Args: dataset_id: Dataset ID document_id: Document ID Returns: True if successful """ logger.info(f"Deleting document: {document_id} from dataset {dataset_id}") self._request("DELETE", f"/datasets/{dataset_id}/documents/{document_id}") return True def update_document( self, dataset_id: str, document_id: str, **kwargs: Any, ) -> dict[str, Any]: """ Update document metadata or configuration. Args: dataset_id: Dataset ID document_id: Document ID **kwargs: Fields to update (e.g., chunk_method, parser_config) Returns: Updated document information """ logger.info(f"Updating document: {document_id}") return self._request( "PUT", f"/datasets/{dataset_id}/documents/{document_id}", json=kwargs, ) def download_document( self, dataset_id: str, document_id: str, ) -> bytes: """ Download document content. Args: dataset_id: Dataset ID document_id: Document ID Returns: Document content as bytes """ logger.debug(f"Downloading document: {document_id}") response = self._client.get( f"/api/v1/datasets/{dataset_id}/documents/{document_id}/download" ) response.raise_for_status() return response.content # Retrieval operations def retrieve( self, dataset_ids: list[str], query: str, top_k: int = 5, document_ids: Optional[list[str]] = None, ) -> list[dict[str, Any]]: """ Retrieve relevant chunks from datasets. Args: dataset_ids: List of dataset IDs to search query: Search query top_k: Number of results to return document_ids: Optional list of document IDs to filter Returns: List of retrieved chunks with scores """ payload = { "dataset_ids": dataset_ids, "query": query, "top_k": top_k, } if document_ids: payload["document_ids"] = document_ids logger.info(f"Retrieving chunks for query: {query[:50]}...") data = self._request("POST", "/retrieval", json=payload) return data.get("chunks", []) if isinstance(data, dict) else data def parse_documents( self, dataset_id: str, document_ids: list[str], ) -> dict[str, Any]: """ Trigger document parsing for specified documents. Args: dataset_id: Dataset ID document_ids: List of document IDs to parse Returns: Dictionary with parsing results for each document """ payload = {"document_ids": document_ids} logger.info(f"Triggering parsing for {len(document_ids)} documents") try: return self._request("POST", f"/datasets/{dataset_id}/chunks", json=payload) except Exception as e: logger.error(f"Parse documents error: {e}") raise def get_document_status(self, dataset_id: str, document_id: str) -> dict[str, Any]: """ Get document status including parsing progress. Args: dataset_id: Dataset ID document_id: Document ID Returns: Document information including progress and run status """ # Use list instead of get since GET /documents/{id} returns empty body docs = self.list_documents(dataset_id=dataset_id, page=1, page_size=100) for doc in docs: if doc.get("id") == document_id: return doc return {} def wait_for_parsing( self, dataset_id: str, document_ids: list[str], timeout: int = 300, poll_interval: int = 2, ) -> dict[str, Any]: """ Wait for document parsing to complete. Args: dataset_id: Dataset ID document_ids: List of document IDs to wait for timeout: Maximum time to wait in seconds poll_interval: Polling interval in seconds Returns: Dictionary with completed, failed, and timed_out document IDs """ import time result = { "completed": [], "failed": [], "timed_out": [], } pending = list(document_ids) start_time = time.time() while pending: elapsed = time.time() - start_time if elapsed >= timeout: result["timed_out"] = pending logger.warning(f"Parsing timeout after {timeout}s for {len(pending)} documents") break for doc_id in pending[:]: doc = self.get_document(dataset_id, doc_id) status = doc.get("run", "") if status == "DONE": result["completed"].append(doc_id) pending.remove(doc_id) logger.debug(f"Document {doc_id} parsing completed") elif status == "FAIL" or status == "4": result["failed"].append(doc_id) pending.remove(doc_id) logger.warning(f"Document {doc_id} parsing failed") if pending: time.sleep(poll_interval) logger.info( f"Parsing finished: {len(result['completed'])} completed, " f"{len(result['failed'])} failed, " f"{len(result['timed_out'])} timed out" ) return result def close(self) -> None: """Close the HTTP client.""" self._client.close() def __enter__(self) -> "RAGFlowClient": return self def __exit__(self, exc_type, exc_val, exc_tb) -> None: self.close()