Skip to Content
Intelligence Service (Python)gRPC Server Interface

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/)

RPCInputOutputNotes
CheckHealthCheckRequest {}HealthCheckResponse { status, version, uptime_seconds }Always returns “healthy” if process is alive
ReadyReadyCheckRequest {}ReadyCheckResponse { ready, dependencies, dependency_status }Checks PostgreSQL + Qdrant + Redis connectivity

Chat Service (features/chat/grpc_servicer.py)

RPCInputOutputNotes
SendMessageChatRequestChatResponseUnary; timeout 1200 s
StreamChatChatRequeststream ChatStreamChunkServer-streaming; timeout 300 s
GetConversationGetConversationRequest { user_id, conversation_id }ConversationResponse { messages[] }Fetches from chat_messages
DeleteConversationDeleteConversationRequest { user_id, conversation_id }DeleteConversationResponse { success }Checks ownership
GenerateTitleGenerateTitleRequest { 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] AddResource is now async in v1.1.0. It publishes the ingestion job to Redis Streams and returns immediately with status=queued. The actual scraping, chunking, embedding and Qdrant upsert happen in the background worker.

RPCTypeInputOutputTimeout
AddResourceUnaryAddResourceRequestAddResourceResponse { resource_id, job_id, status="queued" }30 s
ChunkedUploadClient-streamingstream FileChunkChunkedUploadResponse3000 s
GetResourceStatusUnaryGetResourceStatusRequest { resource_id, user_id }ResourceStatusResponse10 s
ListResourcesUnaryListResourcesRequest { user_id, filters, pagination }ListResourcesResponse { resources[], total, next_cursor }30 s
DeleteResourceUnaryDeleteResourceRequest { resource_id, user_id }DeleteResourceResponse { success }30 s
CancelIngestionUnaryCancelIngestionRequest { job_id, user_id }CancelIngestionResponse { success }10 s
SyncResourceMetadataUnarySyncMetadataRequestSyncMetadataResponse { 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/failed

ChunkedUpload — Client-Streaming

Supports files larger than the 100 MB gRPC message limit:

  1. First FileChunk must contain payload.metadata: ChunkMetadata { user_id, resource_id, filename, content_type, total_size, total_chunks }
  2. Subsequent chunks contain payload.data: bytes (max 10 MB per chunk)
  3. Final chunk has is_last = true
  4. Maximum total file size: 1 GB
  5. Optional checksum in ChunkMetadata for 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 ExceptiongRPC Status CodeHTTP (after Rust mapping)
ConversationNotFoundNOT_FOUND404
PermissionErrorPERMISSION_DENIED403
ValidationError (Pydantic)INVALID_ARGUMENT400
RateLimitErrorRESOURCE_EXHAUSTED429
TimeoutErrorDEADLINE_EXCEEDED504
ServiceUnavailableUNAVAILABLE503
All othersINTERNAL500

Error mapping is centralized in shared/grpc/errors.py.

Last updated on