codekingpro/portable-devtools
114k
1import asyncio2import logging3import warnings4from concurrent.futures import Future, ThreadPoolExecutor5from typing import (6 Any,7 AsyncIterator,8 Dict,9 Iterator,10 List,11 Optional,12 Tuple,13 Union,14 cast,15)16 17import aiohttp18import requests19from langchain_core.documents import Document20 21from langchain_community.document_loaders.base import BaseLoader22from langchain_community.utils.user_agent import get_user_agent23 24logger = logging.getLogger(__name__)25 26default_header_template = {27 "User-Agent": get_user_agent(),28 "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*"29 ";q=0.8",30 "Accept-Language": "en-US,en;q=0.5",31 "Referer": "https://www.google.com/",32 "DNT": "1",33 "Connection": "keep-alive",34 "Upgrade-Insecure-Requests": "1",35}36 37 38def _build_metadata(soup: Any, url: str) -> dict:39 """Build metadata from BeautifulSoup output."""40 metadata = {"source": url}41 if title := soup.find("title"):42 metadata["title"] = title.get_text()43 if description := soup.find("meta", attrs={"name": "description"}):44 metadata["description"] = description.get("content", "No description found.")45 if html := soup.find("html"):46 metadata["language"] = html.get("lang", "No language found.")47 return metadata48 49 50class AsyncHtmlLoader(BaseLoader):51 """Load `HTML` asynchronously."""52 53 def __init__(54 self,55 web_path: Union[str, List[str]],56 header_template: Optional[dict] = None,57 verify_ssl: Optional[bool] = True,58 proxies: Optional[dict] = None,59 autoset_encoding: bool = True,60 encoding: Optional[str] = None,61 default_parser: str = "html.parser",62 requests_per_second: int = 2,63 requests_kwargs: Optional[Dict[str, Any]] = None,64 raise_for_status: bool = False,65 ignore_load_errors: bool = False,66 *,67 preserve_order: bool = True,68 trust_env: bool = False,69 ):70 """Initialize with a webpage path."""71 72 # TODO: Deprecate web_path in favor of web_paths, and remove this73 # left like this because there are a number of loaders that expect single74 # urls75 if isinstance(web_path, str):76 self.web_paths = [web_path]77 elif isinstance(web_path, List):78 self.web_paths = web_path79 80 headers = header_template or default_header_template81 if not headers.get("User-Agent"):82 try:83 from fake_useragent import UserAgent84 85 headers["User-Agent"] = UserAgent().random86 except ImportError:87 logger.info(88 "fake_useragent not found, using default user agent."89 "To get a realistic header for requests, "90 "`pip install fake_useragent`."91 )92 93 self.session = requests.Session()94 self.session.headers = dict(headers)95 self.session.verify = verify_ssl96 97 if proxies:98 self.session.proxies.update(proxies)99 100 self.requests_per_second = requests_per_second101 self.default_parser = default_parser102 self.requests_kwargs = requests_kwargs or {}103 self.raise_for_status = raise_for_status104 self.autoset_encoding = autoset_encoding105 self.encoding = encoding106 self.ignore_load_errors = ignore_load_errors107 self.preserve_order = preserve_order108 109 self.trust_env = trust_env110 111 def _fetch_valid_connection_docs(self, url: str) -> Any:112 if self.ignore_load_errors:113 try:114 return self.session.get(url, **self.requests_kwargs)115 except Exception as e:116 warnings.warn(str(e))117 return None118 119 return self.session.get(url, **self.requests_kwargs)120 121 @staticmethod122 def _check_parser(parser: str) -> None:123 """Check that parser is valid for bs4."""124 valid_parsers = ["html.parser", "lxml", "xml", "lxml-xml", "html5lib"]125 if parser not in valid_parsers:126 raise ValueError(127 "`parser` must be one of " + ", ".join(valid_parsers) + "."128 )129 130 async def _fetch(131 self, url: str, retries: int = 3, cooldown: int = 2, backoff: float = 1.5132 ) -> str:133 async with aiohttp.ClientSession(trust_env=self.trust_env) as session:134 for i in range(retries):135 try:136 kwargs: Dict = dict(137 headers=self.session.headers,138 cookies=self.session.cookies.get_dict(),139 **self.requests_kwargs,140 )141 if not self.session.verify:142 kwargs["ssl"] = False143 async with session.get(144 url,145 **kwargs,146 ) as response:147 try:148 text = await response.text()149 except UnicodeDecodeError:150 logger.error(f"Failed to decode content from {url}")151 text = ""152 return text153 except (aiohttp.ClientConnectionError, TimeoutError) as e:154 if i == retries - 1 and self.ignore_load_errors:155 logger.warning(f"Error fetching {url} after {retries} retries.")156 return ""157 elif i == retries - 1:158 raise159 else:160 logger.warning(161 f"Error fetching {url} with attempt "162 f"{i + 1}/{retries}: {e}. Retrying..."163 )164 await asyncio.sleep(cooldown * backoff**i)165 raise ValueError("retry count exceeded")166 167 async def _fetch_with_rate_limit(168 self, url: str, semaphore: asyncio.Semaphore169 ) -> Tuple[str, str]:170 async with semaphore:171 return url, await self._fetch(url)172 173 async def _lazy_fetch_all(174 self, urls: List[str], preserve_order: bool175 ) -> AsyncIterator[Tuple[str, str]]:176 semaphore = asyncio.Semaphore(self.requests_per_second)177 tasks = [178 asyncio.create_task(self._fetch_with_rate_limit(url, semaphore))179 for url in urls180 ]181 try:182 from tqdm.asyncio import tqdm_asyncio183 184 if preserve_order:185 for task in tqdm_asyncio(186 tasks, desc="Fetching pages", ascii=True, mininterval=1187 ):188 yield await task189 else:190 for task in tqdm_asyncio.as_completed(191 tasks, desc="Fetching pages", ascii=True, mininterval=1192 ):193 yield await task194 except ImportError:195 warnings.warn("For better logging of progress, `pip install tqdm`")196 if preserve_order:197 for result in await asyncio.gather(*tasks):198 yield result199 else:200 for task in asyncio.as_completed(tasks):201 yield await task202 203 async def fetch_all(self, urls: List[str]) -> List[str]:204 """Fetch all urls concurrently with rate limiting."""205 return [doc async for _, doc in self._lazy_fetch_all(urls, True)]206 207 def _to_document(self, url: str, text: str) -> Document:208 from bs4 import BeautifulSoup209 210 if url.endswith(".xml"):211 parser = "xml"212 else:213 parser = self.default_parser214 self._check_parser(parser)215 soup = BeautifulSoup(text, parser)216 metadata = _build_metadata(soup, url)217 return Document(page_content=text, metadata=metadata)218 219 def lazy_load(self) -> Iterator[Document]:220 """Lazy load text from the url(s) in web_path."""221 results: List[str]222 try:223 # Raises RuntimeError if there is no current event loop.224 asyncio.get_running_loop()225 # If there is a current event loop, we need to run the async code226 # in a separate loop, in a separate thread.227 with ThreadPoolExecutor(max_workers=1) as executor:228 future: Future[List[str]] = executor.submit(229 asyncio.run,230 self.fetch_all(self.web_paths),231 )232 results = future.result()233 except RuntimeError:234 results = asyncio.run(self.fetch_all(self.web_paths))235 236 for i, text in enumerate(cast(List[str], results)):237 yield self._to_document(self.web_paths[i], text)238 239 async def alazy_load(self) -> AsyncIterator[Document]:240 """Lazy load text from the url(s) in web_path."""241 async for url, text in self._lazy_fetch_all(242 self.web_paths, self.preserve_order243 ):244 yield self._to_document(url, text)245 