1
import mimetypes
2
import os
3
import asyncio
4
-import aiohttp
4
import json
5
6
from helpers.vector_db import VectorDB
12
from typing import Callable, Sequence, List, Optional, Tuple
13
from datetime import datetime
14
16
-from langchain_community.document_loaders import AsyncHtmlLoader
17
-from langchain_community.document_loaders.text import TextLoader
15
from langchain_community.document_loaders.pdf import PyMuPDFLoader
16
from langchain_community.document_transformers import MarkdownifyTransformer
17
from langchain_community.document_loaders.parsers.images import TesseractBlobParser
21
22
from helpers.print_style import PrintStyle
23
from helpers import files, errors
24
+from helpers.network import HttpFetchResult, fetch_public_http_resource
25
from agent import Agent
26
27
from langchain.text_splitter import RecursiveCharacterTextSplitter
28
29
30
DEFAULT_SEARCH_THRESHOLD = 0.5
31
+MAX_REMOTE_DOCUMENT_BYTES = 50 * 1024 * 1024
32
33
34
class DocumentQueryStore:
449
scheme = url.scheme or "file"
450
mimetype, encoding = mimetypes.guess_type(document_uri)
451
mimetype = mimetype or "application/octet-stream"
452
+ remote_resource: HttpFetchResult | None = None
453
454
- if mimetype == "application/octet-stream":
455
- if url.scheme in ["http", "https"]:
456
- response: aiohttp.ClientResponse | None = None
457
- retries = 0
458
- last_error = ""
459
- while not response and retries < 3:
460
- try:
461
- async with aiohttp.ClientSession() as session:
462
- response = await session.head(
463
- document_uri,
464
- timeout=aiohttp.ClientTimeout(total=2.0),
465
- allow_redirects=True,
466
- )
467
- if response.status > 399:
468
- raise Exception(response.status)
469
- break
470
- except Exception as e:
471
- await asyncio.sleep(1)
472
- last_error = str(e)
473
- retries += 1
474
- await self.agent.handle_intervention()
475
-
476
- if not response:
477
- raise ValueError(
478
- f"DocumentQueryHelper::document_get_content: Document fetch error: {document_uri} ({last_error})"
479
- )
480
-
481
- mimetype = response.headers["content-type"]
482
- if "content-length" in response.headers:
483
- content_length = (
484
- float(response.headers["content-length"]) / 1024 / 1024
485
- ) # MB
486
- if content_length > 50.0:
487
- raise ValueError(
488
- f"Document content length exceeds max. 50MB: {content_length} MB ({document_uri})"
489
- )
490
- if mimetype and "; charset=" in mimetype:
491
- mimetype = mimetype.split("; charset=")[0]
454
+ if scheme in ["http", "https"]:
455
+ remote_resource = await asyncio.to_thread(
456
+ fetch_public_http_resource,
457
+ document_uri,
458
+ max_bytes=MAX_REMOTE_DOCUMENT_BYTES,
459
+ )
460
+ if (
461
+ remote_resource.content_type
462
+ and remote_resource.content_type != "application/octet-stream"
463
+ ):
464
+ mimetype = remote_resource.content_type
465
466
if scheme == "file":
467
try:
488
if not exists:
489
await self.agent.handle_intervention()
490
if mimetype.startswith("image/"):
518
- document_content = self.handle_image_document(document_uri, scheme)
491
+ document_content = self.handle_image_document(
492
+ document_uri, scheme, remote_resource=remote_resource
493
+ )
494
elif mimetype == "text/html":
520
- document_content = self.handle_html_document(document_uri, scheme)
495
+ document_content = self.handle_html_document(
496
+ document_uri, scheme, remote_resource=remote_resource
497
+ )
498
elif mimetype.startswith("text/") or mimetype == "application/json":
522
- document_content = self.handle_text_document(document_uri, scheme)
499
+ document_content = self.handle_text_document(
500
+ document_uri, scheme, remote_resource=remote_resource
501
+ )
502
elif mimetype == "application/pdf":
524
- document_content = self.handle_pdf_document(document_uri, scheme)
503
+ document_content = self.handle_pdf_document(
504
+ document_uri, scheme, remote_resource=remote_resource
505
+ )
506
else:
507
document_content = self.handle_unstructured_document(
527
- document_uri, scheme
508
+ document_uri, scheme, remote_resource=remote_resource
509
)
510
if add_to_db:
511
self.progress_callback(f"Indexing document")
531
)
532
return document_content
533
553
- def handle_image_document(self, document: str, scheme: str) -> str:
554
- return self.handle_unstructured_document(document, scheme)
534
+ @staticmethod
535
+ def _decode_remote_text(remote_resource: HttpFetchResult) -> str:
536
+ encoding = remote_resource.encoding or "utf-8"
537
+ try:
538
+ return remote_resource.content.decode(encoding)
539
+ except (LookupError, UnicodeDecodeError):
540
+ return remote_resource.content.decode("utf-8", errors="replace")
541
556
- def handle_html_document(self, document: str, scheme: str) -> str:
542
+ @staticmethod
543
+ def _get_temp_file_suffix(
544
+ document: str, remote_resource: HttpFetchResult | None = None
545
+ ) -> str:
546
+ parsed = urlparse(document)
547
+ _stem, ext = os.path.splitext(parsed.path or document)
548
+ if ext:
549
+ return ext
550
+
551
+ if remote_resource and remote_resource.content_type:
552
+ guessed_ext = mimetypes.guess_extension(
553
+ remote_resource.content_type, strict=False
554
+ )
555
+ if guessed_ext:
556
+ return guessed_ext
557
+
558
+ return ".bin"
559
+
560
+ def handle_image_document(
561
+ self,
562
+ document: str,
563
+ scheme: str,
564
+ remote_resource: HttpFetchResult | None = None,
565
+ ) -> str:
566
+ return self.handle_unstructured_document(
567
+ document, scheme, remote_resource=remote_resource
568
+ )
569
+
570
+ def handle_html_document(
571
+ self,
572
+ document: str,
573
+ scheme: str,
574
+ remote_resource: HttpFetchResult | None = None,
575
+ ) -> str:
576
if scheme in ["http", "https"]:
558
- loader = AsyncHtmlLoader(web_path=document)
559
- parts: list[Document] = loader.load()
577
+ if remote_resource is None:
578
+ raise ValueError("Missing prefetched remote HTML content")
579
+ html_content = self._decode_remote_text(remote_resource)
580
+ parts = [Document(page_content=html_content, metadata={"source": document})]
581
elif scheme == "file":
582
# Use RFC file operations instead of TextLoader
583
file_content_bytes = files.read_file_bin(document)
594
]
595
)
596
576
- def handle_text_document(self, document: str, scheme: str) -> str:
597
+ def handle_text_document(
598
+ self,
599
+ document: str,
600
+ scheme: str,
601
+ remote_resource: HttpFetchResult | None = None,
602
+ ) -> str:
603
if scheme in ["http", "https"]:
578
- loader = AsyncHtmlLoader(web_path=document)
579
- elements: list[Document] = loader.load()
604
+ if remote_resource is None:
605
+ raise ValueError("Missing prefetched remote text content")
606
+ file_content = self._decode_remote_text(remote_resource)
607
+ elements = [
608
+ Document(page_content=file_content, metadata={"source": document})
609
+ ]
610
elif scheme == "file":
611
# Use RFC file operations instead of TextLoader
612
file_content_bytes = files.read_file_bin(document)
620
621
return "\n".join([element.page_content for element in elements])
622
593
- def handle_pdf_document(self, document: str, scheme: str) -> str:
623
+ def handle_pdf_document(
624
+ self,
625
+ document: str,
626
+ scheme: str,
627
+ remote_resource: HttpFetchResult | None = None,
628
+ ) -> str:
629
temp_file_path = ""
630
if scheme == "file":
631
# Use RFC file operations to read the PDF file as binary
637
temp_file.write(file_content_bytes)
638
temp_file_path = temp_file.name
639
elif scheme in ["http", "https"]:
605
- # download the file from the web url to a temporary file using python libraries for downloading
606
- import requests
640
import tempfile
641
642
+ if remote_resource is None:
643
+ raise ValueError("Missing prefetched remote PDF content")
644
with tempfile.NamedTemporaryFile(delete=False, suffix=".pdf") as temp_file:
610
- response = requests.get(document, timeout=10.0)
611
- if response.status_code != 200:
612
- raise ValueError(
613
- f"DocumentQueryHelper::handle_pdf_document: Failed to download PDF from {document}: {response.status_code}"
614
- )
615
- temp_file.write(response.content)
645
+ temp_file.write(remote_resource.content)
646
temp_file_path = temp_file.name
647
else:
648
raise ValueError(f"Unsupported scheme: {scheme}")
688
finally:
689
os.unlink(temp_file_path)
690
661
- def handle_unstructured_document(self, document: str, scheme: str) -> str:
691
+ def handle_unstructured_document(
692
+ self,
693
+ document: str,
694
+ scheme: str,
695
+ remote_resource: HttpFetchResult | None = None,
696
+ ) -> str:
697
elements: list[Document] = []
698
if scheme in ["http", "https"]:
664
- # loader = UnstructuredURLLoader(urls=[document], mode="single")
665
- loader = UnstructuredLoader(
666
- web_url=document,
667
- mode="single",
668
- partition_via_api=False,
669
- # chunking_strategy="by_page",
670
- strategy="hi_res",
671
- )
672
- elements = loader.load()
699
+ if remote_resource is None:
700
+ raise ValueError("Missing prefetched remote document content")
701
+ import tempfile
702
+
703
+ temp_file_path = ""
704
+ suffix = self._get_temp_file_suffix(document, remote_resource)
705
+ with tempfile.NamedTemporaryFile(delete=False, suffix=suffix) as temp_file:
706
+ temp_file.write(remote_resource.content)
707
+ temp_file_path = temp_file.name
708
+
709
+ try:
710
+ loader = UnstructuredLoader(
711
+ file_path=temp_file_path,
712
+ mode="single",
713
+ partition_via_api=False,
714
+ # chunking_strategy="by_page",
715
+ strategy="hi_res",
716
+ )
717
+ elements = loader.load()
718
+ finally:
719
+ os.unlink(temp_file_path)
720
elif scheme == "file":
721
# Use RFC file operations to read the file as binary
722
file_content_bytes = files.read_file_bin(document)