Skip to Content

Knowledge Ingestion Pipeline

Overview

Ingestion in v1.1.0 is fully asynchronous. The gRPC AddResource call publishes a job to Redis Streams and returns immediately. A dedicated background worker consumes the job and runs the full pipeline.

Browser → Rust API → gRPC KnowledgeServicer → Redis INGESTION_JOBS stream ↓ Worker: handle_ingestion_job() ↓ DocumentProcessor → Qdrant (vectors) + PostgreSQL (metadata)

Module Layout

features/knowledge/ ├── grpc_servicer.py ← AddResource / ListResources / DeleteResource handlers ├── service.py ← KnowledgeService — orchestrates DB + Redis ├── queue.py ← publish_ingestion_job() → Redis INGESTION_JOBS stream ├── repository.py ← documents / ingestion_jobs DB queries ├── worker.py ← handle_ingestion_job() — worker consumer handler └── pipeline/ ├── processor.py ← DocumentProcessor: clean → chunk → embed → upsert ├── chunker.py ← Token-aware text splitter (512 tokens, 50 overlap) ├── crawler.py ← Web crawler with depth limits and rate control ├── retry.py ← tenacity retry decorators ├── validation.py ← Content length / type validation └── scrapers/ ├── browser.py ← Playwright headless Chromium (JS-heavy SPAs) ├── web.py ← httpx lightweight static HTML fetcher └── github.py ← PyGitHub repository content fetcher

Ingestion Flow

1. gRPC AddResource → Redis

features/knowledge/grpc_servicer.py validates the request and delegates to KnowledgeService.add_resource():

# features/knowledge/service.py async def add_resource(self, request: AddResourceRequest) -> AddResourceResponse: # 1. Create document row in PostgreSQL (status=queued) doc_id = await self.repository.create_document(...) job_id = await self.repository.create_ingestion_job(doc_id) # 2. Publish job to Redis INGESTION_JOBS stream await self.queue.publish_ingestion_job(doc_id, job_id, request) # 3. Return immediately — worker handles the rest return AddResourceResponse(resource_id=doc_id, job_id=job_id, status="queued")

2. Worker Receives Job

features/knowledge/worker.py::handle_ingestion_job() is called by the Redis consumer:

async def handle_ingestion_job(envelope: EventEnvelope, msg_id: str) -> None: # Update job status → processing # Run DocumentProcessor # Update job status → completed / failed # Sweep: auto-recover stuck jobs on startup

3. DocumentProcessor (pipeline/processor.py)

Full pipeline for a single document:

  1. Scrape — BrowserScraper (Playwright), WebScraper (httpx), or GitHubScraper
  2. Clean — remove boilerplate, normalize whitespace, strip HTML
  3. Chunk — TextChunker: 512-token chunks with 50-token overlap
  4. Embed — CloudEmbeddingGateway.embed_documents() — text-embedding-3-large (3072d), batched 64 texts/call
  5. Sparse encode — HashedSparseEncoder.encode() — BM25-style hashed unigram TF
  6. Upsert to Qdrant — QdrantVectorStore.upsert_chunks() — stores dense + sparse vectors + payload
  7. Update PostgreSQL — mark document status=completed, record chunk count

4. Qdrant Point Structure

Each chunk becomes one Qdrant point in the knowledge_chunks collection:

PointStruct( id=chunk_uuid, vector={ "dense": [float, ...], # 3072-dim cosine "sparse": SparseVector( indices=[int, ...], # hashed token indices (dim=65536) values=[float, ...], # sublinear TF weights ), }, payload={ "user_id": str, # tenant partitioning "is_global": bool, # visible to all users if True "document_id": str, # UUID → PostgreSQL documents.id "chunk_index": int, # position in document "content": str, # raw text of the chunk "document_title": str, "source_url": str, }, )

Payload indexes: user_id (keyword / tenant), is_global (bool), document_id (keyword), chunk_index (integer).


Startup Sweep

On boot, sweep_stuck_jobs() finds any ingestion jobs in processing state (crashed workers) and re-publishes them to the Redis stream. A periodic sweep runs every 5 minutes.


Configuration

VariableDefaultDescription
INGESTION_CHUNK_SIZE512Tokens per chunk
INGESTION_CHUNK_OVERLAP50Overlap tokens between chunks
INGESTION_MAX_BATCH_SIZE100Max chunks per Qdrant upsert call
INGESTION_AUTO_CLEANtrueRun content cleaner before chunking
INGESTION_GENERATE_EMBEDDINGStrueGenerate + upsert vectors
INGESTION_MAX_CONTENT_LENGTH1000000Max bytes before truncation
SCRAPING_TIMEOUT30HTTP timeout seconds
SCRAPING_MAX_RETRIES3Retry attempts for failed fetches

Supported Resource Types

ResourceTypeScraper UsedNotes
TEXTNone (direct content)Content passed inline in the gRPC request
WEBSITEBrowserScraper → WebScraper fallbackHandles SPAs and static HTML
GITHUB_REPOGitHubScraper (PyGitHub)Fetches README + source files
FILEChunkedUpload gRPC streamUp to 1 GB; streamed in 10 MB chunks
Last updated on