# How r-nacos Handles Service Discovery Subscriptions and Push Notifications via gRPC Bi-Stream

> Discover how r-nacos leverages gRPC bi-stream for efficient service discovery subscriptions and push notifications. Learn about its active NamingActor and delay buffer for streamlined updates without client polling.

- Repository: [Nacos Group/r-nacos](https://github.com/nacos-group/r-nacos)
- Tags: how-to-guide
- Published: 2026-03-07

---

**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`](https://github.com/nacos-group/r-nacos/blob/main/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.

```rust
// 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`](https://github.com/nacos-group/r-nacos/blob/main/src/grpc/handler/naming_subscribe_service.rs). The handler parses the `SubscribeServiceRequest` protobuf message and translates the boolean `subscribe` field into internal commands.

```rust
// 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`](https://github.com/nacos-group/r-nacos/blob/main/src/naming/core.rs)) delegates subscription state management to the `Subscriber` struct in [`src/naming/naming_subscriber.rs`](https://github.com/nacos-group/r-nacos/blob/main/src/naming/naming_subscriber.rs). The `Subscriber` maintains a `listener` map that links `service_key` to a set of subscribed `client_id` values.

```rust
// 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`](https://github.com/nacos-group/r-nacos/blob/main/src/naming/naming_delay_nofity.rs). When `Subscriber::notify(key)` is invoked, it schedules a `NamingDelayEvent` with a default 500ms delay.

```rust
// 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`](https://github.com/nacos-group/r-nacos/blob/main/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`.

```rust
// 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.

```rust
// 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 `NamingActor`** through the `Subscriber` struct in [`src/naming/naming_subscriber.rs`](https://github.com/nacos-group/r-nacos/blob/main/src/naming/naming_subscriber.rs), which maps service keys to sets of client IDs.
- **Change notifications are coalesced** by the `DelayNotifyActor` in [`src/naming/naming_delay_nofity.rs`](https://github.com/nacos-group/r-nacos/blob/main/src/naming/naming_delay_nofity.rs) using a 500ms delay buffer to prevent notification storms during rapid service changes.
- **Push delivery is orchestrated by `BiStreamManage`** in [`src/grpc/bistream_manage.rs`](https://github.com/nacos-group/r-nacos/blob/main/src/grpc/bistream_manage.rs), which constructs `NotifySubscriberRequest` payloads and broadcasts them to active connections cached in `conn_cache`.
- **Low-level transmission** occurs through `BiStreamConn` in [`src/grpc/bistream_conn.rs`](https://github.com/nacos-group/r-nacos/blob/main/src/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`](https://github.com/nacos-group/r-nacos/blob/main/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`](https://github.com/nacos-group/r-nacos/blob/main/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.