meta_client/client/
config.rs1use std::sync::Arc;
16
17use api::v1::meta::config_client::ConfigClient;
18use api::v1::meta::{PullConfigRequest, PullConfigResponse, RequestHeader, Role};
19use common_grpc::channel_manager::ChannelManager;
20use common_meta::util;
21use common_telemetry::tracing_context::TracingContext;
22use snafu::{OptionExt, ResultExt, ensure};
23use tokio::sync::RwLock;
24use tonic::transport::Channel;
25
26use crate::client::{Id, LeaderProviderRef};
27use crate::error;
28use crate::error::{InvalidResponseHeaderSnafu, Result};
29
30#[derive(Clone, Debug)]
31pub struct Client {
32 inner: Arc<RwLock<Inner>>,
33}
34
35impl Client {
36 pub fn new(id: Id, role: Role, channel_manager: ChannelManager) -> Self {
37 let inner = Arc::new(RwLock::new(Inner::new(id, role, channel_manager)));
38 Self { inner }
39 }
40
41 pub(crate) async fn start_with(&self, leader_provider: LeaderProviderRef) -> Result<()> {
42 let mut inner = self.inner.write().await;
43 inner.start_with(leader_provider)
44 }
45
46 pub async fn pull_config(&self) -> Result<PullConfigResponse> {
47 let inner = self.inner.read().await;
48 inner.ask_leader().await?;
49 inner.pull_config().await
50 }
51}
52
53#[derive(Debug)]
54struct Inner {
55 id: Id,
56 role: Role,
57 channel_manager: ChannelManager,
58 leader_provider: Option<LeaderProviderRef>,
59}
60
61impl Inner {
62 fn new(id: Id, role: Role, channel_manager: ChannelManager) -> Self {
63 Self {
64 id,
65 role,
66 channel_manager,
67 leader_provider: None,
68 }
69 }
70
71 fn start_with(&mut self, leader_provider: LeaderProviderRef) -> Result<()> {
72 ensure!(
73 !self.is_started(),
74 error::IllegalGrpcClientStateSnafu {
75 err_msg: "Config client already started"
76 }
77 );
78 self.leader_provider = Some(leader_provider);
79 Ok(())
80 }
81
82 async fn ask_leader(&self) -> Result<String> {
83 let Some(leader_provider) = self.leader_provider.as_ref() else {
84 return error::IllegalGrpcClientStateSnafu {
85 err_msg: "not started",
86 }
87 .fail();
88 };
89 leader_provider.ask_leader().await
90 }
91
92 async fn pull_config(&self) -> Result<PullConfigResponse> {
93 ensure!(
94 self.is_started(),
95 error::IllegalGrpcClientStateSnafu {
96 err_msg: "Config client not start"
97 }
98 );
99
100 let leader_addr = self
101 .leader_provider
102 .as_ref()
103 .unwrap()
104 .leader()
105 .context(error::NoLeaderSnafu)?;
106 let mut client = self.make_client(&leader_addr)?;
107
108 let header = RequestHeader::new(
109 self.id,
110 self.role,
111 TracingContext::from_current_span().to_w3c(),
112 );
113 let req = PullConfigRequest {
114 header: Some(header),
115 };
116
117 let res = client
118 .pull_config(req)
119 .await
120 .map_err(error::Error::from)?
121 .into_inner();
122
123 util::check_response_header(res.header.as_ref()).context(InvalidResponseHeaderSnafu)?;
124
125 Ok(res)
126 }
127
128 fn make_client(&self, addr: impl AsRef<str>) -> Result<ConfigClient<Channel>> {
129 let channel = self
130 .channel_manager
131 .get(addr)
132 .context(error::CreateChannelSnafu)?;
133
134 Ok(common_grpc::configure_tonic_client!(
135 ConfigClient::new(channel),
136 self.channel_manager,
137 ))
138 }
139
140 #[inline]
141 fn is_started(&self) -> bool {
142 self.leader_provider.is_some()
143 }
144}