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 fetcherIngestion 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 startup3. DocumentProcessor (pipeline/processor.py)
Full pipeline for a single document:
- Scrape —
BrowserScraper(Playwright),WebScraper(httpx), orGitHubScraper - Clean — remove boilerplate, normalize whitespace, strip HTML
- Chunk —
TextChunker: 512-token chunks with 50-token overlap - Embed —
CloudEmbeddingGateway.embed_documents()—text-embedding-3-large(3072d), batched 64 texts/call - Sparse encode —
HashedSparseEncoder.encode()— BM25-style hashed unigram TF - Upsert to Qdrant —
QdrantVectorStore.upsert_chunks()— stores dense + sparse vectors + payload - 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
| Variable | Default | Description |
|---|---|---|
INGESTION_CHUNK_SIZE | 512 | Tokens per chunk |
INGESTION_CHUNK_OVERLAP | 50 | Overlap tokens between chunks |
INGESTION_MAX_BATCH_SIZE | 100 | Max chunks per Qdrant upsert call |
INGESTION_AUTO_CLEAN | true | Run content cleaner before chunking |
INGESTION_GENERATE_EMBEDDINGS | true | Generate + upsert vectors |
INGESTION_MAX_CONTENT_LENGTH | 1000000 | Max bytes before truncation |
SCRAPING_TIMEOUT | 30 | HTTP timeout seconds |
SCRAPING_MAX_RETRIES | 3 | Retry attempts for failed fetches |
Supported Resource Types
ResourceType | Scraper Used | Notes |
|---|---|---|
TEXT | None (direct content) | Content passed inline in the gRPC request |
WEBSITE | BrowserScraper → WebScraper fallback | Handles SPAs and static HTML |
GITHUB_REPO | GitHubScraper (PyGitHub) | Fetches README + source files |
FILE | ChunkedUpload gRPC stream | Up to 1 GB; streamed in 10 MB chunks |