From 93fe77f76c1ebaf3709e16f243d8e779bfcfaf50 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Thu, 6 Aug 2026 10:51:00 -0700 Subject: [PATCH 1/5] feat: Add the FDv2 data system configuration API --- launchdarkly-server-sdk/src/client.rs | 29 +- launchdarkly-server-sdk/src/config.rs | 34 ++ .../src/data_system_builders.rs | 479 ++++++++++++++++++ launchdarkly-server-sdk/src/fdv2/mod.rs | 8 +- launchdarkly-server-sdk/src/fdv2/polling.rs | 70 ++- launchdarkly-server-sdk/src/fdv2/streaming.rs | 38 ++ launchdarkly-server-sdk/src/lib.rs | 6 +- 7 files changed, 651 insertions(+), 13 deletions(-) create mode 100644 launchdarkly-server-sdk/src/data_system_builders.rs diff --git a/launchdarkly-server-sdk/src/client.rs b/launchdarkly-server-sdk/src/client.rs index ade25ba..c9dec96 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,21 @@ 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.clone(), &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..f63ce0b 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). /// @@ -348,6 +367,20 @@ impl ConfigBuilder { }; let data_source_builder = data_source_builder_result?; + // The data system is optional; when unset the client uses the data source above. + // Like that 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 event_processor_builder_result: Result, BuildError> = match self.event_processor_builder { None if self.offline => Ok(Box::new(NullEventProcessorBuilder::new())), @@ -401,6 +434,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..6066cef --- /dev/null +++ b/launchdarkly-server-sdk/src/data_system_builders.rs @@ -0,0 +1,479 @@ +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::source::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, + instance_id: &str, + ) -> Result, BuildError>; +} + +impl DataSystemFactory for DataSystemBuilder { + fn build( + &self, + endpoints: &ServiceEndpoints, + sdk_key: &str, + tags: Option, + instance_id: &str, + ) -> Result, BuildError> { + let headers = RequestHeaders::new(sdk_key, tags.as_deref(), 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_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 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)); + } + + #[test] + 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..b9e458a 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; 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/lib.rs b/launchdarkly-server-sdk/src/lib.rs index 3dd353b..a8cb869 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,10 @@ 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. +#[allow(dead_code)] // Some items are consumed only by later data-system phases. mod fdv2; mod feature_requester; mod feature_requester_builders; From 05a9c970a2dc1bd11171a55bc6e184486d06c1e2 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Fri, 7 Aug 2026 23:22:18 -0700 Subject: [PATCH 2/5] fix: Let a configured data system supersede the data source --- launchdarkly-server-sdk/src/config.rs | 33 +++++++++++++++------------ 1 file changed, 19 insertions(+), 14 deletions(-) diff --git a/launchdarkly-server-sdk/src/config.rs b/launchdarkly-server-sdk/src/config.rs index f63ce0b..2b27c61 100644 --- a/launchdarkly-server-sdk/src/config.rs +++ b/launchdarkly-server-sdk/src/config.rs @@ -326,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"); @@ -367,20 +386,6 @@ impl ConfigBuilder { }; let data_source_builder = data_source_builder_result?; - // The data system is optional; when unset the client uses the data source above. - // Like that 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 event_processor_builder_result: Result, BuildError> = match self.event_processor_builder { None if self.offline => Ok(Box::new(NullEventProcessorBuilder::new())), From b19625c55bec4f65eca42dfc6895705fc1c1be75 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Sat, 8 Aug 2026 22:24:10 -0700 Subject: [PATCH 3/5] refactor: Import RequestHeaders from its own module and borrow tags --- launchdarkly-server-sdk/src/client.rs | 9 ++++++--- .../src/data_system_builders.rs | 16 +++++++++------- launchdarkly-server-sdk/src/fdv2/mod.rs | 2 +- 3 files changed, 16 insertions(+), 11 deletions(-) diff --git a/launchdarkly-server-sdk/src/client.rs b/launchdarkly-server-sdk/src/client.rs index c9dec96..21d8241 100644 --- a/launchdarkly-server-sdk/src/client.rs +++ b/launchdarkly-server-sdk/src/client.rs @@ -192,9 +192,12 @@ impl Client { event_processor_builder.build(&endpoints, config.sdk_key(), tags.clone())?; let data_system: Arc = match config.data_system_builder() { - Some(data_system_builder) => { - data_system_builder.build(&endpoints, config.sdk_key(), tags.clone(), &instance_id)? - } + 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); diff --git a/launchdarkly-server-sdk/src/data_system_builders.rs b/launchdarkly-server-sdk/src/data_system_builders.rs index 6066cef..9490689 100644 --- a/launchdarkly-server-sdk/src/data_system_builders.rs +++ b/launchdarkly-server-sdk/src/data_system_builders.rs @@ -9,7 +9,7 @@ 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::source::RequestHeaders; +use crate::fdv2::request_headers::RequestHeaders; use crate::fdv2::streaming::StreamingSynchronizerFactory; use crate::service_endpoints::ServiceEndpoints; @@ -366,7 +366,7 @@ pub(crate) trait DataSystemFactory { &self, endpoints: &ServiceEndpoints, sdk_key: &str, - tags: Option, + tags: Option<&str>, instance_id: &str, ) -> Result, BuildError>; } @@ -376,10 +376,10 @@ impl DataSystemFactory for DataSystemBuilder { &self, endpoints: &ServiceEndpoints, sdk_key: &str, - tags: Option, + tags: Option<&str>, instance_id: &str, ) -> Result, BuildError> { - let headers = RequestHeaders::new(sdk_key, tags.as_deref(), instance_id); + let headers = RequestHeaders::new(sdk_key, tags, instance_id); let initializer_factories: Vec> = self .initializers @@ -398,9 +398,11 @@ impl DataSystemFactory for DataSystemBuilder { 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_err(|e| { - BuildError::InvalidConfig(format!("failed to build FDv1 fallback source: {e}")) - })?; + 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)); } diff --git a/launchdarkly-server-sdk/src/fdv2/mod.rs b/launchdarkly-server-sdk/src/fdv2/mod.rs index b9e458a..cf1880b 100644 --- a/launchdarkly-server-sdk/src/fdv2/mod.rs +++ b/launchdarkly-server-sdk/src/fdv2/mod.rs @@ -3,7 +3,7 @@ pub mod fdv1_adapter; pub mod model; pub mod polling; mod protocol; -mod request_headers; +pub mod request_headers; mod source; pub mod streaming; mod url; From 984d018006f7d0c6c7947d1005e87dff4a278aea Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Sun, 9 Aug 2026 19:51:50 -0700 Subject: [PATCH 4/5] chore: Narrow the FDv2 dead-code allow from the whole module to specific fields --- launchdarkly-server-sdk/src/fdv2/wire.rs | 5 +++++ launchdarkly-server-sdk/src/lib.rs | 1 - 2 files changed, 5 insertions(+), 1 deletion(-) 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 a8cb869..3a6358f 100644 --- a/launchdarkly-server-sdk/src/lib.rs +++ b/launchdarkly-server-sdk/src/lib.rs @@ -65,7 +65,6 @@ mod data_system; mod data_system_builders; mod evaluation; mod events; -#[allow(dead_code)] // Some items are consumed only by later data-system phases. mod fdv2; mod feature_requester; mod feature_requester_builders; From 483b5c53324b811e5f12723974065509a965bb7e Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Mon, 10 Aug 2026 15:36:01 -0700 Subject: [PATCH 5/5] test: Gate the FDv2 default-transport test and add injected-transport coverage --- .../src/data_system_builders.rs | 39 +++++++++++++++++++ 1 file changed, 39 insertions(+) diff --git a/launchdarkly-server-sdk/src/data_system_builders.rs b/launchdarkly-server-sdk/src/data_system_builders.rs index 9490689..a036474 100644 --- a/launchdarkly-server-sdk/src/data_system_builders.rs +++ b/launchdarkly-server-sdk/src/data_system_builders.rs @@ -419,6 +419,9 @@ impl DataSystemFactory for DataSystemBuilder { #[cfg(test)] mod tests { + use bytes::Bytes; + use launchdarkly_sdk_transport::{Request, ResponseFuture}; + use super::*; #[test] @@ -462,7 +465,43 @@ mod tests { 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");