Intelligence Engine Architecture
Two Runtimes
The intelligence service ships two independent Python entry points that share the same Docker image:
| Entry Point | Command | Purpose |
|---|---|---|
main.py | uv run python main.py | gRPC server — handles synchronous chat + resource API calls |
worker.py | uv run python worker.py | Redis 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:
- Load
core/config.py(Pydantic-settings, validates all env vars) - Initialize DB pool (
shared/database/session.py) - Connect to Qdrant, ensure
knowledge_chunkscollection exists - Connect to Redis
- Build
ModelCatalog,CloudEmbeddingGateway,QdrantHybridRetriever - 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_ROLE | Stream | Group | Handler |
|---|---|---|---|
ingestion | ot:intel:ingestion_jobs | workers | features/knowledge/worker.py::handle_ingestion_job |
metering | ot:intel:chat_events | metering | features/billing/consumer.py::handle_chat_completed |
memory | ot:intel:chat_events | memory | features/chat/memory.py::handle_memory_extraction |
heartbeat | ot:intel:heartbeat | workers | log + 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