nacos_rs_sdk/api/
instance.rs1use crate::client::NacosClient;
2use crate::model::instance::{
3 InstanceBeat, InstanceBeatOption, InstanceObject, QueryInstanceOption, RegisterInstanceOption,
4};
5use serde_json;
6use std::error::Error;
13use std::time::Duration;
14use tokio::{task, time};
15
16impl InstanceObject {
17 pub async fn hart(&self, client: &NacosClient, options: &Option<RegisterInstanceOption>) {
18 task::spawn(hart_beat_thread(
19 self.clone(),
20 client.clone(),
21 options.clone(),
22 ));
23 }
24}
25
26impl InstanceBeat {
27 pub async fn hart(
28 &self,
29 client: &NacosClient,
30 instance: &InstanceObject,
31 options: &Option<InstanceBeatOption>,
32 ) -> Result<InstanceBeatOption, Box<dyn Error + Send + Sync>> {
33 let res = client.beat(&self, &instance, options).await?;
34 Ok(res.json::<InstanceBeatOption>().await?)
35 }
36}
37
38impl QueryInstanceOption {
39 pub fn from(register: &RegisterInstanceOption) -> Self {
40 Self {
41 namespace_id: register.namespace_id.clone(),
42 group_name: register.group_name.clone(),
43 ephemeral: register.ephemeral.clone(),
44 cluster: register.cluster_name.clone(),
45 healthy_only: Some(false),
46 }
47 }
48}
49
50async fn hart_beat_thread(
51 instance: InstanceObject,
52 client: NacosClient,
53 instance_options: Option<RegisterInstanceOption>,
54) {
55 let mut beat = InstanceBeat::builder().build().unwrap();
56 let mut options_builder = InstanceBeatOption::builder();
57 if let Some(instance_option) = instance_options.clone() {
58 if let Some(namespace_id) = instance_option.namespace_id() {
59 options_builder = options_builder.namespace_id(namespace_id);
60 }
61 if let Some(group_name) = instance_option.group_name() {
62 options_builder = options_builder.group_name(group_name);
63 }
64 }
65 let options = options_builder.build().unwrap();
66 let ins_option = QueryInstanceOption::from(&instance_options.clone().unwrap());
67 let config = client
68 .instance_with_object(&instance, &Some(ins_option.clone()))
69 .await
70 .unwrap();
71 let mut bt: Option<String> = None;
72 'hb: loop {
73 match &bt {
74 None => (),
75 Some(s) => {
76 beat.set_beat(&s);
77 }
78 };
79 let br = beat.hart(&client, &instance, &Some(options.clone())).await;
80 match br {
81 Ok(o) => {
82 if !o.light_beat_enabled.unwrap() {
83 bt = Some(serde_json::to_string(&config.clone()).unwrap());
85 }
86 let delay = if o.client_beat_interval.unwrap() > 2 {
87 o.client_beat_interval.unwrap() - 2
88 } else {
89 o.client_beat_interval.unwrap()
90 };
91 time::sleep(Duration::from_millis(delay as u64)).await;
92 }
93 Err(e) => {
94 println!(" -- hart beat err : {:?}", e);
95 break 'hb;
96 }
97 }
98 }
99}