diff --git a/launchdarkly-server-sdk/src/client.rs b/launchdarkly-server-sdk/src/client.rs index ade25ba..21d8241 100644 --- a/launchdarkly-server-sdk/src/client.rs +++ b/launchdarkly-server-sdk/src/client.rs @@ -14,6 +14,7 @@ use tokio::sync::{broadcast, Semaphore}; use super::config::Config; use super::data_source_builders::BuildError as DataSourceError; use super::data_system::{DataSystem, FDv1DataSystem}; +use super::data_system_builders::BuildError as DataSystemError; use super::evaluation::{FlagDetail, FlagDetailConfig}; use super::stores::store::DataStore; use super::stores::store_builders::BuildError as DataStoreError; @@ -65,6 +66,12 @@ impl From for BuildError { } } +impl From for BuildError { + fn from(error: DataSystemError) -> Self { + Self::InvalidConfig(error.to_string()) + } +} + impl From for BuildError { fn from(error: DataStoreError) -> Self { Self::InvalidConfig(error.to_string()) @@ -184,13 +191,24 @@ impl Client { let event_processor = event_processor_builder.build(&endpoints, config.sdk_key(), tags.clone())?; - let mut data_source_builder = config.data_source_builder().to_owned(); - data_source_builder.set_instance_id(instance_id); - let data_source = data_source_builder.build(&endpoints, config.sdk_key(), tags.clone())?; - let data_system: Arc = Arc::new(FDv1DataSystem::new( - data_source, - config.data_store_builder(), - )?); + let data_system: Arc = match config.data_system_builder() { + Some(data_system_builder) => data_system_builder.build( + &endpoints, + config.sdk_key(), + tags.as_deref(), + &instance_id, + )?, + None => { + let mut data_source_builder = config.data_source_builder().to_owned(); + data_source_builder.set_instance_id(instance_id); + let data_source = + data_source_builder.build(&endpoints, config.sdk_key(), tags.clone())?; + Arc::new(FDv1DataSystem::new( + data_source, + config.data_store_builder(), + )?) + } + }; let data_store = data_system.store(); let events_default = EventsScope { diff --git a/launchdarkly-server-sdk/src/config.rs b/launchdarkly-server-sdk/src/config.rs index 32413b7..2b27c61 100644 --- a/launchdarkly-server-sdk/src/config.rs +++ b/launchdarkly-server-sdk/src/config.rs @@ -1,6 +1,7 @@ use thiserror::Error; use crate::data_source_builders::{DataSourceFactory, NullDataSourceBuilder}; +use crate::data_system_builders::{DataSystemBuilder, DataSystemFactory}; #[cfg(any( feature = "hyper-rustls-native-roots", @@ -136,6 +137,7 @@ pub struct Config { service_endpoints_builder: ServiceEndpointsBuilder, data_store_builder: Box, data_source_builder: Box, + data_system_builder: Option>, event_processor_builder: Box, application_tag: Option, instance_id: String, @@ -164,6 +166,11 @@ impl Config { self.data_source_builder.borrow() } + /// Returns the DataSystemFactory, if an FDv2 data system was configured. + pub(crate) fn data_system_builder(&self) -> Option<&dyn DataSystemFactory> { + self.data_system_builder.as_deref() + } + /// Returns the EventProcessorFactory pub fn event_processor_builder(&self) -> &dyn EventProcessorFactory { self.event_processor_builder.borrow() @@ -212,6 +219,7 @@ pub struct ConfigBuilder { service_endpoints_builder: Option, data_store_builder: Option>, data_source_builder: Option>, + data_system_builder: Option>, event_processor_builder: Option>, application_info: Option, offline: bool, @@ -226,6 +234,7 @@ impl ConfigBuilder { service_endpoints_builder: None, data_store_builder: None, data_source_builder: None, + data_system_builder: None, event_processor_builder: None, offline: false, daemon_mode: false, @@ -258,6 +267,16 @@ impl ConfigBuilder { self } + /// Set the data system to use for this client. + /// + /// When set, the data system supersedes the [data_source](ConfigBuilder::data_source). + /// If offline mode is enabled, it will be ignored. + pub fn data_system(mut self, builder: &DataSystemBuilder) -> Self { + let factory: Box = Box::new(builder.clone()); + self.data_system_builder = Some(factory); + self + } + /// Set the event processor to use for this client. /// For usage see [EventProcessorBuilder](crate::EventProcessorBuilder). /// @@ -307,8 +326,27 @@ impl ConfigBuilder { Some(_data_store_builder) => self.data_store_builder.unwrap(), }; + // The data system is optional; when set it supersedes the data source. + // Like the data source, it is ignored in offline or daemon mode. + let data_system_builder = match self.data_system_builder { + Some(_) if self.offline => { + warn!("Custom data system builders will be ignored when in offline mode"); + None + } + Some(_) if self.daemon_mode => { + warn!("Custom data system builders will be ignored when in daemon mode"); + None + } + other => other, + }; + let data_source_builder_result: Result, BuildError> = match self.data_source_builder { + None if data_system_builder.is_some() => Ok(Box::new(NullDataSourceBuilder::new())), + Some(_) if data_system_builder.is_some() => { + warn!("Custom data source builders will be ignored when a data system is configured"); + Ok(Box::new(NullDataSourceBuilder::new())) + } None if self.offline => Ok(Box::new(NullDataSourceBuilder::new())), Some(_) if self.offline => { warn!("Custom data source builders will be ignored when in offline mode"); @@ -401,6 +439,7 @@ impl ConfigBuilder { service_endpoints_builder, data_store_builder, data_source_builder, + data_system_builder, event_processor_builder, application_tag, instance_id, diff --git a/launchdarkly-server-sdk/src/data_system_builders.rs b/launchdarkly-server-sdk/src/data_system_builders.rs new file mode 100644 index 0000000..a036474 --- /dev/null +++ b/launchdarkly-server-sdk/src/data_system_builders.rs @@ -0,0 +1,520 @@ +use std::sync::Arc; +use std::time::Duration; + +use launchdarkly_sdk_transport::{HttpTransport, HyperTransport}; +use thiserror::Error; + +use crate::data_source_builders::{DataSourceFactory, StreamingDataSourceBuilder}; +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); + +/// Error returned when a data system configuration cannot be built. +#[non_exhaustive] +#[derive(Debug, Error)] +pub enum BuildError { + /// The data system configuration was invalid. + #[error("data system config failed to build: {0}")] + InvalidConfig(String), +} + +/// A configured FDv2 source usable as a synchronizer. +pub(crate) trait FDv2SynchronizerConfig { + fn build_synchronizer( + &self, + endpoints: &ServiceEndpoints, + headers: &RequestHeaders, + ) -> Result, BuildError>; + + fn to_owned(&self) -> Box; +} + +/// A configured FDv2 source usable as an initializer. +pub(crate) trait FDv2InitializerConfig { + fn build_initializer( + &self, + endpoints: &ServiceEndpoints, + headers: &RequestHeaders, + ) -> Result, BuildError>; + + fn to_owned(&self) -> Box; +} + +/// Builds the default HTTPS transport, or errors if no TLS feature is enabled. +fn default_https_transport() -> Result { + #[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::(BuildError::InvalidConfig( + "https connector required when hyper-rustls-native-roots, hyper-rustls-webpki-roots, or native-tls features are disabled".into(), + )) + } +} + +/// Configures an FDv2 streaming source, which can only act as a synchronizer. +#[derive(Clone)] +pub struct FDv2StreamingBuilder { + initial_reconnect_delay: Duration, + base_url: Option, + transport: Option, +} + +impl FDv2StreamingBuilder { + /// Creates a builder with default values. + pub fn new() -> Self { + Self { + initial_reconnect_delay: DEFAULT_INITIAL_RECONNECT_DELAY, + base_url: None, + transport: None, + } + } + + /// Sets the initial reconnect delay for the streaming connection. + pub fn initial_reconnect_delay(&mut self, duration: Duration) -> &mut Self { + self.initial_reconnect_delay = duration; + self + } + + /// Sets the streaming base URL, overriding the configured service endpoints. + pub fn base_url(&mut self, url: &str) -> &mut Self { + self.base_url = Some(url.to_string()); + self + } + + /// Sets the transport to use, instead of the default HTTPS transport. + pub fn transport(&mut self, transport: T) -> &mut Self { + self.transport = Some(transport); + self + } +} + +impl FDv2SynchronizerConfig + for FDv2StreamingBuilder +{ + fn build_synchronizer( + &self, + endpoints: &ServiceEndpoints, + headers: &RequestHeaders, + ) -> Result, BuildError> { + let base_url = self + .base_url + .clone() + .unwrap_or_else(|| endpoints.streaming_base_url().to_string()); + let factory: Box = match &self.transport { + Some(transport) => Box::new(StreamingSynchronizerFactory::new( + transport.clone(), + base_url, + headers.clone(), + self.initial_reconnect_delay, + )), + None => Box::new(StreamingSynchronizerFactory::new( + default_https_transport()?, + base_url, + headers.clone(), + self.initial_reconnect_delay, + )), + }; + Ok(factory) + } + + fn to_owned(&self) -> Box { + Box::new(self.clone()) + } +} + +impl Default for FDv2StreamingBuilder { + fn default() -> Self { + Self::new() + } +} + +/// Configures an FDv2 polling source, which can act as an initializer or a synchronizer. +#[derive(Clone)] +pub struct FDv2PollingBuilder { + poll_interval: Duration, + base_url: Option, + transport: Option, +} + +impl FDv2PollingBuilder { + /// Creates a builder with default values. + pub fn new() -> Self { + Self { + poll_interval: DEFAULT_POLL_INTERVAL, + base_url: None, + transport: None, + } + } + + /// Sets the interval between polling requests, with an effective minimum of 30 seconds. + pub fn poll_interval(&mut self, poll_interval: Duration) -> &mut Self { + self.poll_interval = poll_interval; + self + } + + /// Sets the polling base URL, overriding the configured service endpoints. + pub fn base_url(&mut self, url: &str) -> &mut Self { + self.base_url = Some(url.to_string()); + self + } + + /// Sets the transport to use, instead of the default HTTPS transport. + pub fn transport(&mut self, transport: T) -> &mut Self { + self.transport = Some(transport); + self + } +} + +impl FDv2SynchronizerConfig + for FDv2PollingBuilder +{ + fn build_synchronizer( + &self, + endpoints: &ServiceEndpoints, + headers: &RequestHeaders, + ) -> Result, BuildError> { + let base_url = self + .base_url + .clone() + .unwrap_or_else(|| endpoints.polling_base_url().to_string()); + let factory: Box = match &self.transport { + Some(transport) => Box::new(PollingSynchronizerFactory::new( + transport.clone(), + base_url, + headers.clone(), + self.poll_interval, + )), + None => Box::new(PollingSynchronizerFactory::new( + default_https_transport()?, + base_url, + headers.clone(), + self.poll_interval, + )), + }; + Ok(factory) + } + + fn to_owned(&self) -> Box { + Box::new(self.clone()) + } +} + +impl FDv2InitializerConfig + for FDv2PollingBuilder +{ + fn build_initializer( + &self, + endpoints: &ServiceEndpoints, + headers: &RequestHeaders, + ) -> Result, BuildError> { + let base_url = self + .base_url + .clone() + .unwrap_or_else(|| endpoints.polling_base_url().to_string()); + let factory: Box = match &self.transport { + Some(transport) => Box::new(PollingInitializerFactory::new( + transport.clone(), + base_url, + headers.clone(), + )), + None => Box::new(PollingInitializerFactory::new( + default_https_transport()?, + base_url, + headers.clone(), + )), + }; + Ok(factory) + } + + fn to_owned(&self) -> Box { + Box::new(self.clone()) + } +} + +impl Default for FDv2PollingBuilder { + fn default() -> Self { + Self::new() + } +} + +/// Configures the FDv2 data system. +pub struct DataSystemBuilder { + initializers: Vec>, + synchronizers: Vec>, + fdv1_fallback: Option>, + fallback_timeout: Duration, + recovery_timeout: Duration, +} + +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()), + fallback_timeout: self.fallback_timeout, + recovery_timeout: self.recovery_timeout, + } + } +} + +impl DataSystemBuilder { + /// Creates an empty builder; the caller adds sources explicitly. + pub fn custom() -> Self { + Self { + initializers: Vec::new(), + synchronizers: Vec::new(), + fdv1_fallback: None, + fallback_timeout: DEFAULT_FALLBACK_TIMEOUT, + recovery_timeout: DEFAULT_RECOVERY_TIMEOUT, + } + } + + /// Appends a polling initializer. + pub fn initializer( + &mut self, + source: FDv2PollingBuilder, + ) -> &mut Self { + self.initializers.push(Box::new(source)); + self + } + + /// Appends a streaming synchronizer, ordered after any already added. + pub fn streaming_synchronizer( + &mut self, + source: FDv2StreamingBuilder, + ) -> &mut Self { + self.synchronizers.push(Box::new(source)); + self + } + + /// Appends a polling synchronizer, ordered after any already added. + pub fn polling_synchronizer( + &mut self, + source: FDv2PollingBuilder, + ) -> &mut Self { + self.synchronizers.push(Box::new(source)); + self + } + + /// Sets the FDv1 source used as a last-resort fallback. + pub fn fdv1_fallback(&mut self, factory: &dyn DataSourceFactory) -> &mut Self { + self.fdv1_fallback = Some(factory.to_owned()); + self + } + + /// Disables the FDv1 fallback. + pub fn disable_fdv1_fallback(&mut self) -> &mut Self { + self.fdv1_fallback = None; + self + } + + /// Sets how long the active synchronizer may stay interrupted before failing over. + pub fn fallback_timeout(&mut self, timeout: Duration) -> &mut Self { + self.fallback_timeout = timeout; + self + } + + /// Sets how long a fallback synchronizer must run before recovering to the primary. + pub fn recovery_timeout(&mut self, timeout: Duration) -> &mut Self { + self.recovery_timeout = timeout; + self + } +} + +impl Default for DataSystemBuilder { + /// The recommended data system setup. + fn default() -> Self { + let mut builder = Self::custom(); + builder.initializer(FDv2PollingBuilder::::new()); + builder.streaming_synchronizer(FDv2StreamingBuilder::::new()); + builder.polling_synchronizer(FDv2PollingBuilder::::new()); + builder.fdv1_fallback(&StreamingDataSourceBuilder::::new()); + builder + } +} + +/// Builds the internal FDv2 data system from a configured source set. +pub(crate) trait DataSystemFactory { + fn build( + &self, + endpoints: &ServiceEndpoints, + sdk_key: &str, + tags: Option<&str>, + instance_id: &str, + ) -> Result, BuildError>; +} + +impl DataSystemFactory for DataSystemBuilder { + fn build( + &self, + endpoints: &ServiceEndpoints, + sdk_key: &str, + tags: Option<&str>, + instance_id: &str, + ) -> Result, BuildError> { + let headers = RequestHeaders::new(sdk_key, tags, instance_id); + + let initializer_factories: Vec> = self + .initializers + .iter() + .map(|c| c.build_initializer(endpoints, &headers)) + .collect::>()?; + + let mut synchronizer_factories: Vec> = self + .synchronizers + .iter() + .map(|c| c.build_synchronizer(endpoints, &headers).map(Arc::from)) + .collect::>()?; + + // Build the FDv1 fallback source once and wrap it as a synchronizer; the + // adapter re-subscribes it whenever the fallback activates. + 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 = Arc::new(FDv2DataSystem::new( + initializer_factories, + synchronizer_factories, + self.fallback_timeout, + self.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()); + assert_eq!(builder.fallback_timeout, DEFAULT_FALLBACK_TIMEOUT); + assert_eq!(builder.recovery_timeout, DEFAULT_RECOVERY_TIMEOUT); + } + + #[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()); + } + + #[test] + fn timeouts_are_configurable() { + let mut builder = DataSystemBuilder::custom(); + + builder.fallback_timeout(Duration::from_secs(7)); + builder.recovery_timeout(Duration::from_secs(11)); + + assert_eq!(builder.fallback_timeout, Duration::from_secs(7)); + assert_eq!(builder.recovery_timeout, Duration::from_secs(11)); + } + + #[derive(Debug, Clone)] + struct TestTransport; + + impl HttpTransport for TestTransport { + fn request(&self, _request: Request>) -> 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"); + + // Each source builds a factory from a configured transport. + assert!(FDv2StreamingBuilder::::new() + .transport(TestTransport) + .build_synchronizer(&endpoints, &headers) + .is_ok()); + assert!(FDv2PollingBuilder::::new() + .transport(TestTransport) + .build_synchronizer(&endpoints, &headers) + .is_ok()); + assert!(FDv2PollingBuilder::::new() + .transport(TestTransport) + .build_initializer(&endpoints, &headers) + .is_ok()); + } + + // The default path builds a real HTTPS transport, which needs a TLS backend, + // so this only runs where one of those features is enabled. + #[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"); + + // Each source builds a factory off the default HTTPS transport. + assert!(FDv2StreamingBuilder::::new() + .build_synchronizer(&endpoints, &headers) + .is_ok()); + assert!(FDv2PollingBuilder::::new() + .build_synchronizer(&endpoints, &headers) + .is_ok()); + assert!(FDv2PollingBuilder::::new() + .build_initializer(&endpoints, &headers) + .is_ok()); + } +} diff --git a/launchdarkly-server-sdk/src/fdv2/mod.rs b/launchdarkly-server-sdk/src/fdv2/mod.rs index f2b2ed3..cf1880b 100644 --- a/launchdarkly-server-sdk/src/fdv2/mod.rs +++ b/launchdarkly-server-sdk/src/fdv2/mod.rs @@ -1,10 +1,10 @@ -mod data_system; -mod fdv1_adapter; +pub mod data_system; +pub mod fdv1_adapter; pub mod model; -mod polling; +pub mod polling; mod protocol; -mod request_headers; +pub mod request_headers; mod source; -mod streaming; +pub mod streaming; mod url; mod wire; diff --git a/launchdarkly-server-sdk/src/fdv2/polling.rs b/launchdarkly-server-sdk/src/fdv2/polling.rs index 7a8e036..028c5d5 100644 --- a/launchdarkly-server-sdk/src/fdv2/polling.rs +++ b/launchdarkly-server-sdk/src/fdv2/polling.rs @@ -6,12 +6,13 @@ use http::{Request, Response, StatusCode}; use launchdarkly_sdk_transport::{ByteStream, HttpTransport}; use serde::Deserialize; +use super::data_system::{InitializerFactory, SynchronizerFactory}; use super::model::{ChangeSetKind, Selector}; use super::protocol::{FDv2ProtocolHandler, ProtocolError, ProtocolResult}; use super::request_headers::RequestHeaders; use super::source::{ read_fallback_directive, ErrorInfo, ErrorKind, FDv1FallbackDirective, FDv2SourceEvent, - FDv2SourceResult, + FDv2SourceResult, Initializer, Synchronizer, }; use super::url::build_fdv2_url; use crate::reqwest::is_http_error_recoverable; @@ -290,6 +291,73 @@ impl super::source::Synchronizer for PollingSynchronizer { } } +/// Builds a fresh `PollingInitializer` each time an initializer run starts. +pub(crate) struct PollingInitializerFactory { + transport: T, + base_url: String, + headers: RequestHeaders, +} + +impl PollingInitializerFactory { + pub(crate) fn new(transport: T, base_url: String, headers: RequestHeaders) -> Self { + Self { + transport, + base_url, + headers, + } + } +} + +impl InitializerFactory + for PollingInitializerFactory +{ + fn create(&self) -> Box { + Box::new(PollingInitializer::new( + self.transport.clone(), + self.base_url.clone(), + self.headers.clone(), + None, + )) + } +} + +/// Builds a fresh `PollingSynchronizer` each time the synchronizer activates. +pub(crate) struct PollingSynchronizerFactory { + transport: T, + base_url: String, + headers: RequestHeaders, + poll_interval: Duration, +} + +impl PollingSynchronizerFactory { + pub(crate) fn new( + transport: T, + base_url: String, + headers: RequestHeaders, + poll_interval: Duration, + ) -> Self { + Self { + transport, + base_url, + headers, + poll_interval, + } + } +} + +impl SynchronizerFactory + for PollingSynchronizerFactory +{ + fn create(&self) -> Box { + Box::new(PollingSynchronizer::new( + self.transport.clone(), + self.base_url.clone(), + self.headers.clone(), + self.poll_interval, + )) + } +} + fn build_poll_request( base_url: &str, headers: &RequestHeaders, diff --git a/launchdarkly-server-sdk/src/fdv2/streaming.rs b/launchdarkly-server-sdk/src/fdv2/streaming.rs index bf8414c..da7a330 100644 --- a/launchdarkly-server-sdk/src/fdv2/streaming.rs +++ b/launchdarkly-server-sdk/src/fdv2/streaming.rs @@ -10,6 +10,7 @@ use http::Uri; use launchdarkly_sdk_transport::HttpTransport; use tokio::sync::watch; +use super::data_system::SynchronizerFactory; use super::model::Selector; use super::protocol::{FDv2ProtocolHandler, ProtocolError, ProtocolResult}; use super::request_headers::RequestHeaders; @@ -293,6 +294,43 @@ impl Synchronizer for Streamin } } +/// Builds a fresh `StreamingSynchronizer` each time the synchronizer activates. +pub(crate) struct StreamingSynchronizerFactory { + transport: T, + base_url: String, + headers: RequestHeaders, + initial_reconnect_delay: Duration, +} + +impl StreamingSynchronizerFactory { + pub(crate) fn new( + transport: T, + base_url: String, + headers: RequestHeaders, + initial_reconnect_delay: Duration, + ) -> Self { + Self { + transport, + base_url, + headers, + initial_reconnect_delay, + } + } +} + +impl SynchronizerFactory + for StreamingSynchronizerFactory +{ + fn create(&self) -> Box { + Box::new(StreamingSynchronizer::new( + self.transport.clone(), + self.base_url.clone(), + self.headers.clone(), + self.initial_reconnect_delay, + )) + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/launchdarkly-server-sdk/src/fdv2/wire.rs b/launchdarkly-server-sdk/src/fdv2/wire.rs index 93ebc67..a8078ac 100644 --- a/launchdarkly-server-sdk/src/fdv2/wire.rs +++ b/launchdarkly-server-sdk/src/fdv2/wire.rs @@ -18,14 +18,19 @@ pub(super) struct ServerIntent { #[derive(Debug, Deserialize)] #[serde(rename_all = "camelCase")] pub(super) struct ServerIntentPayload { + // Parsed from the wire but not used; only intent_code is read. + #[allow(dead_code)] pub(super) id: String, + #[allow(dead_code)] pub(super) target: u64, pub(super) intent_code: IntentCode, + #[allow(dead_code)] pub(super) reason: String, } #[derive(Debug, Deserialize)] pub(super) struct PutObject { + #[allow(dead_code)] // Parsed from the wire but not used. pub(super) version: u64, pub(super) kind: String, pub(super) key: String, diff --git a/launchdarkly-server-sdk/src/lib.rs b/launchdarkly-server-sdk/src/lib.rs index 3dd353b..3a6358f 100644 --- a/launchdarkly-server-sdk/src/lib.rs +++ b/launchdarkly-server-sdk/src/lib.rs @@ -32,6 +32,9 @@ pub use config::{ApplicationInfo, BuildError as ConfigBuildError, Config, Config pub use data_source_builders::{ BuildError as DataSourceBuildError, PollingDataSourceBuilder, StreamingDataSourceBuilder, }; +pub use data_system_builders::{ + BuildError as DataSystemBuildError, DataSystemBuilder, FDv2PollingBuilder, FDv2StreamingBuilder, +}; pub use evaluation::{FlagDetail, FlagDetailConfig, FlagFilter}; pub use events::event::MigrationOpEvent; pub use events::processor::EventProcessor; @@ -59,9 +62,9 @@ mod config; mod data_source; mod data_source_builders; mod data_system; +mod data_system_builders; mod evaluation; mod events; -#[allow(dead_code)] // Tested but not yet reachable from production code. mod fdv2; mod feature_requester; mod feature_requester_builders;