use crate::{
connection::blocking::FalkorSyncConnection,
parser::{redis_value_as_string, redis_value_as_vec},
FalkorDBError, FalkorResult,
};
use std::collections::HashMap;
use std::num::NonZeroU8;
#[cfg(feature = "tokio")]
use std::num::NonZeroUsize;
#[cfg(feature = "tokio")]
use crate::connection::asynchronous::FalkorAsyncConnection;
pub(crate) mod blocking;
pub(crate) mod builder;
#[cfg(feature = "tokio")]
pub(crate) mod asynchronous;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ConnectionStrategy {
Pooled {
size: NonZeroU8,
},
Multiplexed {
connections: NonZeroU8,
},
}
impl ConnectionStrategy {
pub fn connection_count(&self) -> NonZeroU8 {
match self {
ConnectionStrategy::Pooled { size } => *size,
ConnectionStrategy::Multiplexed { connections } => *connections,
}
}
pub(crate) fn with_connection_count(
self,
count: NonZeroU8,
) -> Self {
match self {
ConnectionStrategy::Pooled { .. } => ConnectionStrategy::Pooled { size: count },
ConnectionStrategy::Multiplexed { .. } => {
ConnectionStrategy::Multiplexed { connections: count }
}
}
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub enum ReadPreference {
#[default]
Primary,
PreferReplica,
}
#[allow(clippy::large_enum_variant)]
pub(crate) enum FalkorClientProvider {
#[cfg(test)]
None,
Redis {
client: redis::Client,
sentinel: Option<redis::sentinel::SentinelClient>,
sentinel_replica: Option<redis::sentinel::SentinelClient>,
#[cfg(feature = "embedded-core")]
#[allow(dead_code)]
embedded_server: Option<std::sync::Arc<crate::embedded::EmbeddedServer>>,
response_timeout: Option<std::time::Duration>,
},
}
impl FalkorClientProvider {
pub(crate) fn get_connection(&mut self) -> FalkorResult<FalkorSyncConnection> {
Ok(match self {
FalkorClientProvider::Redis {
sentinel: Some(sentinel),
..
} => FalkorSyncConnection::Redis(
sentinel
.get_connection()
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
),
FalkorClientProvider::Redis { client, .. } => FalkorSyncConnection::Redis(
client
.get_connection()
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
),
#[cfg(test)]
FalkorClientProvider::None => Err(FalkorDBError::UnavailableProvider)?,
})
}
#[cfg(feature = "tokio")]
fn async_connection_config(
response_timeout: Option<std::time::Duration>
) -> redis::AsyncConnectionConfig {
redis::AsyncConnectionConfig::new().set_response_timeout(response_timeout)
}
#[cfg(feature = "tokio")]
pub(crate) async fn get_async_connection(&mut self) -> FalkorResult<FalkorAsyncConnection> {
Ok(match self {
FalkorClientProvider::Redis {
sentinel: Some(sentinel),
response_timeout,
..
} => FalkorAsyncConnection::Redis(
sentinel
.get_async_connection_with_config(&Self::async_connection_config(
*response_timeout,
))
.await
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
),
FalkorClientProvider::Redis {
client,
response_timeout,
..
} => FalkorAsyncConnection::Redis(
client
.get_multiplexed_async_connection_with_config(&Self::async_connection_config(
*response_timeout,
))
.await
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
),
#[cfg(test)]
FalkorClientProvider::None => Err(FalkorDBError::UnavailableProvider)?,
})
}
pub(crate) fn get_replica_connection(&mut self) -> FalkorResult<FalkorSyncConnection> {
match self {
FalkorClientProvider::Redis {
sentinel_replica: Some(replica),
..
} => Ok(FalkorSyncConnection::Redis(
replica
.get_connection()
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
)),
_ => Err(FalkorDBError::UnavailableProvider),
}
}
#[cfg(feature = "tokio")]
pub(crate) async fn get_async_replica_connection(
&mut self
) -> FalkorResult<FalkorAsyncConnection> {
match self {
FalkorClientProvider::Redis {
sentinel_replica: Some(replica),
response_timeout,
..
} => Ok(FalkorAsyncConnection::Redis(
replica
.get_async_connection_with_config(&Self::async_connection_config(
*response_timeout,
))
.await
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
)),
_ => Err(FalkorDBError::UnavailableProvider),
}
}
pub(crate) fn has_sentinel_replica(&self) -> bool {
matches!(
self,
FalkorClientProvider::Redis {
sentinel_replica: Some(_),
..
}
)
}
#[cfg(feature = "tokio")]
pub(crate) fn has_sentinel(&self) -> bool {
matches!(
self,
FalkorClientProvider::Redis {
sentinel: Some(_),
..
}
)
}
#[cfg(feature = "tokio")]
pub(crate) async fn get_async_connection_manager(
&mut self,
max_inflight: Option<NonZeroUsize>,
) -> FalkorResult<FalkorAsyncConnection> {
let response_timeout = self.response_timeout();
let client = match self {
FalkorClientProvider::Redis {
sentinel: Some(sentinel),
..
} => sentinel
.async_get_client()
.await
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?,
FalkorClientProvider::Redis { client, .. } => client.clone(),
#[cfg(test)]
FalkorClientProvider::None => return Err(FalkorDBError::UnavailableProvider),
};
Self::manager_from_client(client, max_inflight, response_timeout).await
}
#[cfg(feature = "tokio")]
pub(crate) async fn get_async_replica_connection_manager(
&mut self,
max_inflight: Option<NonZeroUsize>,
) -> FalkorResult<FalkorAsyncConnection> {
let response_timeout = self.response_timeout();
match self {
FalkorClientProvider::Redis {
sentinel_replica: Some(replica),
..
} => {
let client = replica
.async_get_client()
.await
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?;
Self::manager_from_client(client, max_inflight, response_timeout).await
}
_ => Err(FalkorDBError::UnavailableProvider),
}
}
#[cfg(feature = "tokio")]
fn response_timeout(&self) -> Option<std::time::Duration> {
match self {
FalkorClientProvider::Redis {
response_timeout, ..
} => *response_timeout,
#[cfg(test)]
FalkorClientProvider::None => None,
}
}
#[cfg(feature = "tokio")]
async fn manager_from_client(
client: redis::Client,
max_inflight: Option<NonZeroUsize>,
response_timeout: Option<std::time::Duration>,
) -> FalkorResult<FalkorAsyncConnection> {
let config =
redis::aio::ConnectionManagerConfig::new().set_response_timeout(response_timeout);
let config = match max_inflight {
Some(limit) => config.set_concurrency_limit(limit.get()),
None => config,
};
let manager = redis::aio::ConnectionManager::new_with_config(client, config)
.await
.map_err(|err| FalkorDBError::RedisError(err.to_string()))?;
Ok(FalkorAsyncConnection::Managed(manager))
}
pub(crate) fn set_sentinel(
&mut self,
sentinel_client: redis::sentinel::SentinelClient,
) {
match self {
FalkorClientProvider::Redis { sentinel, .. } => *sentinel = Some(sentinel_client),
#[cfg(test)]
FalkorClientProvider::None => {}
}
}
pub(crate) fn set_sentinel_replica(
&mut self,
sentinel_client: redis::sentinel::SentinelClient,
) {
match self {
FalkorClientProvider::Redis {
sentinel_replica, ..
} => *sentinel_replica = Some(sentinel_client),
#[cfg(test)]
FalkorClientProvider::None => {}
}
}
#[cfg(test)]
pub(crate) fn get_sentinel_client_common(
&self,
connection_info: &redis::ConnectionInfo,
sentinel_masters: Vec<redis::Value>,
) -> FalkorResult<Option<redis::sentinel::SentinelClient>> {
self.build_sentinel_client(
connection_info,
sentinel_masters,
redis::sentinel::SentinelServerType::Master,
)
}
fn build_sentinel_client(
&self,
connection_info: &redis::ConnectionInfo,
sentinel_masters: Vec<redis::Value>,
server_type: redis::sentinel::SentinelServerType,
) -> FalkorResult<Option<redis::sentinel::SentinelClient>> {
if sentinel_masters.len() != 1 {
return Err(FalkorDBError::SentinelMastersCount);
}
let sentinel_master: HashMap<_, _> = sentinel_masters
.into_iter()
.next()
.and_then(|master| master.into_sequence().ok())
.ok_or(FalkorDBError::SentinelMastersCount)?
.chunks_exact(2)
.flat_map(TryInto::<&[redis::Value; 2]>::try_into) .flat_map(|[key, val]| {
redis_value_as_string(key.to_owned())
.and_then(|key| redis_value_as_string(val.to_owned()).map(|val| (key, val)))
})
.collect();
let name = sentinel_master
.get("name")
.ok_or(FalkorDBError::SentinelMastersCount)?;
Ok(Some(
redis::sentinel::SentinelClient::build(
vec![connection_info.to_owned()],
name.clone(),
Some({
let node_info = redis::sentinel::SentinelNodeConnectionInfo::default()
.set_redis_connection_info(connection_info.redis_settings().clone());
match connection_info.addr() {
redis::ConnectionAddr::TcpTls { insecure: true, .. } => {
node_info.set_tls_mode(redis::TlsMode::Insecure)
}
redis::ConnectionAddr::TcpTls {
insecure: false, ..
} => node_info.set_tls_mode(redis::TlsMode::Secure),
_ => node_info,
}
}),
server_type,
)
.map_err(|err| FalkorDBError::SentinelConnection(err.to_string()))?,
))
}
#[cfg_attr(
feature = "tracing",
tracing::instrument(name = "Get Sentinel Clients", skip_all, level = "info")
)]
pub(crate) fn get_sentinel_client(
&mut self,
connection_info: &redis::ConnectionInfo,
) -> FalkorResult<Option<SentinelClients>> {
let mut conn = self.get_connection()?;
if !conn.check_is_redis_sentinel()? {
return Ok(None);
}
let sentinel_masters = conn
.execute_command(None, "SENTINEL", Some("MASTERS"), None)
.and_then(redis_value_as_vec)?;
self.build_sentinel_clients(connection_info, sentinel_masters)
}
#[cfg(feature = "tokio")]
#[cfg_attr(
feature = "tracing",
tracing::instrument(name = "Get Sentinel Clients", skip_all, level = "info")
)]
pub(crate) async fn get_sentinel_client_async(
&mut self,
connection_info: &redis::ConnectionInfo,
) -> FalkorResult<Option<SentinelClients>> {
let mut conn = self.get_async_connection().await?;
if !conn.check_is_redis_sentinel().await? {
return Ok(None);
}
let sentinel_masters = conn
.execute_command(None, "SENTINEL", Some("MASTERS"), None)
.await
.and_then(redis_value_as_vec)?;
self.build_sentinel_clients(connection_info, sentinel_masters)
}
fn build_sentinel_clients(
&self,
connection_info: &redis::ConnectionInfo,
sentinel_masters: Vec<redis::Value>,
) -> FalkorResult<Option<SentinelClients>> {
let master = self
.build_sentinel_client(
connection_info,
sentinel_masters.clone(),
redis::sentinel::SentinelServerType::Master,
)?
.ok_or(FalkorDBError::SentinelMastersCount)?;
let replica = self.build_sentinel_client(
connection_info,
sentinel_masters,
redis::sentinel::SentinelServerType::Replica,
)?;
Ok(Some(SentinelClients { master, replica }))
}
}
pub(crate) struct SentinelClients {
pub(crate) master: redis::sentinel::SentinelClient,
pub(crate) replica: Option<redis::sentinel::SentinelClient>,
}
pub(crate) trait ProvidesSyncConnections: Sync + Send {
fn get_connection(&self) -> FalkorResult<FalkorSyncConnection>;
}
#[cfg(test)]
mod tests {
use super::*;
use std::str::FromStr;
#[test]
fn test_falkor_client_provider_none_connection() {
let mut provider = FalkorClientProvider::None;
let result = provider.get_connection();
assert!(result.is_err());
if let Err(e) = result {
assert!(matches!(e, FalkorDBError::UnavailableProvider));
}
}
#[test]
fn test_has_sentinel_replica_default() {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
assert!(!provider.has_sentinel_replica());
}
#[test]
fn test_get_replica_connection_errors_without_replica() {
let mut provider = FalkorClientProvider::None;
assert!(!provider.has_sentinel_replica());
let result = provider.get_replica_connection();
assert!(matches!(result, Err(FalkorDBError::UnavailableProvider)));
}
#[test]
fn test_set_sentinel_replica() {
let mut provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:26379").unwrap();
let replica = redis::sentinel::SentinelClient::build(
vec![connection_info],
"master".to_string(),
None,
redis::sentinel::SentinelServerType::Replica,
)
.unwrap();
provider.set_sentinel_replica(replica);
assert!(!provider.has_sentinel_replica());
}
#[test]
#[cfg(feature = "tokio")]
fn test_falkor_client_provider_none_async_connection() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let mut provider = FalkorClientProvider::None;
let result = provider.get_async_connection().await;
assert!(result.is_err());
if let Err(e) = result {
assert!(matches!(e, FalkorDBError::UnavailableProvider));
}
});
}
#[test]
fn test_falkor_client_provider_set_sentinel() {
let mut provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:26379").unwrap();
let sentinel = redis::sentinel::SentinelClient::build(
vec![connection_info],
"master".to_string(),
None,
redis::sentinel::SentinelServerType::Master,
)
.unwrap();
provider.set_sentinel(sentinel);
}
#[test]
fn test_get_sentinel_client_common_invalid_count() {
let provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:6379").unwrap();
let result = provider.get_sentinel_client_common(&connection_info, vec![]);
assert!(result.is_err());
if let Err(e) = result {
assert!(matches!(e, FalkorDBError::SentinelMastersCount));
}
let result = provider.get_sentinel_client_common(
&connection_info,
vec![redis::Value::Nil, redis::Value::Nil],
);
assert!(matches!(result, Err(FalkorDBError::SentinelMastersCount)));
}
fn single_master_reply(name: &str) -> Vec<redis::Value> {
vec![redis::Value::Array(vec![
redis::Value::BulkString(b"name".to_vec()),
redis::Value::BulkString(name.as_bytes().to_vec()),
])]
}
#[test]
fn test_build_sentinel_client_happy_path() {
let provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:6379").unwrap();
let client = provider
.get_sentinel_client_common(&connection_info, single_master_reply("mymaster"))
.expect("master client should build");
assert!(client.is_some());
}
#[test]
fn test_build_sentinel_client_missing_name() {
let provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:6379").unwrap();
let reply = vec![redis::Value::Array(vec![
redis::Value::BulkString(b"ip".to_vec()),
redis::Value::BulkString(b"127.0.0.1".to_vec()),
])];
let result = provider.get_sentinel_client_common(&connection_info, reply);
assert!(matches!(result, Err(FalkorDBError::SentinelMastersCount)));
}
#[test]
fn test_build_sentinel_clients_master_and_replica() {
let provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:6379").unwrap();
let clients = provider
.build_sentinel_clients(&connection_info, single_master_reply("mymaster"))
.expect("clients should build")
.expect("a Sentinel reply yields clients");
assert!(clients.replica.is_some());
}
#[test]
fn test_build_sentinel_clients_invalid_count() {
let provider = FalkorClientProvider::None;
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:6379").unwrap();
let result = provider.build_sentinel_clients(&connection_info, vec![]);
assert!(matches!(result, Err(FalkorDBError::SentinelMastersCount)));
}
#[test]
fn test_set_sentinel_replica_on_redis_provider() {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
assert!(!provider.has_sentinel_replica());
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:26379").unwrap();
let replica = redis::sentinel::SentinelClient::build(
vec![connection_info],
"mymaster".to_string(),
None,
redis::sentinel::SentinelServerType::Replica,
)
.unwrap();
provider.set_sentinel_replica(replica);
assert!(provider.has_sentinel_replica());
}
#[test]
#[cfg(feature = "embedded-core")]
fn test_falkor_client_provider_with_embedded_server() {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let _provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
embedded_server: None,
response_timeout: None,
};
}
#[test]
fn test_falkor_client_provider_redis_without_sentinel() {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let _provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
}
#[test]
fn test_get_replica_connection_errors_when_replica_unreachable() {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:1").unwrap();
let replica = redis::sentinel::SentinelClient::build(
vec![connection_info],
"mymaster".to_string(),
None,
redis::sentinel::SentinelServerType::Replica,
)
.unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: Some(replica),
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_replica_connection();
assert!(
matches!(result, Err(FalkorDBError::RedisError(_))),
"error should come from the replica path"
);
}
#[test]
#[cfg(feature = "tokio")]
fn test_get_async_replica_connection_errors_when_replica_unreachable() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:1").unwrap();
let replica = redis::sentinel::SentinelClient::build(
vec![connection_info],
"mymaster".to_string(),
None,
redis::sentinel::SentinelServerType::Replica,
)
.unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: Some(replica),
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_async_replica_connection().await;
assert!(
matches!(result, Err(FalkorDBError::RedisError(_))),
"error should come from the replica path"
);
});
}
#[test]
fn test_get_replica_connection_on_redis_provider_without_replica_does_not_use_primary() {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_replica_connection();
assert!(matches!(result, Err(FalkorDBError::UnavailableProvider)));
}
#[test]
#[cfg(feature = "tokio")]
fn test_get_async_replica_connection_on_redis_provider_without_replica_does_not_use_primary() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_async_replica_connection().await;
assert!(matches!(result, Err(FalkorDBError::UnavailableProvider)));
});
}
#[test]
#[cfg(feature = "tokio")]
fn test_get_async_connection_manager_none_provider_is_unavailable() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let mut provider = FalkorClientProvider::None;
let result = provider.get_async_connection_manager(None).await;
assert!(matches!(result, Err(FalkorDBError::UnavailableProvider)));
});
}
#[test]
#[cfg(feature = "tokio")]
fn test_get_async_connection_manager_routes_through_sentinel() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:1").unwrap();
let sentinel = redis::sentinel::SentinelClient::build(
vec![connection_info],
"mymaster".to_string(),
None,
redis::sentinel::SentinelServerType::Master,
)
.unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: Some(sentinel),
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_async_connection_manager(None).await;
assert!(
matches!(result, Err(FalkorDBError::RedisError(_))),
"error should come from the sentinel resolution path"
);
});
}
#[test]
#[cfg(feature = "tokio")]
fn test_get_async_replica_connection_manager_errors_when_replica_unreachable() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let connection_info = redis::ConnectionInfo::from_str("redis://127.0.0.1:1").unwrap();
let replica = redis::sentinel::SentinelClient::build(
vec![connection_info],
"mymaster".to_string(),
None,
redis::sentinel::SentinelServerType::Replica,
)
.unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: Some(replica),
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_async_replica_connection_manager(None).await;
assert!(
matches!(result, Err(FalkorDBError::RedisError(_))),
"error should come from the replica manager path"
);
});
}
#[test]
#[cfg(feature = "tokio")]
fn test_get_async_replica_connection_manager_without_replica_does_not_use_primary() {
use tokio::runtime::Runtime;
let rt = Runtime::new().unwrap();
rt.block_on(async {
let client = redis::Client::open("redis://127.0.0.1:6379").unwrap();
let mut provider = FalkorClientProvider::Redis {
client,
sentinel: None,
sentinel_replica: None,
#[cfg(feature = "embedded-core")]
embedded_server: None,
response_timeout: None,
};
let result = provider.get_async_replica_connection_manager(None).await;
assert!(matches!(result, Err(FalkorDBError::UnavailableProvider)));
});
}
}