Skip to content

pipeline

pipeline

IngestionPipeline — deduplicate, chunk, and store Documents.

Takes Document objects from connectors, deduplicates by doc_id, splits content using SemanticChunker, and persists chunks to a KnowledgeStore.

Typical usage::

store = KnowledgeStore(db_path=":memory:")
pipeline = IngestionPipeline(store)
n_chunks = pipeline.ingest(connector.sync())

Classes

IngestionPipeline

IngestionPipeline(store: KnowledgeStore, *, max_tokens: int = 512, attachment_store: Optional[AttachmentStore] = None, embedder: Optional[OllamaEmbedder] = None)

Deduplicate, chunk, and index documents into a KnowledgeStore.

PARAMETER DESCRIPTION
store

The KnowledgeStore instance to write chunks into.

TYPE: KnowledgeStore

max_tokens

Soft upper-limit on chunk size passed to SemanticChunker.

TYPE: int DEFAULT: 512

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.

TYPE: Optional[AttachmentStore] DEFAULT: None

embedder

Optional embedding client (e.g. OllamaEmbedder). When provided, every chunk is embedded at ingest time and the resulting float32 vector is written to the embedding BLOB column alongside embedding_model_version. None (default) skips embedding so in-memory tests and offline runs don't depend on a sidecar daemon.

TYPE: Optional[OllamaEmbedder] DEFAULT: None

Source code in src/openjarvis/connectors/pipeline.py
def __init__(
    self,
    store: KnowledgeStore,
    *,
    max_tokens: int = 512,
    attachment_store: Optional[AttachmentStore] = None,
    embedder: Optional[OllamaEmbedder] = None,
) -> None:
    self._store = store
    self._chunker = SemanticChunker(max_tokens=max_tokens)
    self._attachment_store = attachment_store
    self._embedder = embedder
Functions
ingest
ingest(documents: Iterable[Document]) -> int

Ingest an iterable of documents into the knowledge store.

Duplicate doc_id values within a call are skipped. Documents already in the store are compared by chunk hash: unchanged documents are skipped, while edited documents are rewritten.

PARAMETER DESCRIPTION
documents

An iterable of Document objects (e.g. from a connector's sync() method).

TYPE: Iterable[Document]

RETURNS DESCRIPTION
int

The total number of chunks written to the store in this call.

Source code in src/openjarvis/connectors/pipeline.py
def ingest(self, documents: Iterable[Document]) -> int:
    """Ingest an iterable of documents into the knowledge store.

    Duplicate ``doc_id`` values within a call are skipped. Documents
    already in the store are compared by chunk hash: unchanged documents
    are skipped, while edited documents are rewritten.

    Parameters
    ----------
    documents:
        An iterable of ``Document`` objects (e.g. from a connector's
        ``sync()`` method).

    Returns
    -------
    int
        The total number of chunks written to the store in this call.
    """
    chunks_stored = 0
    # Keep dedup local to this call. A long-lived pipeline is reused by
    # SyncEngine/SyncScheduler across syncs, where a previously seen
    # document may now contain edits that need hash comparison below.
    seen_doc_ids: set[str] = set()

    for doc in documents:
        if doc.doc_id in seen_doc_ids:
            continue

        # Compute v1 provenance fields once per document.
        namespaced_thread = _namespace_thread_id(doc.source, doc.thread_id)
        source_id = _derive_source_id(doc)
        ingest_epoch = time.time()

        # Build the parent metadata dict that will be inherited by every
        # chunk produced from this document.
        parent_meta = {
            "title": doc.title,
            "author": doc.author,
            "source": doc.source,
            "source_id": source_id,
            "doc_type": doc.doc_type,
            "url": doc.url or "",
            "thread_id": namespaced_thread or "",
            "channel": doc.channel or "",
        }
        # Merge any extra connector-level metadata (without overwriting
        # the standard provenance fields set above).
        parent_meta.update(doc.metadata)

        # Normalise the timestamp to a string once.
        if hasattr(doc.timestamp, "isoformat"):
            timestamp_str = doc.timestamp.isoformat()
        else:
            timestamp_str = str(doc.timestamp)

        # Chunk the document content using the type-aware strategy.
        chunks = self._chunker.chunk(
            doc.content,
            doc_type=doc.doc_type,
            metadata=parent_meta,
        )

        # Change detection against what is already stored. Unchanged
        # documents are skipped outright; changed ones are deleted
        # wholesale (body and attachment chunks) and rewritten so a
        # document that shrank does not leave orphaned trailing chunks
        # behind. Compared before the chunks are embedded so an unchanged
        # vault costs no embedder calls on re-sync. Only body chunks are
        # compared: attachments never change without their parent.
        #
        # An unchanged document is still rewritten when an embedder is
        # configured and its stored rows carry a different
        # ``embedding_model_version`` -- including "" for rows ingested
        # before any embedder was available -- so pulling an embedding
        # model after the first sync backfills vectors instead of leaving
        # hybrid search with nothing to score. With no embedder configured
        # (daemon down, model not pulled) existing vectors are left alone.
        new_hashes = {c.index: _content_hash(c.content) for c in chunks}
        existing_hashes = self._store.chunk_hashes(doc.doc_id, source_id)
        if existing_hashes:
            stale_vectors = (
                self._embedder is not None
                and self._store.embedding_versions(doc.doc_id, source_id)
                != {self._embedder.model_version}
            )
            if existing_hashes == new_hashes and not stale_vectors:
                seen_doc_ids.add(doc.doc_id)
                continue
            self._store.delete(doc.doc_id)

        for chunk in chunks:
            embedding_bytes, embedding_version = self._embed_chunk(chunk.content)
            self._store.store(
                content=chunk.content,
                source=doc.source,
                source_id=source_id,
                doc_type=doc.doc_type,
                doc_id=doc.doc_id,
                title=doc.title,
                author=doc.author,
                participants=doc.participants,
                participants_raw=doc.participants_raw,
                timestamp=timestamp_str,
                thread_id=namespaced_thread,
                channel=doc.channel,
                url=doc.url,
                metadata=chunk.metadata,
                chunk_index=chunk.index,
                content_hash=_content_hash(chunk.content),
                embedding=embedding_bytes,
                embedding_model_version=embedding_version,
                last_synced=ingest_epoch,
            )
            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,
                        },
                    )
                    # Synthetic source_id keeps attachment chunks distinct
                    # from body chunks under the UNIQUE(source, source_id,
                    # chunk_index) constraint while still letting them share
                    # a parent doc_id for dedup and blob linkage.
                    att_source_id = f"{source_id}#{att.filename}"
                    for chunk in att_chunks:
                        embedding_bytes, embedding_version = self._embed_chunk(
                            chunk.content
                        )
                        self._store.store(
                            content=chunk.content,
                            source=doc.source,
                            source_id=att_source_id,
                            doc_type=doc.doc_type,
                            doc_id=doc.doc_id,
                            title=f"{doc.title} [{att.filename}]",
                            author=doc.author,
                            participants=doc.participants,
                            participants_raw=doc.participants_raw,
                            timestamp=timestamp_str,
                            thread_id=namespaced_thread,
                            channel=doc.channel,
                            url=doc.url,
                            metadata=chunk.metadata,
                            chunk_index=chunk.index,
                            content_hash=_content_hash(chunk.content),
                            embedding=embedding_bytes,
                            embedding_model_version=embedding_version,
                            last_synced=ingest_epoch,
                        )
                        chunks_stored += 1

        seen_doc_ids.add(doc.doc_id)

    return chunks_stored