use crate::app_state::AppState;
use crate::client::stream::ClientStreamReq;
use crate::{Client, Error, Node, NodeId};
use openraft::RaftMetrics;
use std::clone::Clone;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::time::Duration;
use tokio::sync::RwLock;
use tokio::time;
use tracing::{debug, error, warn};
impl Client {
#[inline(always)]
pub(crate) async fn build_addr(
&self,
path: &str,
leader: &Arc<RwLock<(NodeId, String)>>,
) -> String {
let scheme = if self.inner.tls_config.is_some() {
"https"
} else {
"http"
};
let url = {
let lock = leader.read().await;
format!("{}://{}{}", scheme, lock.1, path)
};
debug!("request url: {}", url);
url
}
pub(crate) async fn find_set_active_leader(&self) {
if let Some(state) = &self.inner.state {
#[cfg(feature = "sqlite")]
{
let metrics = state.raft_db.raft.metrics().borrow().clone();
let mut find_leader = Self::find_set_leader(metrics, &self.inner.leader_db).await;
while let Err(err) = find_leader {
warn!("Find DB leader error: {}", err);
time::sleep(Duration::from_millis(500)).await;
let metrics = state.raft_db.raft.metrics().borrow().clone();
find_leader = Self::find_set_leader(metrics, &self.inner.leader_db).await;
}
}
#[cfg(feature = "cache")]
{
let metrics = state.raft_cache.raft.metrics().borrow().clone();
let mut find_leader =
Self::find_set_leader(metrics, &self.inner.leader_cache).await;
while let Err(err) = find_leader {
warn!("Find cache leader error: {}", err);
time::sleep(Duration::from_millis(500)).await;
let metrics = state.raft_cache.raft.metrics().borrow().clone();
find_leader = Self::find_set_leader(metrics, &self.inner.leader_cache).await;
}
}
} else {
#[cfg(feature = "sqlite")]
{
let mut metrics = self.remote_metrics_loop_db().await;
loop {
match Self::find_set_leader(metrics, &self.inner.leader_db).await {
Ok(_) => {
break;
}
Err(_) => {
metrics = self.remote_metrics_loop_db().await;
}
}
}
}
#[cfg(feature = "cache")]
{
let mut metrics = self.remote_metrics_loop_cache().await;
loop {
match Self::find_set_leader(metrics, &self.inner.leader_cache).await {
Ok(_) => {
break;
}
Err(_) => {
metrics = self.remote_metrics_loop_cache().await;
}
}
}
}
}
}
#[cfg(feature = "cache")]
async fn remote_metrics_loop_cache(&self) -> RaftMetrics<NodeId, Node> {
loop {
for addr in &self.inner.nodes {
{
let mut lock = self.inner.leader_cache.write().await;
*lock = (lock.0, addr.clone());
}
match self.metrics_cache().await {
Ok(metrics) => {
return metrics;
}
Err(err) => {
error!("Error looking up Cache metrics: {}", err);
}
}
}
time::sleep(Duration::from_millis(500)).await;
}
}
#[cfg(feature = "sqlite")]
async fn remote_metrics_loop_db(&self) -> RaftMetrics<NodeId, Node> {
loop {
for addr in &self.inner.nodes {
{
let mut lock = self.inner.leader_db.write().await;
*lock = (lock.0, addr.clone());
}
match self.metrics_db().await {
Ok(metrics) => {
return metrics;
}
Err(err) => {
error!("Error looking up DB metrics: {}", err);
}
}
}
time::sleep(Duration::from_millis(500)).await;
}
}
async fn find_set_leader(
metrics: RaftMetrics<NodeId, Node>,
leader: &Arc<RwLock<(NodeId, String)>>,
) -> Result<(), Error> {
let leader_id = match metrics.current_leader {
None => {
return Err(Error::Connect("Leader vote is in progress".to_string()));
}
Some(leader_id) => leader_id,
};
let leader_filtered = metrics
.membership_config
.nodes()
.filter(|(id, _)| *id == &leader_id)
.collect::<Vec<_>>();
assert_eq!(leader_filtered.len(), 1);
let mut lock = leader.write().await;
*lock = (*leader_filtered[0].0, leader_filtered[0].1.addr_api.clone());
Ok(())
}
#[cfg(feature = "sqlite")]
pub async fn is_leader_db(&self) -> bool {
if let Some(state) = &self.inner.state
&& state.id == self.inner.leader_db.read().await.0
{
return true;
}
false
}
#[cfg(feature = "cache")]
pub async fn is_leader_cache(&self) -> bool {
if let Some(state) = &self.inner.state
&& state.id == self.inner.leader_cache.read().await.0
{
return true;
}
false
}
#[cfg(feature = "sqlite")]
#[inline(always)]
pub(crate) async fn is_leader_db_with_state(&self) -> Option<&Arc<AppState>> {
if let Some(state) = &self.inner.state
&& state.id == self.inner.leader_db.read().await.0
{
return Some(state);
}
None
}
#[cfg(feature = "cache")]
#[inline(always)]
pub(crate) async fn is_leader_cache_with_state(&self) -> Option<&Arc<AppState>> {
if let Some(state) = &self.inner.state
&& state.id == self.inner.leader_cache.read().await.0
{
return Some(state);
}
None
}
#[cfg(not(feature = "dashboard"))]
#[inline(always)]
pub(crate) fn new_request_id(&self) -> usize {
self.inner.request_id.fetch_add(1, Ordering::Relaxed)
}
#[cfg(feature = "dashboard")]
#[inline(always)]
pub(crate) fn new_request_id(&self) -> usize {
if let Some(st) = &self.inner.state {
st.new_request_id()
} else {
self.inner.request_id.fetch_add(1, Ordering::Relaxed)
}
}
#[inline]
pub(crate) async fn was_leader_update_error(
&self,
err: &Error,
lock: &Arc<RwLock<(NodeId, String)>>,
tx: &flume::Sender<ClientStreamReq>,
) -> bool {
let mut was_leader_error = false;
if let Some((id, node)) = err.is_forward_to_leader()
&& let Some(leader_id) = id
&& let Some(node) = node
{
was_leader_error = true;
let api_addr = node.addr_api.clone();
{
let mut lock = lock.write().await;
if lock.0 != leader_id {
*lock = (leader_id, api_addr.clone());
}
}
if was_leader_error {
tx.send_async(ClientStreamReq::LeaderChange((id, Some(node.clone()))))
.await
.expect("the Client API WebSocket Manager to always be running");
}
}
was_leader_error
}
}