How r-nacos Implements Config Change Notifications to Subscribed Clients
r-nacos implements config change notifications using an Actix actor-based pub/sub architecture that processes gRPC subscription requests through ConfigActor and pushes real-time updates via BiStreamManage to connected clients.
r-nacos is a Rust implementation of the Nacos configuration and service discovery platform. When clients need to react to configuration updates in real-time, r-nacos config change notifications ensure immediate delivery through an internal messaging system. This article examines the complete flow from subscription registration to push delivery, referencing the actual source implementation in the nacos-group/r-nacos repository.
Architecture Overview of r-nacos Config Change Notifications
The notification system relies on three primary components coordinated through Actix actors:
- ConfigActor (
src/config/core.rs): The central configuration cache and command processor - Subscriber (
src/config/config_subscribe.rs): The registry maintaining client-to-config-key mappings - BiStreamManage (
src/grpc/bistream_manage.rs): The manager handling active gRPC bi-directional streams
When a configuration changes, the system traverses this chain: ConfigActor detects the change → Subscriber identifies affected clients → BiStreamManage pushes the payload through persistent gRPC streams.
Stage 1: Handling Client Subscription Requests
Client subscriptions begin with a gRPC ConfigBatchListenRequest sent to the ConfigBatchListen handler.
Parsing the gRPC Request
The handler at src/grpc/handler/config_change_batch_listen.rs processes incoming requests:
// Inside ConfigChangeBatchListenRequestHandler::handle
let cmd = if request.listen {
ConfigCmd::Subscribe(listener_items, request_meta.connection_id)
} else {
ConfigCmd::RemoveSubscribe(listener_items, request_meta.connection_id)
};
self.app_data.config_addr.send(cmd).await?;
Each config_listen_context in the request becomes a ListenerItem containing the ConfigKey (data_id, group, tenant) and the client's current MD5 hash. The handler forwards these as ConfigCmd messages to the ConfigActor.
Stage 2: Registering Subscribers in ConfigActor
When ConfigActor receives a subscription command, it updates the internal subscriber registry.
The Subscriber Registry Structure
Located in src/config/config_subscribe.rs, the Subscriber struct maintains bidirectional mappings:
// Subscriber maintains:
listener: HashMap<ConfigKey, HashSet<Arc<String>>>, // key -> client IDs
client_keys: HashMap<Arc<String>, HashSet<ConfigKey>>; // client ID -> keys
The Subscriber::add_subscribe method populates these maps when ConfigActor calls it from src/config/core.rs:
// Inside ConfigActor's command handler
ConfigCmd::Subscribe(items, client_id) => {
self.subscriber.add_subscribe(client_id, items);
}
During actor initialization (injection phase), ConfigActor receives the BiStreamManage address and stores it in the subscriber (self.subscriber.set_conn_manage(conn_manage)), enabling downstream push notifications.
Stage 3: Detecting Changes and Pushing Notifications
When configuration data is modified, the system triggers the notification chain.
Change Detection in ConfigActor
After updating the configuration cache, ConfigActor::set_config (or del_config) invokes:
// Notify different listener types
self.listener.notify(param.key.clone()); // Long-polling listeners
self.subscriber.notify(param.key); // Push subscribers (gRPC streams)
The Notification Push Flow
The Subscriber::notify method (lines 119-127 in src/config/config_subscribe.rs) forwards the notification to BiStreamManage:
pub fn notify(&self, key: ConfigKey) {
if let Some(conn_manage) = &self.conn_manage {
if let Some(client_ids) = self.listener.get(&key) {
conn_manage.do_send(BiStreamManageCmd::NotifyConfig(key, client_ids.clone()));
}
}
}
In src/grpc/bistream_manage.rs (lines 274-306), BiStreamManage constructs the payload:
// Building the notification payload
let mut request = ConfigChangeNotifyRequest {
group: config_key.group,
data_id: config_key.data_id,
tenant: config_key.tenant,
request_id: Some(self.next_request_id()),
module: Some(CONFIG_MODEL.to_string()),
..Default::default()
};
let payload = Arc::new(PayloadUtils::build_payload(
"ConfigChangeNotifyRequest",
serde_json::to_string(&request)?,
));
// Push to all subscribed clients
for client_id in client_id_set {
if let Some(conn) = self.conn_cache.get(client_id) {
conn.conn.do_send(BiStreamSenderCmd::Send(payload.clone()));
}
}
The system handles special tenancy logic (lines 92-103) to ensure compatibility with default-namespace clients before transmission.
Key Source Files and Components
Understanding r-nacos config change notifications requires familiarity with these specific implementation files:
src/config/core.rs: ContainsConfigActorwhich processes subscription commands and triggers notifications after config modificationssrc/config/config_subscribe.rs: Implements theSubscriberstruct that manages client-to-key mappings and forwards notifications toBiStreamManagesrc/grpc/handler/config_change_batch_listen.rs: Handles incomingConfigBatchListenRequestgRPC calls and converts them to internal commandssrc/grpc/bistream_manage.rs: Manages active bi-directional gRPC streams and constructsConfigChangeNotifyRequestpayloads for push deliverysrc/grpc/api_model.rs: Defines theConfigChangeNotifyRequeststructure used in notification payloads
Summary
r-nacos implements config change notifications through a carefully orchestrated actor system:
- Subscription handling: Clients send
ConfigBatchListenRequestvia gRPC, processed byConfigChangeBatchListenRequestHandlerinsrc/grpc/handler/config_change_batch_listen.rs - Registry management:
ConfigActormaintains subscriptions through theSubscriberstruct insrc/config/config_subscribe.rs, tracking which client IDs listen to which configuration keys - Push delivery: When
set_configordel_configis called,ConfigActortriggersSubscriber::notify, which forwards toBiStreamManageinsrc/grpc/bistream_manage.rsto construct and sendConfigChangeNotifyRequestpayloads over persistent gRPC bi-directional streams
This architecture separates concerns between connection management, configuration storage, and notification delivery while maintaining high concurrency through Actix actors.
Frequently Asked Questions
What is the difference between long-polling listeners and push subscribers in r-nacos?
r-nacos maintains two distinct notification mechanisms. Long-polling listeners are handled by the self.listener.notify() call in ConfigActor, which responds to open HTTP long-polling connections. Push subscribers use the self.subscriber.notify() path, which leverages the gRPC bi-directional stream infrastructure to actively push ConfigChangeNotifyRequest payloads to connected clients without requiring a pending request.
How does r-nacos handle client disconnections in the notification system?
When a client disconnects, the BiStreamManage actor detects the broken gRPC stream and removes the entry from conn_cache. The Subscriber struct maintains a client_keys HashMap that maps client IDs to their subscribed keys. While the analysis does not show explicit cleanup code in the provided snippets, the bidirectional mapping structure (listener and client_keys) in src/config/config_subscribe.rs enables efficient removal of stale subscriptions when the system detects a disconnection.
What protocol does r-nacos use for config change notifications?
r-nacos uses gRPC bi-directional streaming for config change notifications. The system constructs a ConfigChangeNotifyRequest payload (defined in src/grpc/api_model.rs) and sends it over persistent gRPC streams managed by BiStreamManage in src/grpc/bistream_manage.rs. This approach allows the server to push updates immediately when configuration changes occur, rather than requiring clients to poll.
How does r-nacos ensure compatibility with different namespace configurations when sending notifications?
The BiStreamManage actor in src/grpc/bistream_manage.rs contains specific logic (lines 92-103) to handle tenancy compatibility. When constructing the ConfigChangeNotifyRequest, the system checks if the tenant matches the default namespace and adjusts the payload accordingly. This ensures that clients using the default namespace receive properly formatted notifications compatible with the Nacos protocol specification, while maintaining correct isolation for multi-tenant scenarios.
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 →