Skip to Content
Communication ModelProtocols & Contracts

Communication Model — Protocols & Contracts

Protocol Layers


REST API (Browser ↔ Rust Gateway)

All REST communication uses standard JSON with bearer session authentication:

  • Content-Type: application/json
  • Authorization: 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 payload

Active Streams

Stream NameProducersConsumers / GroupsPurpose
opentier:intel:ingestion_jobsResourceServiceworkersWeb scraping, chunking, and Qdrant upserts.
opentier:intel:chat_eventsChatServicemeteringDeducts tokens from user credit balances.
opentier:intel:chat_eventsChatServicememoryAsynchronously extracts user facts and updates user_memories.
opentier:ops:heartbeatWorkerRuntimeworkersHeartbeat 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.py to 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 = 3 attempts, it automatically diverts the envelope to {stream}:dlq and acknowledges the original message, preventing poison-pill blocking of the consumer group.
Last updated on