gRPC Server Interface
Service Registration
The Intelligence service registers 3 gRPC servicers against the shared grpc.aio.server in main.py:
from features.health.grpc_servicer import HealthServicer
from features.chat.grpc_servicer import ChatServicer
from features.knowledge.grpc_servicer import ResourceServiceServicer
health_pb2_grpc.add_HealthServicer_to_server(HealthServicer(), server)
intelligence_pb2_grpc.add_ChatServicer_to_server(ChatServicer(catalog, retriever), server)
intelligence_pb2_grpc.add_ResourceServiceServicer_to_server(ResourceServiceServicer(svc), server)Health Service (features/health/)
| RPC | Input | Output | Notes |
|---|---|---|---|
Check | HealthCheckRequest {} | HealthCheckResponse { status, version, uptime_seconds } | Always returns “healthy” if process is alive |
Ready | ReadyCheckRequest {} | ReadyCheckResponse { ready, dependencies, dependency_status } | Checks PostgreSQL + Qdrant + Redis connectivity |
Chat Service (features/chat/grpc_servicer.py)
| RPC | Input | Output | Notes |
|---|---|---|---|
SendMessage | ChatRequest | ChatResponse | Unary; timeout 1200 s |
StreamChat | ChatRequest | stream ChatStreamChunk | Server-streaming; timeout 300 s |
GetConversation | GetConversationRequest { user_id, conversation_id } | ConversationResponse { messages[] } | Fetches from chat_messages |
DeleteConversation | DeleteConversationRequest { user_id, conversation_id } | DeleteConversationResponse { success } | Checks ownership |
GenerateTitle | GenerateTitleRequest { user_id, conversation_id, user_message, assistant_message } | GenerateTitleResponse { title } | LLM call |
ChatRequest Structure
message ChatRequest {
string user_id = 1;
string conversation_id = 2;
string message = 3;
optional ChatConfig config = 4;
map<string, string> metadata = 5;
}ChatStreamChunk Structure
message ChatStreamChunk {
string conversation_id = 1;
string message_id = 2;
oneof chunk_type {
TokenChunk token = 3; // incremental LLM output
SourcesChunk sources = 4; // RAG source citations
MetricsChunk metrics = 5; // token counts, latency
ErrorChunk error = 6; // terminal error
}
bool is_final = 7;
}The Rust layer reads each ChatStreamChunk from the gRPC stream and converts it to an SSE data: event. is_final=true triggers the SSE stream close.
ResourceService Service (features/knowledge/grpc_servicer.py)
[!IMPORTANT]
AddResourceis now async in v1.1.0. It publishes the ingestion job to Redis Streams and returns immediately withstatus=queued. The actual scraping, chunking, embedding and Qdrant upsert happen in the background worker.
| RPC | Type | Input | Output | Timeout |
|---|---|---|---|---|
AddResource | Unary | AddResourceRequest | AddResourceResponse { resource_id, job_id, status="queued" } | 30 s |
ChunkedUpload | Client-streaming | stream FileChunk | ChunkedUploadResponse | 3000 s |
GetResourceStatus | Unary | GetResourceStatusRequest { resource_id, user_id } | ResourceStatusResponse | 10 s |
ListResources | Unary | ListResourcesRequest { user_id, filters, pagination } | ListResourcesResponse { resources[], total, next_cursor } | 30 s |
DeleteResource | Unary | DeleteResourceRequest { resource_id, user_id } | DeleteResourceResponse { success } | 30 s |
CancelIngestion | Unary | CancelIngestionRequest { job_id, user_id } | CancelIngestionResponse { success } | 10 s |
SyncResourceMetadata | Unary | SyncMetadataRequest | SyncMetadataResponse { synced, conflicts[] } | 60 s |
AddResource — Async Flow
gRPC AddResource
→ KnowledgeService.add_resource()
→ repository.create_document() # PostgreSQL: status=queued
→ queue.publish_ingestion_job() # Redis: INGESTION_JOBS stream
← returns { resource_id, job_id, status="queued" }
(Background — opentier-worker container)
→ handle_ingestion_job() # Redis consumer
→ DocumentProcessor.process() # scrape → chunk → embed → upsert Qdrant
→ repository.update_job_status() # PostgreSQL: status=completed/failedChunkedUpload — Client-Streaming
Supports files larger than the 100 MB gRPC message limit:
- First
FileChunkmust containpayload.metadata: ChunkMetadata { user_id, resource_id, filename, content_type, total_size, total_chunks } - Subsequent chunks contain
payload.data: bytes(max 10 MB per chunk) - Final chunk has
is_last = true - Maximum total file size: 1 GB
- Optional
checksuminChunkMetadatafor integrity verification
Correlation ID Propagation
The Chat servicer extracts or generates a x-correlation-id from gRPC metadata:
# features/chat/grpc_servicer.py
def _get_correlation_id(context: grpc.aio.ServicerContext) -> str:
metadata = dict(context.invocation_metadata())
return metadata.get('x-correlation-id', str(uuid.uuid4()))The correlation ID is included in all log lines within the request scope, enabling distributed tracing across the Rust→Python boundary.
Error Code Mapping
| Python Exception | gRPC Status Code | HTTP (after Rust mapping) |
|---|---|---|
ConversationNotFound | NOT_FOUND | 404 |
PermissionError | PERMISSION_DENIED | 403 |
ValidationError (Pydantic) | INVALID_ARGUMENT | 400 |
RateLimitError | RESOURCE_EXHAUSTED | 429 |
TimeoutError | DEADLINE_EXCEEDED | 504 |
ServiceUnavailable | UNAVAILABLE | 503 |
| All others | INTERNAL | 500 |
Error mapping is centralized in shared/grpc/errors.py.