codekingpro/portable-devtools
114k
1"""Pebblo's safe dataloader is a wrapper for document loaders"""2 3import logging4import os5import uuid6from importlib.metadata import version7from typing import Any, Dict, Iterable, Iterator, List, Optional8 9from langchain_core.documents import Document10 11from langchain_community.document_loaders.base import BaseLoader12from langchain_community.utilities.pebblo import (13 BATCH_SIZE_BYTES,14 PLUGIN_VERSION,15 App,16 Framework,17 IndexedDocument,18 PebbloLoaderAPIWrapper,19 generate_size_based_batches,20 get_full_path,21 get_loader_full_path,22 get_loader_type,23 get_runtime,24 get_source_size,25)26 27logger = logging.getLogger(__name__)28 29 30class PebbloSafeLoader(BaseLoader):31 """Pebblo Safe Loader class is a wrapper around document loaders enabling the data32 to be scrutinized.33 """34 35 _discover_sent: bool = False36 37 def __init__(38 self,39 langchain_loader: BaseLoader,40 name: str,41 owner: str = "",42 description: str = "",43 api_key: Optional[str] = None,44 load_semantic: bool = False,45 classifier_url: Optional[str] = None,46 *,47 classifier_location: str = "local",48 anonymize_snippets: bool = False,49 ):50 if not name or not isinstance(name, str):51 raise NameError("Must specify a valid name.")52 self.app_name = name53 self.load_id = str(uuid.uuid4())54 self.loader = langchain_loader55 self.load_semantic = os.environ.get("PEBBLO_LOAD_SEMANTIC") or load_semantic56 self.owner = owner57 self.description = description58 self.source_path = get_loader_full_path(self.loader)59 self.docs: List[Document] = []60 self.docs_with_id: List[IndexedDocument] = []61 loader_name = str(type(self.loader)).split(".")[-1].split("'")[0]62 self.source_type = get_loader_type(loader_name)63 self.source_path_size = get_source_size(self.source_path)64 self.batch_size = BATCH_SIZE_BYTES65 self.loader_details = {66 "loader": loader_name,67 "source_path": self.source_path,68 "source_type": self.source_type,69 **(70 {"source_path_size": str(self.source_path_size)}71 if self.source_path_size > 072 else {}73 ),74 }75 # generate app76 self.app = self._get_app_details()77 # initialize Pebblo Loader API client78 self.pb_client = PebbloLoaderAPIWrapper(79 api_key=api_key,80 classifier_location=classifier_location,81 classifier_url=classifier_url,82 anonymize_snippets=anonymize_snippets,83 )84 self.pb_client.send_loader_discover(self.app)85 86 def load(self) -> List[Document]:87 """Load Documents.88 89 Returns:90 list: Documents fetched from load method of the wrapped `loader`.91 """92 self.docs = self.loader.load()93 # Classify docs in batches94 self.classify_in_batches()95 return self.docs96 97 def classify_in_batches(self) -> None:98 """99 Classify documents in batches.100 This is to avoid API timeouts when sending large number of documents.101 Batches are generated based on the page_content size.102 """103 batches: List[List[Document]] = generate_size_based_batches(104 self.docs, self.batch_size105 )106 107 processed_docs: List[Document] = []108 109 total_batches = len(batches)110 for i, batch in enumerate(batches):111 is_last_batch: bool = i == total_batches - 1112 self.docs = batch113 self.docs_with_id = self._index_docs()114 classified_docs = self.pb_client.classify_documents(115 self.docs_with_id,116 self.app,117 self.loader_details,118 loading_end=is_last_batch,119 )120 self._add_pebblo_specific_metadata(classified_docs)121 if self.load_semantic:122 batch_processed_docs = self._add_semantic_to_docs(classified_docs)123 else:124 batch_processed_docs = self._unindex_docs()125 processed_docs.extend(batch_processed_docs)126 127 self.docs = processed_docs128 129 def lazy_load(self) -> Iterator[Document]:130 """Load documents in lazy fashion.131 132 Raises:133 NotImplementedError: raised when lazy_load id not implemented134 within wrapped loader.135 136 Yields:137 list: Documents from loader's lazy loading.138 """139 try:140 doc_iterator = self.loader.lazy_load()141 except NotImplementedError as exc:142 err_str = f"{self.loader.__class__.__name__} does not implement lazy_load()"143 logger.error(err_str)144 raise NotImplementedError(err_str) from exc145 while True:146 try:147 doc = next(doc_iterator)148 except StopIteration:149 self.docs = []150 break151 self.docs = list((doc,))152 self.docs_with_id = self._index_docs()153 classified_doc = self.pb_client.classify_documents(154 self.docs_with_id, self.app, self.loader_details155 )156 self._add_pebblo_specific_metadata(classified_doc)157 if self.load_semantic:158 self.docs = self._add_semantic_to_docs(classified_doc)159 else:160 self.docs = self._unindex_docs()161 yield self.docs[0]162 163 @classmethod164 def set_discover_sent(cls) -> None:165 cls._discover_sent = True166 167 def _get_app_details(self) -> App:168 """Fetch app details. Internal method.169 170 Returns:171 App: App details.172 """173 framework, runtime = get_runtime()174 app = App(175 name=self.app_name,176 owner=self.owner,177 description=self.description,178 load_id=self.load_id,179 runtime=runtime,180 framework=framework,181 plugin_version=PLUGIN_VERSION,182 client_version=Framework(183 name="langchain_community",184 version=version("langchain_community"),185 ),186 )187 return app188 189 def _index_docs(self) -> List[IndexedDocument]:190 """191 Indexes the documents and returns a list of IndexedDocument objects.192 193 Returns:194 List[IndexedDocument]: A list of IndexedDocument objects with unique IDs.195 """196 docs_with_id = [197 IndexedDocument(pb_id=str(i), **doc.dict())198 for i, doc in enumerate(self.docs)199 ]200 return docs_with_id201 202 def _add_semantic_to_docs(self, classified_docs: Dict) -> List[Document]:203 """204 Adds semantic metadata to the given list of documents.205 206 Args:207 classified_docs (Dict): A dictionary of dictionaries containing the208 classified documents with pb_id as key.209 210 Returns:211 A list of `Document` objects with added semantic metadata.212 """213 indexed_docs = {214 doc.pb_id: Document(page_content=doc.page_content, metadata=doc.metadata)215 for doc in self.docs_with_id216 }217 218 for classified_doc in classified_docs.values():219 doc_id = classified_doc.get("pb_id")220 if doc_id in indexed_docs:221 self._add_semantic_to_doc(indexed_docs[doc_id], classified_doc)222 223 semantic_metadata_docs = [doc for doc in indexed_docs.values()]224 225 return semantic_metadata_docs226 227 def _unindex_docs(self) -> List[Document]:228 """229 Converts a list of `IndexedDocument` objects to a list of `Document` objects.230 231 Returns:232 A list of `Document` objects.233 """234 docs = [235 Document(page_content=doc.page_content, metadata=doc.metadata)236 for i, doc in enumerate(self.docs_with_id)237 ]238 return docs239 240 def _add_semantic_to_doc(self, doc: Document, classified_doc: dict) -> Document:241 """242 Adds semantic metadata to the given document in-place.243 244 Args:245 doc (Document): A Document object.246 classified_doc: `dict` containing the classified document.247 248 Returns:249 Document: The Document object with added semantic metadata.250 """251 doc.metadata["pebblo_semantic_entities"] = list(252 classified_doc.get("entities", {}).keys()253 )254 doc.metadata["pebblo_semantic_topics"] = list(255 classified_doc.get("topics", {}).keys()256 )257 return doc258 259 def _add_pebblo_specific_metadata(self, classified_docs: dict) -> None:260 """Add Pebblo specific metadata to documents."""261 for doc in self.docs_with_id:262 doc_metadata = doc.metadata263 if self.loader.__class__.__name__ == "SharePointLoader":264 doc_metadata["full_path"] = get_full_path(265 doc_metadata.get("source", self.source_path)266 )267 else:268 doc_metadata["full_path"] = get_full_path(269 doc_metadata.get(270 "full_path", doc_metadata.get("source", self.source_path)271 )272 )273 doc_metadata["pb_checksum"] = classified_docs.get(doc.pb_id, {}).get(274 "pb_checksum", None275 )276 277 278class PebbloTextLoader(BaseLoader):279 """280 Loader for text data.281 282 Since PebbloSafeLoader is a wrapper around document loaders, this loader is283 used to load text data directly into Documents.284 """285 286 def __init__(287 self,288 texts: Iterable[str],289 *,290 source: Optional[str] = None,291 ids: Optional[List[str]] = None,292 metadata: Optional[Dict[str, Any]] = None,293 metadatas: Optional[List[Dict[str, Any]]] = None,294 ) -> None:295 """296 Args:297 texts: Iterable of text data.298 source: Source of the text data.299 Optional. Defaults to None.300 ids: List of unique identifiers for each text.301 Optional. Defaults to None.302 metadata: Metadata for all texts.303 Optional. Defaults to None.304 metadatas: List of metadata for each text.305 Optional. Defaults to None.306 """307 self.texts = texts308 self.source = source309 self.ids = ids310 self.metadata = metadata311 self.metadatas = metadatas312 313 def lazy_load(self) -> Iterator[Document]:314 """315 Lazy load text data into Documents.316 317 Returns:318 Iterator of Documents319 """320 for i, text in enumerate(self.texts):321 _id = None322 metadata = self.metadata or {}323 if self.metadatas and i < len(self.metadatas) and self.metadatas[i]:324 metadata.update(self.metadatas[i])325 if self.ids and i < len(self.ids):326 _id = self.ids[i]327 yield Document(id=_id, page_content=text, metadata=metadata)328 329 def load(self) -> List[Document]:330 """331 Load text data into Documents.332 333 Returns:334 List of Documents335 """336 documents = []337 for doc in self.lazy_load():338 documents.append(doc)339 return documents340 