1"""
2Simplified preprocessing pipeline after modular refactoring.
3
4Keeps only the essential pipeline functionality.
5"""
6
7from __future__ import annotations
8
9from lexigram.ai.rag.preprocessing.base import AbstractPreprocessor
10from lexigram.ai.rag.preprocessing.document import PreprocessedDocument
11from lexigram.ai.rag.preprocessing.types import DocumentType
12
13
14class PreprocessingPipeline:
15 """Simplified document preprocessing pipeline.
16
17 Provides a clean interface for document preprocessing
18 with support for OCR, table extraction, and metadata enrichment.
19 """
20
21 def __init__(self, preprocessors: list[AbstractPreprocessor] | None = None):
22 """Initialize preprocessing pipeline.
23
24 Args:
25 preprocessors: List of preprocessors to run in sequence.
26 """
27 self.preprocessors = preprocessors or []
28
29 def add_preprocessor(self, preprocessor: AbstractPreprocessor) -> None:
30 """Add a preprocessor to the pipeline.
31
32 Args:
33 preprocessor: Preprocessor to add.
34 """
35 self.preprocessors.append(preprocessor)
36
37 async def process(
38 self,
39 content: str,
40 doc_type: DocumentType = DocumentType.UNKNOWN,
41 **kwargs,
42 ) -> PreprocessedDocument:
43 """Process document through pipeline.
44
45 Args:
46 content: Document content to process.
47 doc_type: Type of document (auto-detected if unknown).
48 **kwargs: Additional processing parameters.
49
50 Returns:
51 Preprocessed document with extracted information.
52 """
53 current_content = content
54 current_metadata = None
55 images = []
56 tables = []
57
58 # Run each preprocessor in sequence
59 for preprocessor in self.preprocessors:
60 result = await preprocessor.preprocess(current_content, **kwargs)
61 current_content = result.content
62 if not current_metadata:
63 current_metadata = result.metadata
64 # Merge metadata
65 elif hasattr(current_metadata, "merge"):
66 current_metadata.merge(result.metadata)
67
68 if hasattr(result, "images") and result.images:
69 images.extend(result.images)
70 if hasattr(result, "tables") and result.tables:
71 tables.extend(result.tables)
72
73 if not self.preprocessors:
74 return PreprocessedDocument(content=content, raw_content=content)
75
76 return PreprocessedDocument(
77 content=current_content,
78 metadata=current_metadata,
79 raw_content=content,
80 images=images,
81 tables=tables,
82 )
83
84 async def process_batch(
85 self,
86 documents: list[tuple[str, DocumentType]],
87 **kwargs,
88 ) -> list[PreprocessedDocument]:
89 """Process multiple documents in parallel.
90
91 Args:
92 documents: List of (content, doc_type) tuples.
93 **kwargs: Additional processing parameters.
94
95 Returns:
96 List of preprocessed documents.
97 """
98 import asyncio
99
100 tasks = [
101 self.process(content, doc_type, **kwargs) for content, doc_type in documents
102 ]
103 results = await asyncio.gather(*tasks, return_exceptions=True)
104 processed = []
105 for result in results:
106 if isinstance(result, BaseException):
107 raise result
108 processed.append(result)
109 return processed