use std::sync::Arc;
use std::time::Duration;
use launchdarkly_sdk_transport::{HttpTransport, HyperTransport};
use thiserror::Error;
use crate::data_source_builders::{DataSourceFactory, PollingDataSourceBuilder};
use crate::data_system::DataSystem;
use crate::fdv2::data_system::{FDv2DataSystem, InitializerFactory, SynchronizerFactory};
use crate::fdv2::fdv1_adapter::FDv1AdapterFactory;
use crate::fdv2::polling::{PollingInitializerFactory, PollingSynchronizerFactory};
use crate::fdv2::request_headers::RequestHeaders;
use crate::fdv2::streaming::StreamingSynchronizerFactory;
use crate::service_endpoints::ServiceEndpoints;
const DEFAULT_INITIAL_RECONNECT_DELAY: Duration = Duration::from_secs(1);
const DEFAULT_POLL_INTERVAL: Duration = Duration::from_secs(30);
const DEFAULT_FALLBACK_TIMEOUT: Duration = Duration::from_secs(120);
const DEFAULT_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
#[non_exhaustive]
#[derive(Debug, Error)]
pub enum BuildError {
#[error("data system config failed to build: {0}")]
InvalidConfig(String),
}
#[non_exhaustive]
pub struct DataSourceBuildContext<'a> {
pub endpoints: &'a ServiceEndpoints,
pub headers: &'a RequestHeaders,
}
pub trait FDv2SynchronizerConfig {
fn build_synchronizer(
&self,
context: &DataSourceBuildContext,
) -> Result<Box<dyn SynchronizerFactory>, BuildError>;
fn to_owned(&self) -> Box<dyn FDv2SynchronizerConfig>;
}
pub trait FDv2InitializerConfig {
fn build_initializer(
&self,
context: &DataSourceBuildContext,
) -> Result<Box<dyn InitializerFactory>, BuildError>;
fn to_owned(&self) -> Box<dyn FDv2InitializerConfig>;
}
fn default_https_transport() -> Result<impl HttpTransport + 'static, BuildError> {
#[cfg(any(
feature = "hyper-rustls-native-roots",
feature = "hyper-rustls-webpki-roots",
feature = "native-tls"
))]
{
HyperTransport::new_https().map_err(|e| {
BuildError::InvalidConfig(format!("failed to create default https transport: {e:?}"))
})
}
#[cfg(not(any(
feature = "hyper-rustls-native-roots",
feature = "hyper-rustls-webpki-roots",
feature = "native-tls"
)))]
{
Err::<HyperTransport, _>(BuildError::InvalidConfig(
"https connector required when hyper-rustls-native-roots, hyper-rustls-webpki-roots, or native-tls features are disabled".into(),
))
}
}
#[derive(Clone)]
pub struct FDv2StreamingBuilder<T: HttpTransport = HyperTransport> {
initial_reconnect_delay: Duration,
base_url: Option<String>,
transport: Option<T>,
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> FDv2StreamingBuilder<T> {
pub fn new() -> Self {
Self {
initial_reconnect_delay: DEFAULT_INITIAL_RECONNECT_DELAY,
base_url: None,
transport: None,
}
}
pub fn initial_reconnect_delay(&mut self, duration: Duration) -> &mut Self {
self.initial_reconnect_delay = duration;
self
}
pub fn base_url(&mut self, url: &str) -> &mut Self {
self.base_url = Some(url.to_string());
self
}
pub fn transport(&mut self, transport: T) -> &mut Self {
self.transport = Some(transport);
self
}
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> FDv2SynchronizerConfig
for FDv2StreamingBuilder<T>
{
fn build_synchronizer(
&self,
context: &DataSourceBuildContext,
) -> Result<Box<dyn SynchronizerFactory>, BuildError> {
let base_url = self
.base_url
.clone()
.unwrap_or_else(|| context.endpoints.streaming_base_url().to_string());
let factory: Box<dyn SynchronizerFactory> = match &self.transport {
Some(transport) => Box::new(StreamingSynchronizerFactory::new(
transport.clone(),
base_url,
context.headers.clone(),
self.initial_reconnect_delay,
)),
None => Box::new(StreamingSynchronizerFactory::new(
default_https_transport()?,
base_url,
context.headers.clone(),
self.initial_reconnect_delay,
)),
};
Ok(factory)
}
fn to_owned(&self) -> Box<dyn FDv2SynchronizerConfig> {
Box::new(self.clone())
}
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> Default for FDv2StreamingBuilder<T> {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct FDv2PollingBuilder<T: HttpTransport = HyperTransport> {
poll_interval: Duration,
base_url: Option<String>,
transport: Option<T>,
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> FDv2PollingBuilder<T> {
pub fn new() -> Self {
Self {
poll_interval: DEFAULT_POLL_INTERVAL,
base_url: None,
transport: None,
}
}
pub fn poll_interval(&mut self, poll_interval: Duration) -> &mut Self {
self.poll_interval = poll_interval;
self
}
pub fn base_url(&mut self, url: &str) -> &mut Self {
self.base_url = Some(url.to_string());
self
}
pub fn transport(&mut self, transport: T) -> &mut Self {
self.transport = Some(transport);
self
}
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> FDv2SynchronizerConfig
for FDv2PollingBuilder<T>
{
fn build_synchronizer(
&self,
context: &DataSourceBuildContext,
) -> Result<Box<dyn SynchronizerFactory>, BuildError> {
let base_url = self
.base_url
.clone()
.unwrap_or_else(|| context.endpoints.polling_base_url().to_string());
let factory: Box<dyn SynchronizerFactory> = match &self.transport {
Some(transport) => Box::new(PollingSynchronizerFactory::new(
transport.clone(),
base_url,
context.headers.clone(),
self.poll_interval,
)),
None => Box::new(PollingSynchronizerFactory::new(
default_https_transport()?,
base_url,
context.headers.clone(),
self.poll_interval,
)),
};
Ok(factory)
}
fn to_owned(&self) -> Box<dyn FDv2SynchronizerConfig> {
Box::new(self.clone())
}
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> FDv2InitializerConfig
for FDv2PollingBuilder<T>
{
fn build_initializer(
&self,
context: &DataSourceBuildContext,
) -> Result<Box<dyn InitializerFactory>, BuildError> {
let base_url = self
.base_url
.clone()
.unwrap_or_else(|| context.endpoints.polling_base_url().to_string());
let factory: Box<dyn InitializerFactory> = match &self.transport {
Some(transport) => Box::new(PollingInitializerFactory::new(
transport.clone(),
base_url,
context.headers.clone(),
)),
None => Box::new(PollingInitializerFactory::new(
default_https_transport()?,
base_url,
context.headers.clone(),
)),
};
Ok(factory)
}
fn to_owned(&self) -> Box<dyn FDv2InitializerConfig> {
Box::new(self.clone())
}
}
impl<T: HttpTransport + Clone + Send + Sync + 'static> Default for FDv2PollingBuilder<T> {
fn default() -> Self {
Self::new()
}
}
pub struct DataSystemBuilder {
initializers: Vec<Box<dyn FDv2InitializerConfig>>,
synchronizers: Vec<Box<dyn FDv2SynchronizerConfig>>,
fdv1_fallback: Option<Box<dyn DataSourceFactory>>,
}
impl Clone for DataSystemBuilder {
fn clone(&self) -> Self {
Self {
initializers: self.initializers.iter().map(|c| (**c).to_owned()).collect(),
synchronizers: self
.synchronizers
.iter()
.map(|c| (**c).to_owned())
.collect(),
fdv1_fallback: self.fdv1_fallback.as_ref().map(|f| (**f).to_owned()),
}
}
}
impl DataSystemBuilder {
pub fn custom() -> Self {
Self {
initializers: Vec::new(),
synchronizers: Vec::new(),
fdv1_fallback: None,
}
}
pub fn initializer(&mut self, source: impl FDv2InitializerConfig + 'static) -> &mut Self {
self.initializers.push(Box::new(source));
self
}
pub fn synchronizer(&mut self, source: impl FDv2SynchronizerConfig + 'static) -> &mut Self {
self.synchronizers.push(Box::new(source));
self
}
pub fn fdv1_fallback(&mut self, factory: &dyn DataSourceFactory) -> &mut Self {
self.fdv1_fallback = Some(factory.to_owned());
self
}
pub fn disable_fdv1_fallback(&mut self) -> &mut Self {
self.fdv1_fallback = None;
self
}
}
impl Default for DataSystemBuilder {
fn default() -> Self {
let mut builder = Self::custom();
builder.initializer(FDv2PollingBuilder::<HyperTransport>::new());
builder.synchronizer(FDv2StreamingBuilder::<HyperTransport>::new());
builder.synchronizer(FDv2PollingBuilder::<HyperTransport>::new());
builder.fdv1_fallback(&PollingDataSourceBuilder::<HyperTransport>::new());
builder
}
}
pub(crate) trait DataSystemFactory {
fn build(
&self,
endpoints: &ServiceEndpoints,
sdk_key: &str,
tags: Option<&str>,
instance_id: &str,
) -> Result<Arc<dyn DataSystem>, BuildError>;
}
impl DataSystemFactory for DataSystemBuilder {
fn build(
&self,
endpoints: &ServiceEndpoints,
sdk_key: &str,
tags: Option<&str>,
instance_id: &str,
) -> Result<Arc<dyn DataSystem>, BuildError> {
let headers = RequestHeaders::new(sdk_key, tags, instance_id);
let context = DataSourceBuildContext {
endpoints,
headers: &headers,
};
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = self
.initializers
.iter()
.map(|c| c.build_initializer(&context).map(Arc::from))
.collect::<Result<_, _>>()?;
let mut synchronizer_factories: Vec<Arc<dyn SynchronizerFactory>> = self
.synchronizers
.iter()
.map(|c| c.build_synchronizer(&context).map(Arc::from))
.collect::<Result<_, _>>()?;
if let Some(fdv1_factory) = &self.fdv1_fallback {
let mut fdv1_factory = (**fdv1_factory).to_owned();
fdv1_factory.set_instance_id(instance_id.to_string());
let source = fdv1_factory
.build(endpoints, sdk_key, tags.map(|t| t.to_string()))
.map_err(|e| {
BuildError::InvalidConfig(format!("failed to build FDv1 fallback source: {e}"))
})?;
let adapter = FDv1AdapterFactory::new(Box::new(move || source.clone()));
synchronizer_factories.push(Arc::new(adapter));
}
let system: Arc<dyn DataSystem> = Arc::new(FDv2DataSystem::new(
initializer_factories,
synchronizer_factories,
DEFAULT_FALLBACK_TIMEOUT,
DEFAULT_RECOVERY_TIMEOUT,
));
Ok(system)
}
}
#[cfg(test)]
mod tests {
use bytes::Bytes;
use launchdarkly_sdk_transport::{Request, ResponseFuture};
use super::*;
#[test]
fn custom_starts_empty() {
let builder = DataSystemBuilder::custom();
assert!(builder.initializers.is_empty());
assert!(builder.synchronizers.is_empty());
assert!(builder.fdv1_fallback.is_none());
}
#[test]
fn default_has_recommended_sources() {
let builder = DataSystemBuilder::default();
assert_eq!(builder.initializers.len(), 1);
assert_eq!(builder.synchronizers.len(), 2);
assert!(builder.fdv1_fallback.is_some());
}
#[test]
fn disable_fdv1_fallback_clears_it() {
let mut builder = DataSystemBuilder::default();
assert!(builder.fdv1_fallback.is_some());
builder.disable_fdv1_fallback();
assert!(builder.fdv1_fallback.is_none());
}
#[derive(Debug, Clone)]
struct TestTransport;
impl HttpTransport for TestTransport {
fn request(&self, _request: Request<Option<Bytes>>) -> ResponseFuture {
unreachable!();
}
}
#[test]
fn builders_build_factories_with_injected_transport() {
let endpoints = crate::ServiceEndpointsBuilder::new().build().unwrap();
let headers = RequestHeaders::new("sdk-key", None, "test-instance");
let context = DataSourceBuildContext {
endpoints: &endpoints,
headers: &headers,
};
assert!(FDv2StreamingBuilder::<TestTransport>::new()
.transport(TestTransport)
.build_synchronizer(&context)
.is_ok());
assert!(FDv2PollingBuilder::<TestTransport>::new()
.transport(TestTransport)
.build_synchronizer(&context)
.is_ok());
assert!(FDv2PollingBuilder::<TestTransport>::new()
.transport(TestTransport)
.build_initializer(&context)
.is_ok());
}
#[test]
#[cfg(any(
feature = "hyper-rustls-native-roots",
feature = "hyper-rustls-webpki-roots",
feature = "native-tls"
))]
fn builders_build_factories_with_default_transport() {
let endpoints = crate::ServiceEndpointsBuilder::new().build().unwrap();
let headers = RequestHeaders::new("sdk-key", None, "test-instance");
let context = DataSourceBuildContext {
endpoints: &endpoints,
headers: &headers,
};
assert!(FDv2StreamingBuilder::<HyperTransport>::new()
.build_synchronizer(&context)
.is_ok());
assert!(FDv2PollingBuilder::<HyperTransport>::new()
.build_synchronizer(&context)
.is_ok());
assert!(FDv2PollingBuilder::<HyperTransport>::new()
.build_initializer(&context)
.is_ok());
}
}