How r-nacos Handles Service Discovery Subscriptions and Push Notifications via gRPC Bi-Stream
r-nacos implements service discovery subscriptions by maintaining a long-lived gRPC bi-directional stream (BiStreamConn) that delegates subscription state to a NamingActor, coalesces change notifications through a 500ms delay buffer, and pushes ServiceInfo snapshots to clients via BiStreamManage without requiring client polling.
The r-nacos project is a high-performance Rust implementation of the Nacos service discovery and configuration management platform. When clients subscribe to service changes, the system leverages gRPC bi-streaming capabilities to establish persistent connections that support real-time push notifications, eliminating the need for repeated polling and reducing network overhead.
The gRPC Bi-Stream Architecture
Clients initiate service discovery subscriptions by establishing a long-lived bi-directional gRPC stream. This connection persists for the duration of the client session and serves as the dedicated channel for both subscription requests and subsequent push notifications.
Establishing the Bi-Directional Stream
The core connection abstraction is the BiStreamConn struct defined in src/grpc/bistream_conn.rs. When a client connects, the server creates a BiStreamManage entry that caches the connection metadata and provides a handle to the active stream.
// Conceptual flow from src/grpc/bistream_conn.rs
match cmd {
BiStreamSenderCmd::Send(payload) => {
// Serialise Payload → bytes → write on the HTTP/2 stream
self.sender.send(payload).await?;
}
BiStreamSenderCmd::Close => { /* close stream */ }
}
SubscribeServiceRequest Handler
Incoming subscription requests are processed by SubscribeServiceRequestHandler in src/grpc/handler/naming_subscribe_service.rs. The handler parses the SubscribeServiceRequest protobuf message and translates the boolean subscribe field into internal commands.
// From src/grpc/handler/naming_subscribe_service.rs
let subscribe_cmd = self.build_subscribe_cmd(
request.subscribe,
key.clone(),
request_meta.connection_id.clone(),
);
self.app_data.naming_addr.do_send(subscribe_cmd);
The build_subscribe_cmd method constructs either NamingCmd::Subscribe or NamingCmd::RemoveSubscribe, embedding the service_key and connection_id required to route notifications to the correct client.
The Event-Driven Subscription Flow
Once the NamingActor receives a subscription command, the system maintains an in-memory registry of active subscriptions and triggers notifications whenever service instances change.
NamingActor and Subscriber Management
The NamingActor (defined in src/naming/core.rs) delegates subscription state management to the Subscriber struct in src/naming/naming_subscriber.rs. The Subscriber maintains a listener map that links service_key to a set of subscribed client_id values.
// From src/naming/naming_subscriber.rs
pub fn add_subscribe(&mut self, client_id: Arc<String>, items: Vec<NamingListenerItem>) {
// update client_keys and listener maps
}
pub fn remove_subscribe(&mut self, client_id: Arc<String>, items: Vec<NamingListenerItem>) {
// Remove client_id from specific service_key entries
}
When a service change occurs (instance registration, deregistration, or health-check status change), the NamingActor calls Subscriber::notify(service_key), which initiates the push notification sequence.
Delayed Notification Coalescing
To prevent notification storms during rapid service changes, r-nacos implements a delay-and-coalesce mechanism via DelayNotifyActor in src/naming/naming_delay_nofity.rs. When Subscriber::notify(key) is invoked, it schedules a NamingDelayEvent with a default 500ms delay.
// From src/naming/naming_delay_nofity.rs
let event = NamingDelayEvent {
key,
client_id_set,
service_info: None,
conn_manage: Some(conn_manage), // injected BiStreamManage address
};
self.inner_delay_notify.add_event(self.delay, event.key.clone(), event)?;
The delay actor aggregates multiple rapid changes to the same service key into a single notification event. When the delay expires, the actor fetches the latest ServiceInfo from NamingActor and invokes event.on_event(), which forwards the notification to BiStreamManage for delivery.
Push Notification Delivery Mechanism
The final stage converts internal events into gRPC push messages and transmits them over the persistent bi-directional streams.
BiStreamManage Broadcasting
BiStreamManage (defined in src/grpc/bistream_manage.rs) receives BiStreamManageCmd::NotifyNaming commands containing the service key, target client IDs, and current ServiceInfo. It constructs a NotifySubscriberRequest protobuf message and retrieves the active connection for each subscribed client from its internal conn_cache.
// From src/grpc/bistream_manage.rs
BiStreamManageCmd::NotifyNaming(service_key, client_id_set, service_info) => {
let request = NotifySubscriberRequest {
service_info: Some(service_info.into()),
// ... other fields
};
let payload = Arc::new(PayloadUtils::build_payload(
"NotifySubscriberRequest",
serde_json::to_string(&request)?,
));
for client_id in client_id_set {
if let Some(item) = self.conn_cache.get(&client_id) {
item.conn.do_send(BiStreamSenderCmd::Send(payload.clone()));
}
}
}
The NotifySubscriberRequest contains a complete snapshot of the service instances, allowing clients to update their local caches atomically.
BiStreamConn and gRPC Payload Transmission
Each active client connection is represented by a BiStreamConn instance that holds the gRPC sender channel. When BiStreamManage dispatches a BiStreamSenderCmd::Send command, the connection serializes the payload and writes it to the HTTP/2 stream.
// From src/grpc/bistream_conn.rs
match cmd {
BiStreamSenderCmd::Send(payload) => {
// Serialize to bytes and push to the open TCP/HTTP2 stream
self.sender.send(payload).await?;
}
}
Because the gRPC bi-stream remains open indefinitely, the client receives the NotifySubscriberResponse immediately as a server-side push, completing the subscription notification cycle without requiring client polling.
Summary
- r-nacos uses a persistent gRPC bi-directional stream (
BiStreamConn) to maintain long-lived connections with clients, enabling real-time server push capabilities. - Subscription state is managed by the
NamingActorthrough theSubscriberstruct insrc/naming/naming_subscriber.rs, which maps service keys to sets of client IDs. - Change notifications are coalesced by the
DelayNotifyActorinsrc/naming/naming_delay_nofity.rsusing a 500ms delay buffer to prevent notification storms during rapid service changes. - Push delivery is orchestrated by
BiStreamManageinsrc/grpc/bistream_manage.rs, which constructsNotifySubscriberRequestpayloads and broadcasts them to active connections cached inconn_cache. - Low-level transmission occurs through
BiStreamConninsrc/grpc/bistream_conn.rs, which serializes payloads and writes them to the open HTTP/2 stream, completing the push notification flow.
Frequently Asked Questions
How does r-nacos prevent clients from being overwhelmed by rapid service changes?
The system implements a delay-and-coalesce mechanism via the DelayNotifyActor in src/naming/naming_delay_nofity.rs. When service changes occur, notifications are delayed by 500ms (configurable) and aggregated, ensuring that multiple rapid updates to the same service result in a single push notification containing the latest state.
What happens when a client disconnects unexpectedly?
The BiStreamConn detects the dropped HTTP/2 stream and triggers cleanup via BiStreamManageCmd. The BiStreamManage removes the client ID from its conn_cache, and subsequent NotifyNaming commands will skip disconnected clients. The NamingActor eventually cleans up orphaned subscriptions through timeout mechanisms.
Can clients subscribe to multiple services on a single gRPC connection?
Yes. The Subscriber in src/naming/naming_subscriber.rs maintains a listener map where each service key points to a set of client IDs. A single connection_id (client ID) can appear in multiple service key sets, allowing one bi-directional stream to receive push notifications for numerous services simultaneously.
How does r-nacos ensure message ordering during push notifications?
The gRPC bi-stream guarantees ordered delivery within a single HTTP/2 stream. The BiStreamConn processes BiStreamSenderCmd::Send commands sequentially, and the DelayNotifyActor ensures that notifications for the same service key are coalesced and sent as atomic snapshots, preventing out-of-order state updates.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →