codekingpro/portable-devtools
114k
1import os2import tempfile3import urllib.parse4from typing import Any, List, Optional5from urllib.parse import urljoin6 7import requests8from langchain_core.documents import Document9from requests.auth import HTTPBasicAuth10 11from langchain_community.document_loaders.base import BaseLoader12from langchain_community.document_loaders.unstructured import UnstructuredBaseLoader13 14 15class LakeFSClient:16 """Client for lakeFS."""17 18 def __init__(19 self,20 lakefs_access_key: str,21 lakefs_secret_key: str,22 lakefs_endpoint: str,23 ):24 self.__endpoint = "/".join([lakefs_endpoint, "api", "v1/"])25 self.__auth = HTTPBasicAuth(lakefs_access_key, lakefs_secret_key)26 try:27 health_check = requests.get(28 urljoin(self.__endpoint, "healthcheck"), auth=self.__auth29 )30 health_check.raise_for_status()31 except Exception:32 raise ValueError(33 "lakeFS server isn't accessible. Make sure lakeFS is running."34 )35 36 def ls_objects(37 self, repo: str, ref: str, path: str, presign: Optional[bool]38 ) -> List:39 qp = {"prefix": path, "presign": presign}40 eqp = urllib.parse.urlencode(qp)41 objects_ls_endpoint = urljoin(42 self.__endpoint, f"repositories/{repo}/refs/{ref}/objects/ls?{eqp}"43 )44 olsr = requests.get(objects_ls_endpoint, auth=self.__auth)45 olsr.raise_for_status()46 olsr_json = olsr.json()47 return list(48 map(49 lambda res: (res["path"], res["physical_address"]), olsr_json["results"]50 )51 )52 53 def is_presign_supported(self) -> bool:54 config_endpoint = self.__endpoint + "config"55 response = requests.get(config_endpoint, auth=self.__auth)56 response.raise_for_status()57 config = response.json()58 return config["storage_config"]["pre_sign_support"]59 60 61class LakeFSLoader(BaseLoader):62 """Load from `lakeFS`."""63 64 repo: str65 ref: str66 path: str67 68 def __init__(69 self,70 lakefs_access_key: str,71 lakefs_secret_key: str,72 lakefs_endpoint: str,73 repo: Optional[str] = None,74 ref: Optional[str] = "main",75 path: Optional[str] = "",76 ):77 """78 79 :param lakefs_access_key: [required] lakeFS server's access key80 :param lakefs_secret_key: [required] lakeFS server's secret key81 :param lakefs_endpoint: [required] lakeFS server's endpoint address,82 ex: https://example.my-lakefs.com83 :param repo: [optional, default = ''] target repository84 :param ref: [optional, default = 'main'] target ref (branch name,85 tag, or commit ID)86 :param path: [optional, default = ''] target path87 """88 89 self.__lakefs_client = LakeFSClient(90 lakefs_access_key, lakefs_secret_key, lakefs_endpoint91 )92 self.repo = "" if repo is None or repo == "" else str(repo)93 self.ref = "main" if ref is None or ref == "" else str(ref)94 self.path = "" if path is None else str(path)95 96 def set_path(self, path: str) -> None:97 self.path = path98 99 def set_ref(self, ref: str) -> None:100 self.ref = ref101 102 def set_repo(self, repo: str) -> None:103 self.repo = repo104 105 def load(self) -> List[Document]:106 self.__validate_instance()107 presigned = self.__lakefs_client.is_presign_supported()108 docs: List[Document] = []109 objs = self.__lakefs_client.ls_objects(110 repo=self.repo, ref=self.ref, path=self.path, presign=presigned111 )112 for obj in objs:113 lakefs_unstructured_loader = UnstructuredLakeFSLoader(114 obj[1], self.repo, self.ref, obj[0], presigned115 )116 docs.extend(lakefs_unstructured_loader.load())117 return docs118 119 def __validate_instance(self) -> None:120 if self.repo is None or self.repo == "":121 raise ValueError(122 "no repository was provided. use `set_repo` to specify a repository"123 )124 if self.ref is None or self.ref == "":125 raise ValueError("no ref was provided. use `set_ref` to specify a ref")126 if self.path is None:127 raise ValueError("no path was provided. use `set_path` to specify a path")128 129 130class UnstructuredLakeFSLoader(UnstructuredBaseLoader):131 """Load from `lakeFS` as unstructured data."""132 133 def __init__(134 self,135 url: str,136 repo: str,137 ref: str = "main",138 path: str = "",139 presign: bool = True,140 **unstructured_kwargs: Any,141 ):142 """Initialize UnstructuredLakeFSLoader.143 144 Args:145 146 :param lakefs_access_key:147 :param lakefs_secret_key:148 :param lakefs_endpoint:149 :param repo:150 :param ref:151 """152 153 super().__init__(**unstructured_kwargs)154 self.url = url155 self.repo = repo156 self.ref = ref157 self.path = path158 self.presign = presign159 160 def _get_metadata(self) -> dict:161 return {"repo": self.repo, "ref": self.ref, "path": self.path}162 163 def _get_elements(self) -> List:164 from unstructured.partition.auto import partition165 166 local_prefix = "local://"167 168 if self.presign:169 with tempfile.TemporaryDirectory() as temp_dir:170 file_path = f"{temp_dir}/{self.path.split('/')[-1]}"171 os.makedirs(os.path.dirname(file_path), exist_ok=True)172 response = requests.get(self.url)173 response.raise_for_status()174 with open(file_path, mode="wb") as file:175 file.write(response.content)176 return partition(filename=file_path)177 elif not self.url.startswith(local_prefix):178 raise ValueError(179 "Non pre-signed URLs are supported only with 'local' blockstore"180 )181 else:182 local_path = self.url[len(local_prefix) :]183 return partition(filename=local_path)184 