Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
lakefs.py184 linesDownload Raw Back to document_loaders
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 
codekingpro/portable-devtools · Team Ai