use crate::client::ClientWrapper;
use redis::PushInfo;
pub use redis::{
ErrorKind, PubSubChannelOrPattern, PubSubSubscriptionInfo, PubSubSubscriptionKind,
PubSubSynchronizer, RedisError,
};
use std::sync::{Arc, Weak};
use std::time::Duration;
use tokio::sync::{RwLock, mpsc};
#[cfg(feature = "mock-pubsub")]
mod mock;
#[cfg(feature = "mock-pubsub")]
pub use mock::MockPubSubBroker;
#[cfg(not(feature = "mock-pubsub"))]
pub mod synchronizer;
pub async fn create_pubsub_synchronizer(
_push_sender: Option<mpsc::UnboundedSender<PushInfo>>,
initial_subscriptions: Option<redis::PubSubSubscriptionInfo>,
is_cluster: bool,
internal_client: Weak<RwLock<ClientWrapper>>,
reconciliation_interval: Option<Duration>,
_request_timeout: Duration,
) -> Arc<dyn PubSubSynchronizer> {
#[cfg(feature = "mock-pubsub")]
{
let sync = mock::MockPubSubSynchronizer::create(
_push_sender,
initial_subscriptions,
is_cluster,
reconciliation_interval,
)
.await;
if internal_client.upgrade().is_some() {
sync.as_any()
.downcast_ref::<mock::MockPubSubSynchronizer>()
.expect("Expected MockPubSubSynchronizer")
.set_internal_client(internal_client);
}
sync
}
#[cfg(not(feature = "mock-pubsub"))]
{
let sync = synchronizer::GlidePubSubSynchronizer::new(
initial_subscriptions,
is_cluster,
reconciliation_interval,
_request_timeout,
);
if internal_client.upgrade().is_some() {
sync.as_any()
.downcast_ref::<synchronizer::GlidePubSubSynchronizer>()
.expect("Expected GlidePubSubSynchronizer")
.set_internal_client(internal_client);
}
sync
}
}