| """ |
| Document Loader |
| Scans PDF folder, processes PDFs, and stores in database + Qdrant |
| """ |
|
|
| import sys |
| from pathlib import Path |
|
|
| |
| project_root = Path(__file__).parent.parent.parent |
| sys.path.insert(0, str(project_root)) |
|
|
| from typing import Dict |
| from datetime import datetime |
| from src.cdms.pdf_processor import PDFProcessor |
| from src.cdms.schema import DatabaseManager, Document, DocumentChunk |
| from src.rag.embeddings import OpenAIEmbeddingService |
| from src.rag.vector_store import QdrantVectorStore, get_shared_vector_store |
| from src.config.credentials import CredentialsManager |
| from src.config.paths import PDF_DIR |
| from src.cdms.product_catalog import normalize_filename |
|
|
|
|
| class DocumentLoader: |
| """ |
| Loads and indexes PDF documents for RAG |
| |
| Usage: |
| loader = DocumentLoader() |
| loader.load_all_pdfs() # Process all PDFs in data/pdfs/ |
| """ |
| |
| def __init__(self, pdf_folder: str = None, vector_store=None): |
| """ |
| Initialize document loader |
| |
| Args: |
| pdf_folder: Folder containing PDF files (defaults to the CDMS PDF dir |
| PDF_DIR = <root>/data/pdfs/cdms, where the labels live) |
| vector_store: Optional shared QdrantVectorStore. The embedded (on-disk) |
| Qdrant allows only ONE client per path per process, so in |
| online mode the loader MUST reuse the RAG searcher's store |
| instead of opening a second one (which would lock-fail and |
| silently drop the freshly fetched label's chunks). |
| """ |
| self.pdf_folder = Path(pdf_folder) if pdf_folder is not None else PDF_DIR |
| self.pdf_folder.mkdir(parents=True, exist_ok=True) |
| |
| |
| self.pdf_processor = PDFProcessor() |
| self.db_manager = DatabaseManager() |
| |
| |
| try: |
| creds = CredentialsManager() |
| openai_key = creds.get_api_key("openai") |
| self.embedding_service = OpenAIEmbeddingService(api_key=openai_key) |
| except Exception as e: |
| print(f"⚠️ Warning: Could not initialize OpenAI embeddings: {e}") |
| self.embedding_service = None |
| |
| |
| |
| if vector_store is not None: |
| self.vector_store = vector_store |
| else: |
| try: |
| self.vector_store = get_shared_vector_store() |
| except Exception as e: |
| print(f"⚠️ Warning: Could not initialize Qdrant: {e}") |
| self.vector_store = None |
| |
| def load_pdf(self, pdf_path: str, force_reprocess: bool = False, pdf_url: str = None) -> Dict: |
| """ |
| Load a single PDF file |
| |
| Args: |
| pdf_path: Path to PDF file |
| force_reprocess: If True, reprocess even if already indexed |
| pdf_url: Optional original PDF URL (for CDMS labels from Tavily) |
| |
| Returns: |
| Dict with processing result |
| """ |
| pdf_path = Path(pdf_path) |
| |
| if not pdf_path.exists(): |
| return { |
| "success": False, |
| "error": f"File not found: {pdf_path}" |
| } |
| |
| |
| session = self.db_manager.get_session() |
| try: |
| doc_id = Document.generate_id(str(pdf_path)) |
| existing_doc = session.query(Document).filter_by(id=doc_id).first() |
| |
| if existing_doc and existing_doc.processed == 1 and not force_reprocess: |
| return { |
| "success": True, |
| "message": f"PDF already processed: {pdf_path.name}", |
| "document_id": doc_id, |
| "skipped": True |
| } |
| |
| |
| if force_reprocess and existing_doc: |
| |
| session.query(DocumentChunk).filter_by(document_id=doc_id).delete() |
| |
| session.delete(existing_doc) |
| session.commit() |
| print(f"🔄 Re-processing: {pdf_path.name}") |
| |
| |
| print(f"📄 Processing: {pdf_path.name}") |
| result = self.pdf_processor.process_pdf(str(pdf_path)) |
| |
| if not result["success"]: |
| return result |
| |
| |
| |
| doc = Document( |
| id=doc_id, |
| filename=pdf_path.name, |
| filepath=str(pdf_path), |
| file_size=pdf_path.stat().st_size, |
| num_pages=result["num_pages"], |
| num_chunks=result["num_chunks"], |
| processed=1, |
| last_processed=datetime.utcnow(), |
| doc_metadata={"pdf_url": pdf_url} if pdf_url else None |
| ) |
|
|
| session.merge(doc) |
| |
| |
| chunks = result.get("chunks", []) |
| page_numbers = result.get("page_numbers", []) |
|
|
| |
| if not page_numbers: |
| print(f"⚠️ Warning: No page numbers provided for {pdf_path.name}, estimating...") |
| page_numbers = [(idx // 3) + 1 for idx in range(len(chunks))] |
| elif len(page_numbers) != len(chunks): |
| print(f"⚠️ Warning: Page numbers count ({len(page_numbers)}) doesn't match chunks count ({len(chunks)}) for {pdf_path.name}") |
| |
| if len(page_numbers) < len(chunks): |
| last_page = page_numbers[-1] if page_numbers else 1 |
| page_numbers.extend([last_page] * (len(chunks) - len(page_numbers))) |
| else: |
| page_numbers = page_numbers[:len(chunks)] |
|
|
| |
| url_hash = "" |
| if pdf_url: |
| import hashlib |
| url_hash = hashlib.md5(pdf_url.encode()).hexdigest()[:12] |
|
|
| |
| |
| |
| product = normalize_filename(pdf_path.name) |
|
|
| |
| |
| |
| chunk_ids = [] |
| chunk_texts = [] |
| chunk_pages = [] |
| chunk_indices = [] |
| for idx, (chunk_text, page_num) in enumerate(zip(chunks, page_numbers)): |
| |
| if page_num <= 0: |
| page_num = (idx // 3) + 1 |
| page_numbers[idx] = page_num |
| chunk_id = DocumentChunk.generate_id(doc_id, idx) |
|
|
| |
| token_count = len(chunk_text) // 4 |
|
|
| chunk = DocumentChunk( |
| id=chunk_id, |
| document_id=doc_id, |
| chunk_index=idx, |
| content=chunk_text, |
| page_number=page_num, |
| char_count=len(chunk_text), |
| token_count=token_count |
| ) |
| session.merge(chunk) |
|
|
| chunk_ids.append(chunk_id) |
| chunk_texts.append(chunk_text) |
| chunk_pages.append(page_num) |
| chunk_indices.append(idx) |
|
|
| chunks_stored = len(chunk_ids) |
| embeddings_generated = 0 |
|
|
| |
| if self.embedding_service and self.vector_store and chunk_texts: |
| try: |
| embeddings = self.embedding_service.generate_embeddings_batch(chunk_texts) |
| for chunk_id, embedding, chunk_text, page_num, idx in zip( |
| chunk_ids, embeddings, chunk_texts, chunk_pages, chunk_indices |
| ): |
| metadata = { |
| "document_id": doc_id, |
| "document_name": pdf_path.name, |
| "chunk_index": idx, |
| "content": chunk_text, |
| "page_number": page_num, |
| "source_file": pdf_path.name, |
| "product": product, |
| "pdf_url": pdf_url if pdf_url else "", |
| "url_hash": url_hash |
| } |
| self.vector_store.add_document_chunk( |
| chunk_id=chunk_id, |
| embedding=embedding, |
| metadata=metadata |
| ) |
| embeddings_generated += 1 |
| except Exception as e: |
| print(f"⚠️ Warning: Could not generate embeddings for {pdf_path.name}: {e}") |
|
|
| session.commit() |
| |
| return { |
| "success": True, |
| "document_id": doc_id, |
| "filename": pdf_path.name, |
| "chunks_stored": chunks_stored, |
| "embeddings_generated": embeddings_generated, |
| "num_pages": result["num_pages"] |
| } |
| |
| finally: |
| session.close() |
| |
| def load_all_pdfs(self, force_reprocess: bool = False) -> Dict: |
| """ |
| Load all PDFs from the pdf_folder |
| |
| Returns: |
| Dict with summary of processing |
| """ |
| if not self.pdf_folder.exists(): |
| return { |
| "success": False, |
| "error": f"PDF folder not found: {self.pdf_folder}" |
| } |
| |
| pdf_files = list(self.pdf_folder.glob("*.pdf")) |
| |
| if not pdf_files: |
| return { |
| "success": False, |
| "error": f"No PDF files found in {self.pdf_folder}" |
| } |
| |
| print(f"📚 Found {len(pdf_files)} PDF file(s)") |
| print("-" * 70) |
| |
| results = [] |
| for pdf_file in pdf_files: |
| result = self.load_pdf(str(pdf_file), force_reprocess=force_reprocess) |
| results.append(result) |
| if result.get("success") and not result.get("skipped"): |
| print(f"✅ {pdf_file.name}: {result.get('chunks_stored', 0)} chunks, {result.get('embeddings_generated', 0)} embeddings") |
| elif result.get("skipped"): |
| print(f"⏭️ {pdf_file.name}: Already processed (skipped)") |
| print() |
| |
| |
| successful = sum(1 for r in results if r.get("success")) |
| total_chunks = sum(r.get("chunks_stored", 0) for r in results) |
| total_embeddings = sum(r.get("embeddings_generated", 0) for r in results) |
| |
| return { |
| "success": True, |
| "total_files": len(pdf_files), |
| "successful": successful, |
| "total_chunks": total_chunks, |
| "total_embeddings": total_embeddings, |
| "results": results |
| } |
|
|
|
|
| |
| if __name__ == "__main__": |
| print("Testing Document Loader...") |
| print("=" * 70) |
| |
| loader = DocumentLoader() |
| |
| |
| result = loader.load_all_pdfs() |
| |
| if result.get("success"): |
| print("\n" + "=" * 70) |
| print("✅ Processing Complete!") |
| print(f" Files processed: {result['successful']}/{result['total_files']}") |
| print(f" Total chunks: {result['total_chunks']}") |
| print(f" Embeddings generated: {result['total_embeddings']}") |
| else: |
| print(f"\n❌ Error: {result.get('error')}") |
| print("\n💡 To use:") |
| print(" 1. Create data/pdfs/ folder") |
| print(" 2. Add PDF files") |
| print(" 3. Make sure OPENAI_API_KEY is in .env") |
| print(" 4. Make sure Qdrant is running (Docker)") |
|
|
|
|