from __future__ import annotations import time from dataclasses import dataclass from typing import Any import requests class LightRAGCircuitOpen(RuntimeError): pass @dataclass(frozen=True) class ProjectionReceipt: track_id: str @dataclass(frozen=True) class DeleteReceipt: verified: bool class _RequestsTransport: def request(self, **kwargs): return requests.request(**kwargs) class LightRAGClient: def __init__( self, *, base_url: str, api_key: str, workspace: str, transport=None, timeout_seconds: float = 15, failure_threshold: int = 3, cooldown_seconds: float = 30, ) -> None: if not base_url or not api_key or not workspace: raise ValueError("LightRAG base URL, API key and workspace are required") self._base_url = base_url.rstrip("/") self._api_key = api_key self.workspace = workspace self._transport = transport or _RequestsTransport() self._timeout = timeout_seconds self._failure_threshold = failure_threshold self._cooldown = cooldown_seconds self._failures = 0 self._opened_at: float | None = None def _request( self, method: str, path: str, *, json: dict[str, Any] | None = None, idempotency_key: str | None = None, ) -> dict[str, Any]: now = time.monotonic() if self._opened_at is not None: if now - self._opened_at < self._cooldown: raise LightRAGCircuitOpen("LightRAG circuit breaker is open") self._opened_at = None self._failures = 0 headers = {"X-API-Key": self._api_key} if idempotency_key: headers["Idempotency-Key"] = idempotency_key try: response = self._transport.request( method=method, url=f"{self._base_url}{path}", json=json, headers=headers, timeout=self._timeout, ) response.raise_for_status() payload = response.json() except Exception: self._failures += 1 if self._failures >= self._failure_threshold: self._opened_at = time.monotonic() raise self._failures = 0 return payload if isinstance(payload, dict) else {"response": payload} def health(self) -> dict[str, Any]: return self._request("GET", "/health") def insert( self, *, external_document_id: str, content: str, metadata: dict[str, Any], ) -> ProjectionReceipt: payload = self._request( "POST", "/documents/text", json={ "text": f"[DATAOPS_SOURCE {external_document_id}]\n{content}", "file_source": external_document_id, "metadata": metadata, }, idempotency_key=external_document_id, ) track_id = payload.get("track_id") or payload.get("data", {}).get("track_id") if not track_id: raise RuntimeError("LightRAG insert response has no track_id") return ProjectionReceipt(track_id=str(track_id)) def track_status(self, track_id: str) -> str: payload = self._request("GET", f"/documents/track_status/{track_id}") value = payload.get("status") or payload.get("data", {}).get("status") return str(value or "unknown") def delete(self, external_document_id: str) -> DeleteReceipt: payload = self._request( "DELETE", "/documents/delete_document", json={"doc_id": external_document_id}, idempotency_key=f"delete:{external_document_id}", ) return DeleteReceipt(verified=payload.get("verified") is True) def query_context(self, query: str, *, mode: str = "mix", limit: int = 20) -> str: if mode not in {"mix", "hybrid", "local", "global", "naive"}: raise ValueError("unsupported LightRAG query mode") payload = self._request( "POST", "/query", json={ "query": query, "mode": mode, "only_need_context": True, "top_k": min(max(limit, 1), 100), }, ) value = payload.get("response", payload.get("data", "")) return str(value)