Add parser engine: base with subprocess isolation, HTML/PDF/email/cable/text parsers
This commit is contained in:
+161
@@ -0,0 +1,161 @@
|
||||
"""Base parser with subprocess isolation for Mosaic document processing."""
|
||||
import os
|
||||
import logging
|
||||
import multiprocessing
|
||||
from multiprocessing import Process, Queue
|
||||
from typing import Optional
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Default limits
|
||||
DEFAULT_MEMORY_LIMIT_MB = 512
|
||||
DEFAULT_TIMEOUT_SECONDS = 60
|
||||
|
||||
|
||||
@dataclass
|
||||
class ParseResult:
|
||||
"""Result from a parser subprocess."""
|
||||
success: bool = False
|
||||
text: str = ""
|
||||
metadata: dict = field(default_factory=dict)
|
||||
doc_type: str = ""
|
||||
title: str = ""
|
||||
char_count: int = 0
|
||||
error: Optional[str] = None
|
||||
needs_ocr: bool = False
|
||||
|
||||
|
||||
class ParserManager:
|
||||
"""Manages document parsing with subprocess isolation.
|
||||
|
||||
All parsers run in subprocess workers with:
|
||||
- Memory limits (resource.setrlimit)
|
||||
- Timeouts
|
||||
- No access to config, SQLite, or network
|
||||
"""
|
||||
|
||||
def __init__(self, memory_limit_mb: int = DEFAULT_MEMORY_LIMIT_MB,
|
||||
timeout_seconds: int = DEFAULT_TIMEOUT_SECONDS,
|
||||
min_chars_per_page: int = 50,
|
||||
max_pdf_pages: int = 200):
|
||||
self.memory_limit_mb = memory_limit_mb
|
||||
self.timeout_seconds = timeout_seconds
|
||||
self.min_chars_per_page = min_chars_per_page
|
||||
self.max_pdf_pages = max_pdf_pages
|
||||
self._parsers = {}
|
||||
self._register_parsers()
|
||||
|
||||
def _register_parsers(self):
|
||||
"""Register available parsers by doc_type."""
|
||||
from .html_parser import HTMLDocParser
|
||||
from .pdf_parser import PDFDocParser
|
||||
from .email_parser import EmailDocParser
|
||||
from .cable_parser import CableDocParser
|
||||
from .text_parser import TextDocParser
|
||||
|
||||
self._parsers = {
|
||||
'html': HTMLDocParser(),
|
||||
'pdf': PDFDocParser(max_pages=self.max_pdf_pages,
|
||||
min_chars_per_page=self.min_chars_per_page),
|
||||
'email': EmailDocParser(),
|
||||
'cable': CableDocParser(),
|
||||
'text': TextDocParser(),
|
||||
}
|
||||
|
||||
def parse(self, file_path: str, doc_type: str,
|
||||
content_hash: Optional[str] = None) -> ParseResult:
|
||||
"""Parse a document in a subprocess worker.
|
||||
|
||||
Args:
|
||||
file_path: Path to the document file
|
||||
doc_type: Document type (html, pdf, email, cable, text)
|
||||
content_hash: Expected hash for integrity verification
|
||||
|
||||
Returns:
|
||||
ParseResult with extracted text and metadata
|
||||
"""
|
||||
if not os.path.exists(file_path):
|
||||
return ParseResult(error=f"File not found: {file_path}")
|
||||
|
||||
# Verify hash before parsing
|
||||
if content_hash:
|
||||
from utils.cache import DownloadCache
|
||||
actual_hash = DownloadCache.compute_hash(file_path)
|
||||
if actual_hash != content_hash:
|
||||
return ParseResult(
|
||||
error=f"Hash mismatch: expected {content_hash}, got {actual_hash}"
|
||||
)
|
||||
|
||||
parser = self._parsers.get(doc_type)
|
||||
if not parser:
|
||||
# Try to detect type
|
||||
parser = self._detect_parser(file_path)
|
||||
if not parser:
|
||||
return ParseResult(error=f"No parser for doc_type: {doc_type}")
|
||||
|
||||
# Run in subprocess
|
||||
result_queue = Queue()
|
||||
proc = Process(
|
||||
target=self._worker,
|
||||
args=(parser, file_path, result_queue, self.memory_limit_mb)
|
||||
)
|
||||
proc.start()
|
||||
proc.join(timeout=self.timeout_seconds)
|
||||
|
||||
if proc.is_alive():
|
||||
proc.terminate()
|
||||
proc.join(timeout=5)
|
||||
if proc.is_alive():
|
||||
proc.kill()
|
||||
return ParseResult(error=f"Parser timed out after {self.timeout_seconds}s")
|
||||
|
||||
if result_queue.empty():
|
||||
return ParseResult(error="Parser process returned no result")
|
||||
|
||||
result = result_queue.get_nowait()
|
||||
|
||||
# Check for OCR needed (low char count in PDFs)
|
||||
if doc_type == 'pdf' and result.success:
|
||||
if result.char_count > 0 and result.metadata.get('page_count', 1) > 0:
|
||||
chars_per_page = result.char_count / result.metadata.get('page_count', 1)
|
||||
if chars_per_page < self.min_chars_per_page:
|
||||
result.needs_ocr = True
|
||||
logger.info("Document may need OCR: %.1f chars/page", chars_per_page)
|
||||
|
||||
return result
|
||||
|
||||
@staticmethod
|
||||
def _worker(parser, file_path: str, result_queue: Queue, memory_limit_mb: int):
|
||||
"""Subprocess worker that runs the parser with resource limits."""
|
||||
try:
|
||||
# Set memory limit
|
||||
try:
|
||||
import resource
|
||||
mem_bytes = memory_limit_mb * 1024 * 1024
|
||||
resource.setrlimit(resource.RLIMIT_AS, (mem_bytes, mem_bytes))
|
||||
except (ImportError, ValueError, OSError) as e:
|
||||
# resource module may not be available on all platforms
|
||||
pass
|
||||
|
||||
result = parser.parse(file_path)
|
||||
result_queue.put(result)
|
||||
except MemoryError:
|
||||
result_queue.put(ParseResult(error="Parser exceeded memory limit"))
|
||||
except Exception as e:
|
||||
result_queue.put(ParseResult(error=f"Parser error: {str(e)}"))
|
||||
|
||||
def _detect_parser(self, file_path: str) -> Optional[object]:
|
||||
"""Try to detect the right parser from file content."""
|
||||
try:
|
||||
from utils.sanitize import validate_content_type
|
||||
content_type = validate_content_type(file_path)
|
||||
if 'html' in content_type:
|
||||
return self._parsers['html']
|
||||
elif 'pdf' in content_type:
|
||||
return self._parsers['pdf']
|
||||
elif 'text' in content_type:
|
||||
return self._parsers['text']
|
||||
except Exception:
|
||||
pass
|
||||
return self._parsers.get('text')
|
||||
Reference in New Issue
Block a user