use redis::sentinel::{
SentinelClient, SentinelClientBuilder, SentinelNodeConnectionInfo, SentinelServerType,
};
use redis::{Connection, ControlFlow, IntoConnectionInfo, Msg, ToRedisArgs};
use std::fmt::Debug;
use super::retry::Retry;
#[derive(Debug)]
pub enum RedisPubSubSentinelError {
SentinelConfigAlreadyConsumed,
SentinelClientBuilderAlreadyConsumed,
FailToBuildSentinelClient,
FailToGetConnection,
FailToSubscribeToChannels,
FailToSubscribeToChannelsWithPatterns,
FailToGetMessage,
}
pub struct RedisPubSubSentinel<I, A>
where
I: IntoConnectionInfo + Sized + Debug + Clone,
A: ToRedisArgs,
{
config: Config<I>,
channels: Vec<A>,
patterns: Vec<A>,
retry: Retry,
}
enum Config<I>
where
I: IntoConnectionInfo + Sized + Debug + Clone,
{
Client(Option<Box<SentinelConfig<I>>>),
Builder(Option<Box<SentinelClientBuilder>>),
}
struct SentinelConfig<I> {
addrs: Vec<I>,
service_name: String,
node_conn_info: Option<SentinelNodeConnectionInfo>,
server_type: SentinelServerType,
}
impl<I, A> RedisPubSubSentinel<I, A>
where
I: IntoConnectionInfo + Sized + Debug + Clone,
A: ToRedisArgs,
{
pub fn new(addrs: Vec<I>, service_name: impl AsRef<str>) -> Self {
Self {
config: Config::Client(Some(Box::new(SentinelConfig {
addrs,
service_name: service_name.as_ref().to_string(),
node_conn_info: None,
server_type: SentinelServerType::Master,
}))),
channels: Vec::new(),
patterns: Vec::new(),
retry: Retry::new(),
}
}
pub fn with_client_params(
addrs: Vec<I>,
service_name: impl AsRef<str>,
node_conn_info: SentinelNodeConnectionInfo,
server_type: SentinelServerType,
) -> Self {
Self {
config: Config::Client(Some(Box::new(SentinelConfig {
addrs,
service_name: service_name.as_ref().to_string(),
node_conn_info: Some(node_conn_info),
server_type,
}))),
channels: Vec::new(),
patterns: Vec::new(),
retry: Retry::new(),
}
}
pub fn with_client_builder(client_builder: SentinelClientBuilder) -> Self {
Self {
config: Config::Builder(Some(Box::new(client_builder))),
channels: Vec::new(),
patterns: Vec::new(),
retry: Retry::new(),
}
}
pub fn set_retry(&mut self, max_count: u32, init_delay_ms: u64, max_delay_ms: u64) {
self.retry = Retry::with_params(max_count, init_delay_ms, max_delay_ms);
}
pub fn subscribe(&mut self, channels: A) {
self.channels.push(channels);
}
pub fn psubscribe(&mut self, patterns: A) {
self.patterns.push(patterns);
}
pub fn receive<F, U>(mut self, mut f: F) -> errs::Result<U>
where
F: FnMut(Msg) -> ControlFlow<U>,
{
let mut client = match &mut self.config {
Config::Client(ref mut c) => {
let cfg = c.take().ok_or_else(|| {
errs::Err::new(RedisPubSubSentinelError::SentinelConfigAlreadyConsumed)
})?;
SentinelClient::build(
cfg.addrs,
cfg.service_name,
cfg.node_conn_info,
cfg.server_type,
)
.map_err(|e| {
errs::Err::with_source(RedisPubSubSentinelError::FailToBuildSentinelClient, e)
})
}
Config::Builder(ref mut b) => {
let builder = b.take().ok_or_else(|| {
errs::Err::new(RedisPubSubSentinelError::SentinelClientBuilderAlreadyConsumed)
})?;
builder.build().map_err(|e| {
errs::Err::with_source(RedisPubSubSentinelError::FailToBuildSentinelClient, e)
})
}
}?;
loop {
let mut conn: Connection = match client.get_connection() {
Ok(c) => c,
Err(e) => {
if self.retry.wait_with_backoff() {
continue;
}
return Err(errs::Err::with_source(
RedisPubSubSentinelError::FailToGetConnection,
e,
));
}
};
let mut pubsub = conn.as_pubsub();
for c in self.channels.iter() {
pubsub.subscribe(c).map_err(|e| {
errs::Err::with_source(RedisPubSubSentinelError::FailToSubscribeToChannels, e)
})?;
}
for p in self.patterns.iter() {
pubsub.psubscribe(p).map_err(|e| {
errs::Err::with_source(
RedisPubSubSentinelError::FailToSubscribeToChannelsWithPatterns,
e,
)
})?;
}
loop {
match pubsub.get_message() {
Ok(msg) => {
self.retry.reset();
if let ControlFlow::Break(value) = f(msg) {
return Ok(value);
}
}
Err(e) => {
if self.retry.wait_with_backoff() {
continue;
}
return Err(errs::Err::with_source(
RedisPubSubSentinelError::FailToGetMessage,
e,
));
}
}
}
}
}
}
#[cfg(test)]
mod unit_tests {
use super::*;
use crate::pubsub::{RedisPubSubMsgDataConn, RedisPubSubMsgDataSrc};
use crate::sentinel_sync::{RedisSentinelDataConn, RedisSentinelDataSrc};
use override_macro::{overridable, override_with};
use redis::{ControlFlow, TypedCommands};
use sabi::{DataAcc, DataHub};
#[overridable]
trait PublishData {
fn say_greet(&mut self, s: &str) -> errs::Result<()>;
}
fn publish_logic(data: &mut impl PublishData) -> errs::Result<()> {
data.say_greet("Hello")?;
Ok(())
}
#[overridable]
trait SubscribeData {
fn receive_greet(&mut self) -> errs::Result<String>;
}
fn subscribe_logic(data: &mut impl SubscribeData) -> errs::Result<()> {
assert_eq!(data.receive_greet()?, "Hello");
Ok(())
}
#[overridable]
trait RedisPubSubDataAcc: DataAcc {
fn say_greet(&mut self, s: &str) -> errs::Result<()> {
let data_conn = self.get_data_conn::<RedisSentinelDataConn>("redis")?;
let conn = data_conn.get_connection();
std::thread::sleep(std::time::Duration::from_millis(100));
conn.publish("channel-2", s).unwrap();
Ok(())
}
fn receive_greet(&mut self) -> errs::Result<String> {
let data_conn = self.get_data_conn::<RedisPubSubMsgDataConn>("redis/pubsub")?;
let msg = data_conn.get_message();
let payload: String = msg.get_payload().unwrap();
Ok(payload)
}
}
impl RedisPubSubDataAcc for DataHub {}
#[override_with(RedisPubSubDataAcc)]
impl PublishData for DataHub {}
#[override_with(RedisPubSubDataAcc)]
impl SubscribeData for DataHub {}
#[test]
fn test() -> errs::Result<()> {
{
let _ = std::thread::spawn(move || {
let mut data = DataHub::new();
data.uses(
"redis",
RedisSentinelDataSrc::new(
vec![
"redis://127.0.0.1:26479",
"redis://127.0.0.1:26480",
"redis://127.0.0.1:26481",
],
"mymaster",
),
);
data.run(publish_logic).unwrap();
});
}
{
let mut pubsub = RedisPubSubSentinel::new(
vec![
"redis://127.0.0.1:26479",
"redis://127.0.0.1:26480",
"redis://127.0.0.1:26481",
],
"mymaster",
);
pubsub.subscribe("channel-2");
let n = pubsub.receive(|msg| {
let mut data = DataHub::new();
data.uses("redis/pubsub", RedisPubSubMsgDataSrc::new(msg));
data.run(subscribe_logic).unwrap();
ControlFlow::Break(1)
})?;
assert_eq!(n, 1);
}
Ok(())
}
}