Skip to Content

Intelligence Engine Architecture

Two Runtimes

The intelligence service ships two independent Python entry points that share the same Docker image:

Entry PointCommandPurpose
main.pyuv run python main.pygRPC server — handles synchronous chat + resource API calls
worker.pyuv run python worker.pyRedis Streams consumer — runs durable background tasks

This separation means the gRPC server and background workers scale and fail independently.


Module Layout (server/intelligence/)

server/intelligence/ ├── main.py ← gRPC server entry point ├── worker.py ← Worker runtime entry point ├── core/ │ ├── config.py ← Pydantic-settings configuration │ ├── events.py ← EventEnvelope + Streams constants │ ├── lifecycle.py ← Startup / shutdown hooks │ └── logging.py ← Structured logging setup ├── features/ │ ├── chat/ │ │ ├── grpc_servicer.py ← ChatServicer (gRPC handler) │ │ ├── service.py ← ChatService orchestration │ │ ├── graph.py ← LangGraph chat pipeline │ │ ├── pipeline.py ← Query pipeline (RAG flow) │ │ ├── memory.py ← Memory extraction worker handler │ │ └── repository.py ← chat_messages DB queries │ ├── knowledge/ │ │ ├── grpc_servicer.py ← ResourceServiceServicer (gRPC handler) │ │ ├── service.py ← KnowledgeService orchestration │ │ ├── queue.py ← Publish ingestion jobs to Redis │ │ ├── repository.py ← documents DB queries │ │ ├── worker.py ← Ingestion job consumer handler │ │ └── pipeline/ │ │ ├── processor.py ← DocumentProcessor (chunk + embed + upsert) │ │ ├── chunker.py ← Token-aware text chunker │ │ ├── crawler.py ← Web crawler │ │ ├── retry.py ← tenacity retry decorators │ │ └── scrapers/ ← Browser + web + GitHub scrapers │ ├── billing/ │ │ ├── consumer.py ← Metering event consumer handler │ │ ├── ledger.py ← Credit transaction ledger │ │ ├── calculator.py ← Token cost calculator │ │ └── publisher.py ← Publish chat events to Redis │ └── health/ │ ├── grpc_servicer.py ← HealthServicer │ └── service.py ← Dependency health checks ├── shared/ │ ├── database/ ← SQLAlchemy session, models │ ├── llm/ ← LLMGateway, ModelCatalog, embeddings │ ├── qdrant/ ← QdrantVectorStore, QdrantHybridRetriever │ ├── redis/ ← EventBus, ConsumerGroupRunner, WorkerRuntime │ ├── grpc/ ← gRPC context helpers, error mapping │ └── crypto/ ← AES-256-GCM secrets (ENCRYPTION_MASTER_KEY) ├── generated/ ← protoc-generated stubs ├── prompts/ ← chat_system.md + loader.py └── scripts/ ← backfill_qdrant.py, dlq.py, smoke_qdrant.py, etc.

gRPC Server (main.py)

Registers 3 servicers against a shared grpc.aio.server:

# features/health/grpc_servicer.py health_pb2_grpc.add_HealthServicer_to_server(HealthServicer(), server) # features/chat/grpc_servicer.py intelligence_pb2_grpc.add_ChatServicer_to_server(ChatServicer(catalog, retriever), server) # features/knowledge/grpc_servicer.py intelligence_pb2_grpc.add_ResourceServiceServicer_to_server(ResourceServicer(svc), server)

Startup sequence in main.py:

  1. Load core/config.py (Pydantic-settings, validates all env vars)
  2. Initialize DB pool (shared/database/session.py)
  3. Connect to Qdrant, ensure knowledge_chunks collection exists
  4. Connect to Redis
  5. Build ModelCatalog, CloudEmbeddingGateway, QdrantHybridRetriever
  6. Register gRPC servicers and start server on 0.0.0.0:50051

Worker Runtime (worker.py)

Consumes Redis Streams events. Controlled by WORKER_ROLE env var:

WORKER_ROLEStreamGroupHandler
ingestionot:intel:ingestion_jobsworkersfeatures/knowledge/worker.py::handle_ingestion_job
meteringot:intel:chat_eventsmeteringfeatures/billing/consumer.py::handle_chat_completed
memoryot:intel:chat_eventsmemoryfeatures/chat/memory.py::handle_memory_extraction
heartbeatot:intel:heartbeatworkerslog + ack
all (default)all of the above—all handlers

Startup sweep: on boot, the worker calls sweep_stuck_jobs() to re-enqueue any ingestion jobs left in processing state by a previously crashed worker. A periodic sweep runs every 5 minutes.

Stream metrics: every 60 seconds, logs length, pending, and dlq_depth for each stream.

DLQ: after max_retries=3 failed handler attempts, the message is written to {stream}:dlq and acknowledged from the original stream.


Architecture Diagram

Last updated on