| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141 |
- 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)
|