From ee1805ceaf7c9feb487341a29dcfb730d2ccc722 Mon Sep 17 00:00:00 2001 From: krypticmouse Date: Thu, 26 Mar 2026 21:18:47 +0000 Subject: [PATCH] feat: wire attachment text extraction into IngestionPipeline Add optional `attachment_store` parameter to IngestionPipeline. When provided, attachments are stored as blobs in AttachmentStore and their text (plain/markdown/csv via decode, PDF via pdfplumber) is chunked and indexed in the KnowledgeStore with attachment provenance metadata. Co-Authored-By: Claude Sonnet 4.6 --- src/openjarvis/connectors/pipeline.py | 83 ++++++++- tests/connectors/test_pipeline_attachments.py | 163 ++++++++++++++++++ 2 files changed, 243 insertions(+), 3 deletions(-) create mode 100644 tests/connectors/test_pipeline_attachments.py diff --git a/src/openjarvis/connectors/pipeline.py b/src/openjarvis/connectors/pipeline.py index 2585d0e0..4df3a030 100644 --- a/src/openjarvis/connectors/pipeline.py +++ b/src/openjarvis/connectors/pipeline.py @@ -13,12 +13,15 @@ Typical usage:: from __future__ import annotations -from typing import Iterable +from typing import TYPE_CHECKING, Iterable, Optional -from openjarvis.connectors._stubs import Document +from openjarvis.connectors._stubs import Attachment, Document from openjarvis.connectors.chunker import SemanticChunker from openjarvis.connectors.store import KnowledgeStore +if TYPE_CHECKING: + from openjarvis.connectors.attachment_store import AttachmentStore + class IngestionPipeline: """Deduplicate, chunk, and index documents into a KnowledgeStore. @@ -29,11 +32,22 @@ class IngestionPipeline: The ``KnowledgeStore`` instance to write chunks into. max_tokens: Soft upper-limit on chunk size passed to ``SemanticChunker``. + attachment_store: + Optional ``AttachmentStore`` for persisting attachment blobs and + extracting text from supported MIME types (PDF, plain text, etc.). + When ``None`` (default) attachments are silently ignored. """ - def __init__(self, store: KnowledgeStore, *, max_tokens: int = 512) -> None: + def __init__( + self, + store: KnowledgeStore, + *, + max_tokens: int = 512, + attachment_store: Optional[AttachmentStore] = None, + ) -> None: self._store = store self._chunker = SemanticChunker(max_tokens=max_tokens) + self._attachment_store = attachment_store self._seen_doc_ids: set[str] = set() self._load_existing_doc_ids() @@ -48,6 +62,26 @@ class IngestionPipeline: ).fetchall() self._seen_doc_ids = {r[0] for r in rows} + def _extract_attachment_text(self, att: Attachment) -> str: + """Extract text from an attachment. + + Returns the extracted text, or an empty string if the MIME type is + unsupported or extraction fails. + """ + if att.mime_type == "application/pdf": + try: + import io + + import pdfplumber + + with pdfplumber.open(io.BytesIO(att.content)) as pdf: + return "\n".join(page.extract_text() or "" for page in pdf.pages) + except Exception: # noqa: BLE001 + return "" + if att.mime_type in ("text/plain", "text/markdown", "text/csv"): + return att.content.decode("utf-8", errors="replace") + return "" + # ------------------------------------------------------------------ # Public API # ------------------------------------------------------------------ @@ -119,6 +153,49 @@ class IngestionPipeline: ) chunks_stored += 1 + # Process attachments when an attachment store is configured. + if self._attachment_store and doc.attachments: + for att in doc.attachments: + if not att.content: + continue + + # Persist the raw blob and obtain its SHA-256. + sha = self._attachment_store.store( + content=att.content, + filename=att.filename, + mime_type=att.mime_type, + source_doc_id=doc.doc_id, + ) + + # Extract searchable text and index it as additional chunks. + extracted = self._extract_attachment_text(att) + if extracted: + att_chunks = self._chunker.chunk( + extracted, + doc_type=doc.doc_type, + metadata={ + **parent_meta, + "attachment": att.filename, + "sha256": sha, + }, + ) + for chunk in att_chunks: + self._store.store( + content=chunk.content, + source=doc.source, + doc_type=doc.doc_type, + doc_id=doc.doc_id, + title=f"{doc.title} [{att.filename}]", + author=doc.author, + participants=doc.participants, + timestamp=timestamp_str, + thread_id=doc.thread_id, + url=doc.url, + metadata=chunk.metadata, + chunk_index=chunk.index, + ) + chunks_stored += 1 + self._seen_doc_ids.add(doc.doc_id) return chunks_stored diff --git a/tests/connectors/test_pipeline_attachments.py b/tests/connectors/test_pipeline_attachments.py new file mode 100644 index 00000000..912cabd0 --- /dev/null +++ b/tests/connectors/test_pipeline_attachments.py @@ -0,0 +1,163 @@ +"""Tests for attachment processing in IngestionPipeline.""" + +from __future__ import annotations + +from datetime import datetime, timezone +from pathlib import Path + +import pytest + +from openjarvis.connectors._stubs import Attachment, Document +from openjarvis.connectors.attachment_store import AttachmentStore +from openjarvis.connectors.pipeline import IngestionPipeline +from openjarvis.connectors.store import KnowledgeStore + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _make_doc(**kwargs) -> Document: # type: ignore[type-arg] + """Build a Document with sensible defaults.""" + defaults = dict( + doc_id="doc:att:001", + source="gmail", + doc_type="email", + content="Main body of the email.", + title="Test Email", + author="sender@example.com", + participants=[], + timestamp=datetime(2025, 3, 1, tzinfo=timezone.utc), + thread_id=None, + url=None, + attachments=[], + metadata={}, + ) + defaults.update(kwargs) + return Document(**defaults) # type: ignore[arg-type] + + +# --------------------------------------------------------------------------- +# Fixtures +# --------------------------------------------------------------------------- + + +@pytest.fixture() +def store(tmp_path: Path) -> KnowledgeStore: + return KnowledgeStore(db_path=tmp_path / "att_pipeline.db") + + +@pytest.fixture() +def att_store(tmp_path: Path) -> AttachmentStore: + return AttachmentStore(base_dir=str(tmp_path / "blobs")) + + +@pytest.fixture() +def pipeline(store: KnowledgeStore, att_store: AttachmentStore) -> IngestionPipeline: + return IngestionPipeline(store, attachment_store=att_store) + + +# --------------------------------------------------------------------------- +# Test 1: text/plain attachment text is indexed and searchable +# --------------------------------------------------------------------------- + + +def test_plain_text_attachment_indexed( + pipeline: IngestionPipeline, store: KnowledgeStore +) -> None: + """Ingesting a doc with a text/plain attachment indexes the attachment text.""" + att_text = b"Deep learning quarterly report: model accuracy improved by 15%." + att = Attachment( + filename="report.txt", + mime_type="text/plain", + size_bytes=len(att_text), + content=att_text, + ) + doc = _make_doc( + doc_id="doc:att:plain:001", + content="Please see the attached report.", + title="Q1 Report Email", + attachments=[att], + ) + + n = pipeline.ingest([doc]) + + # Expect at least 2 chunks: 1 from the main body + 1 from the attachment + assert n >= 2 + + # Attachment text must be retrievable + results = store.retrieve("model accuracy improved", top_k=5) + assert len(results) >= 1 + + # At least one result should come from the attachment chunk + att_results = [r for r in results if r.metadata.get("attachment") == "report.txt"] + assert len(att_results) >= 1 + + # Title for attachment chunks includes the filename in brackets + assert "report.txt" in att_results[0].metadata.get("title", "") + + +# --------------------------------------------------------------------------- +# Test 2: blob is stored in AttachmentStore +# --------------------------------------------------------------------------- + + +def test_attachment_blob_stored( + pipeline: IngestionPipeline, att_store: AttachmentStore +) -> None: + """The raw bytes of an attachment are persisted in the AttachmentStore.""" + att_content = b"Confidential: merger details enclosed." + att = Attachment( + filename="merger.txt", + mime_type="text/plain", + size_bytes=len(att_content), + content=att_content, + ) + doc = _make_doc( + doc_id="doc:att:blob:001", + content="See attached for details.", + attachments=[att], + ) + + pipeline.ingest([doc]) + + # The blob must be retrievable from the attachment store + import hashlib + + expected_sha = hashlib.sha256(att_content).hexdigest() + stored_bytes = att_store.get_content(expected_sha) + assert stored_bytes == att_content + + # Metadata must link back to the source document + meta = att_store.get_metadata(expected_sha) + assert meta is not None + assert meta["filename"] == "merger.txt" + assert meta["mime_type"] == "text/plain" + assert "doc:att:blob:001" in meta["source_doc_ids"] + + +# --------------------------------------------------------------------------- +# Test 3: no regression when doc has no attachments +# --------------------------------------------------------------------------- + + +def test_no_attachments_no_regression( + store: KnowledgeStore, att_store: AttachmentStore +) -> None: + """Pipeline with attachment_store handles docs without attachments correctly.""" + pipeline = IngestionPipeline(store, attachment_store=att_store) + + doc = _make_doc( + doc_id="doc:no:att:001", + content="Just a regular email with no attachments.", + attachments=[], + ) + + n = pipeline.ingest([doc]) + + assert n == 1 + assert store.count() == 1 + + # No blobs should be stored + rows = att_store._conn.execute("SELECT COUNT(*) FROM attachments").fetchone() + assert rows[0] == 0