Skip to main content

nacos_rs_sdk/api/
instance.rs

1use crate::client::NacosClient;
2use crate::model::instance::{
3    InstanceBeat, InstanceBeatOption, InstanceObject, QueryInstanceOption, RegisterInstanceOption,
4};
5use serde_json;
6// use crate::model::Put;
7// use std::borrow::Borrow;
8// use nacos_rs_sdk_macro::{Delete, Get, Nacos, Post, Put, Value, Builder};
9// use reqwest::Response;
10// use serde::{Deserialize, Serialize};
11// use std::sync::{Arc, RwLock};
12use 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(config.clone());
84                    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}