Communication Model — Protocols & Contracts
Protocol Layers
REST API (Browser ↔ Rust Gateway)
All REST communication uses standard JSON with bearer session authentication:
Content-Type: application/jsonAuthorization: Bearer <64-character session token>
Standardized error envelopes:
{
"message": "Human-readable error description",
"code": "ERROR_CODE_OPTIONAL"
}Server-Sent Events (SSE) Chat Streaming
- Endpoint:
POST /chat/conversations/{id}/stream - Response header:
Content-Type: text/event-stream - Keep-alive heartbeat: Axum sends comment ping
:every 15 seconds to prevent intermediate proxy timeout. - Token delta format:
data: {"conversation_id":"...","message_id":"...","token":{"delta":"Hello "},"is_final":false} - Terminal event includes sources citation array and token usage metrics:
data: {"conversation_id":"...","message_id":"...","sources":{...},"metrics":{...},"is_final":true}
gRPC Contract (Rust ↔ Python)
- Package:
opentier.intelligence.v1 - Transport: HTTP/2 with Protobuf v3
- Message Limit: 100 MB max payload size
- Keepalive: 60 seconds ping intervals
// Core chat service
service Chat {
rpc SendMessage(ChatRequest) returns (ChatResponse);
rpc StreamChat(ChatRequest) returns (stream ChatStreamChunk);
rpc GetConversation(GetConversationRequest) returns (ConversationResponse);
rpc DeleteConversation(DeleteConversationRequest) returns (DeleteConversationResponse);
rpc GenerateTitle(GenerateTitleRequest) returns (GenerateTitleResponse);
}
// Asynchronous resource ingestion
service ResourceService {
rpc AddResource(AddResourceRequest) returns (AddResourceResponse);
rpc ChunkedUpload(stream FileChunk) returns (ChunkedUploadResponse);
rpc GetResourceStatus(GetResourceStatusRequest) returns (ResourceStatusResponse);
rpc ListResources(ListResourcesRequest) returns (ListResourcesResponse);
rpc DeleteResource(DeleteResourceRequest) returns (DeleteResourceResponse);
}Redis Streams Event Envelope (core.events.EventEnvelope)
Decoupled inter-service and background communication uses typed envelopes:
@dataclass(frozen=True)
class EventEnvelope:
event_id: str # Unique event UUID
event_type: str # e.g., 'ingestion.job.enqueued.v1', 'chat.completed.v1'
correlation_id: str # Distributed trace correlation ID
timestamp: str # ISO 8601 UTC timestamp
payload: dict[str, Any] # JSON-serializable event payloadActive Streams
| Stream Name | Producers | Consumers / Groups | Purpose |
|---|---|---|---|
opentier:intel:ingestion_jobs | ResourceService | workers | Web scraping, chunking, and Qdrant upserts. |
opentier:intel:chat_events | ChatService | metering | Deducts tokens from user credit balances. |
opentier:intel:chat_events | ChatService | memory | Asynchronously extracts user facts and updates user_memories. |
opentier:ops:heartbeat | WorkerRuntime | workers | Heartbeat roundtrips and consumer liveness tracking. |
Sync vs Async Communication Workflows
Error Propagation & Dead-Letter Queues (DLQ)
- gRPC Error Mapping: Python exceptions are mapped in
shared/grpc/errors.pyto canonical gRPC status codes (NOT_FOUND,RESOURCE_EXHAUSTED,DEADLINE_EXCEEDED,UNAVAILABLE) and then translated by Axum into standard HTTP status codes (404, 429, 504, 503). - Worker DLQ Routing: If a background consumer fails to process an event after
max_retries = 3attempts, it automatically diverts the envelope to{stream}:dlqand acknowledges the original message, preventing poison-pill blocking of the consumer group.
Last updated on