1use eval::Context;
2use parking_lot::RwLock;
3use std::io;
4use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
5use std::sync::{Arc, Mutex};
6use std::time::Duration;
7use tokio::runtime::Runtime;
8
9use launchdarkly_server_sdk_evaluation::{self as eval, Detail, FlagValue, PrerequisiteEvent};
10use serde::Serialize;
11use thiserror::Error;
12use tokio::sync::{broadcast, Semaphore};
13
14use super::config::Config;
15use super::data_source_builders::BuildError as DataSourceError;
16use super::data_system::{DataSystem, FDv1DataSystem};
17use super::evaluation::{FlagDetail, FlagDetailConfig};
18use super::stores::store::DataStore;
19use super::stores::store_builders::BuildError as DataStoreError;
20use crate::config::BuildError as ConfigBuildError;
21use crate::events::event::EventFactory;
22use crate::events::event::InputEvent;
23use crate::events::processor::EventProcessor;
24use crate::events::processor_builders::BuildError as EventProcessorError;
25use crate::{MigrationOpTracker, Stage};
26
27struct EventsScope {
28 disabled: bool,
29 event_factory: EventFactory,
30 prerequisite_event_recorder: Box<dyn eval::PrerequisiteEventRecorder + Send + Sync>,
31}
32
33struct PrerequisiteEventRecorder {
34 event_factory: EventFactory,
35 event_processor: Arc<dyn EventProcessor>,
36}
37
38impl eval::PrerequisiteEventRecorder for PrerequisiteEventRecorder {
39 fn record(&self, event: PrerequisiteEvent) {
40 let evt = self.event_factory.new_eval_event(
41 &event.prerequisite_flag.key,
42 event.context.clone(),
43 &event.prerequisite_flag,
44 event.prerequisite_result,
45 FlagValue::Json(serde_json::Value::Null),
46 Some(event.target_flag_key),
47 );
48
49 self.event_processor.send(evt);
50 }
51}
52
53#[non_exhaustive]
55#[derive(Debug, Error)]
56pub enum BuildError {
57 #[error("invalid client config: {0}")]
59 InvalidConfig(String),
60}
61
62impl From<DataSourceError> for BuildError {
63 fn from(error: DataSourceError) -> Self {
64 Self::InvalidConfig(error.to_string())
65 }
66}
67
68impl From<DataStoreError> for BuildError {
69 fn from(error: DataStoreError) -> Self {
70 Self::InvalidConfig(error.to_string())
71 }
72}
73
74impl From<EventProcessorError> for BuildError {
75 fn from(error: EventProcessorError) -> Self {
76 Self::InvalidConfig(error.to_string())
77 }
78}
79
80impl From<ConfigBuildError> for BuildError {
81 fn from(error: ConfigBuildError) -> Self {
82 Self::InvalidConfig(error.to_string())
83 }
84}
85
86#[non_exhaustive]
88#[derive(Debug, Error)]
89pub enum StartError {
90 #[error("couldn't spawn background thread for client: {0}")]
92 SpawnFailed(io::Error),
93}
94
95#[derive(PartialEq, Copy, Clone, Debug)]
96enum ClientInitState {
97 Initializing = 0,
98 Initialized = 1,
99 InitializationFailed = 2,
100}
101
102impl PartialEq<usize> for ClientInitState {
103 fn eq(&self, other: &usize) -> bool {
104 *self as usize == *other
105 }
106}
107
108impl From<usize> for ClientInitState {
109 fn from(val: usize) -> Self {
110 match val {
111 0 => ClientInitState::Initializing,
112 1 => ClientInitState::Initialized,
113 2 => ClientInitState::InitializationFailed,
114 _ => unreachable!(),
115 }
116 }
117}
118
119pub struct Client {
149 event_processor: Arc<dyn EventProcessor>,
150 data_system: Arc<dyn DataSystem>,
151 data_store: Arc<RwLock<dyn DataStore>>,
152 events_default: EventsScope,
153 events_with_reasons: EventsScope,
154 init_notify: Arc<Semaphore>,
155 init_state: Arc<AtomicUsize>,
156 started: AtomicBool,
157 offline: bool,
158 daemon_mode: bool,
159 #[cfg_attr(
160 not(any(feature = "crypto-openssl", feature = "crypto-aws-lc-rs")),
161 allow(dead_code)
162 )]
163 sdk_key: String,
164 shutdown_broadcast: broadcast::Sender<()>,
165 runtime: RwLock<Option<Runtime>>,
166}
167
168impl Client {
169 pub fn build(config: Config) -> Result<Self, BuildError> {
171 if config.offline() {
172 info!("Started LaunchDarkly Client in offline mode");
173 } else if config.daemon_mode() {
174 info!("Started LaunchDarkly Client in daemon mode");
175 }
176
177 let tags = config.application_tag();
178 let instance_id = config.instance_id().to_string();
179
180 let endpoints = config.service_endpoints_builder().build()?;
181
182 let mut event_processor_builder = config.event_processor_builder().to_owned();
183 event_processor_builder.set_instance_id(instance_id.clone());
184 let event_processor =
185 event_processor_builder.build(&endpoints, config.sdk_key(), tags.clone())?;
186
187 let mut data_source_builder = config.data_source_builder().to_owned();
188 data_source_builder.set_instance_id(instance_id);
189 let data_source = data_source_builder.build(&endpoints, config.sdk_key(), tags.clone())?;
190 let data_system: Arc<dyn DataSystem> = Arc::new(FDv1DataSystem::new(
191 data_source,
192 config.data_store_builder(),
193 )?);
194 let data_store = data_system.store();
195
196 let events_default = EventsScope {
197 disabled: config.offline(),
198 event_factory: EventFactory::new(false),
199 prerequisite_event_recorder: Box::new(PrerequisiteEventRecorder {
200 event_factory: EventFactory::new(false),
201 event_processor: event_processor.clone(),
202 }),
203 };
204
205 let events_with_reasons = EventsScope {
206 disabled: config.offline(),
207 event_factory: EventFactory::new(true),
208 prerequisite_event_recorder: Box::new(PrerequisiteEventRecorder {
209 event_factory: EventFactory::new(true),
210 event_processor: event_processor.clone(),
211 }),
212 };
213
214 let (shutdown_tx, _) = broadcast::channel(1);
215
216 Ok(Client {
217 event_processor,
218 data_system,
219 data_store,
220 events_default,
221 events_with_reasons,
222 init_notify: Arc::new(Semaphore::new(0)),
223 init_state: Arc::new(AtomicUsize::new(ClientInitState::Initializing as usize)),
224 started: AtomicBool::new(false),
225 offline: config.offline(),
226 daemon_mode: config.daemon_mode(),
227 sdk_key: config.sdk_key().into(),
228 shutdown_broadcast: shutdown_tx,
229 runtime: RwLock::new(None),
230 })
231 }
232
233 pub fn start_with_default_executor(&self) {
235 if self.started.load(Ordering::SeqCst) {
236 return;
237 }
238 self.started.store(true, Ordering::SeqCst);
239 self.start_with_default_executor_internal();
240 }
241
242 fn start_with_default_executor_internal(&self) {
243 let notify = self.init_notify.clone();
247 let init_state = self.init_state.clone();
248
249 self.data_system.start(
250 Arc::new(move |success| {
251 init_state.store(
252 (if success {
253 ClientInitState::Initialized
254 } else {
255 ClientInitState::InitializationFailed
256 }) as usize,
257 Ordering::SeqCst,
258 );
259 notify.add_permits(1);
260 }),
261 self.shutdown_broadcast.subscribe(),
262 );
263 }
264
265 pub fn start_with_runtime(&self) -> Result<bool, StartError> {
271 if self.started.load(Ordering::SeqCst) {
272 return Ok(true);
273 }
274 self.started.store(true, Ordering::SeqCst);
275
276 let runtime = Runtime::new().map_err(StartError::SpawnFailed)?;
277 let _guard = runtime.enter();
278 self.runtime.write().replace(runtime);
279
280 self.start_with_default_executor_internal();
281
282 Ok(true)
283 }
284
285 pub async fn wait_for_initialization(&self, timeout: Duration) -> Option<bool> {
291 if timeout > Duration::from_secs(60) {
292 warn!("wait_for_initialization was configured to block for up to {} seconds. We recommend blocking no longer than 60 seconds.", timeout.as_secs());
293 }
294
295 let initialized = tokio::time::timeout(timeout, self.initialized_async_internal()).await;
296 initialized.ok()
297 }
298
299 async fn initialized_async_internal(&self) -> bool {
300 if self.offline || self.daemon_mode {
301 return true;
302 }
303
304 if ClientInitState::Initialized != self.init_state.load(Ordering::SeqCst) {
310 let _permit = self.init_notify.acquire().await;
311 }
312 ClientInitState::Initialized == self.init_state.load(Ordering::SeqCst)
313 }
314
315 pub fn initialized(&self) -> bool {
319 self.offline
320 || self.daemon_mode
321 || ClientInitState::Initialized == self.init_state.load(Ordering::SeqCst)
322 }
323
324 pub fn close(&self) {
328 self.event_processor.close();
329
330 if !self.offline && !self.daemon_mode {
333 if let Err(e) = self.shutdown_broadcast.send(()) {
334 error!("Failed to shutdown client appropriately: {e}");
335 }
336 }
337
338 self.runtime.write().take();
341 }
342
343 pub fn flush(&self) {
351 self.event_processor.flush();
352 }
353
354 pub async fn flush_blocking(&self, timeout: Duration) -> bool {
390 let event_processor = self.event_processor.clone();
391
392 let flush_future =
393 tokio::task::spawn_blocking(move || event_processor.flush_blocking(timeout));
394
395 if timeout == Duration::ZERO {
396 flush_future.await.unwrap_or(false)
398 } else {
399 match tokio::time::timeout(timeout, flush_future).await {
401 Ok(Ok(result)) => result,
402 Ok(Err(_)) => false, Err(_) => false, }
405 }
406 }
407
408 pub fn identify(&self, context: Context) {
413 if self.events_default.disabled {
414 return;
415 }
416
417 self.send_internal(self.events_default.event_factory.new_identify(context));
418 }
419
420 pub fn bool_variation(&self, context: &Context, flag_key: &str, default: bool) -> bool {
428 let val = self.variation(context, flag_key, default);
429 if let Some(b) = val.as_bool() {
430 b
431 } else {
432 warn!("bool_variation called for a non-bool flag {flag_key:?} (got {val:?})");
433 default
434 }
435 }
436
437 pub fn str_variation(&self, context: &Context, flag_key: &str, default: String) -> String {
445 let val = self.variation(context, flag_key, default.clone());
446 if let Some(s) = val.as_string() {
447 s
448 } else {
449 warn!("str_variation called for a non-string flag {flag_key:?} (got {val:?})");
450 default
451 }
452 }
453
454 pub fn float_variation(&self, context: &Context, flag_key: &str, default: f64) -> f64 {
462 let val = self.variation(context, flag_key, default);
463 if let Some(f) = val.as_float() {
464 f
465 } else {
466 warn!("float_variation called for a non-float flag {flag_key:?} (got {val:?})");
467 default
468 }
469 }
470
471 pub fn int_variation(&self, context: &Context, flag_key: &str, default: i64) -> i64 {
479 let val = self.variation(context, flag_key, default);
480 if let Some(f) = val.as_int() {
481 f
482 } else {
483 warn!("int_variation called for a non-int flag {flag_key:?} (got {val:?})");
484 default
485 }
486 }
487
488 pub fn json_variation(
498 &self,
499 context: &Context,
500 flag_key: &str,
501 default: serde_json::Value,
502 ) -> serde_json::Value {
503 self.variation(context, flag_key, default.clone())
504 .as_json()
505 .unwrap_or(default)
506 }
507
508 pub fn bool_variation_detail(
515 &self,
516 context: &Context,
517 flag_key: &str,
518 default: bool,
519 ) -> Detail<bool> {
520 self.variation_detail(context, flag_key, default).try_map(
521 |val| val.as_bool(),
522 default,
523 eval::Error::WrongType,
524 )
525 }
526
527 pub fn str_variation_detail(
534 &self,
535 context: &Context,
536 flag_key: &str,
537 default: String,
538 ) -> Detail<String> {
539 self.variation_detail(context, flag_key, default.clone())
540 .try_map(|val| val.as_string(), default, eval::Error::WrongType)
541 }
542
543 pub fn float_variation_detail(
550 &self,
551 context: &Context,
552 flag_key: &str,
553 default: f64,
554 ) -> Detail<f64> {
555 self.variation_detail(context, flag_key, default).try_map(
556 |val| val.as_float(),
557 default,
558 eval::Error::WrongType,
559 )
560 }
561
562 pub fn int_variation_detail(
569 &self,
570 context: &Context,
571 flag_key: &str,
572 default: i64,
573 ) -> Detail<i64> {
574 self.variation_detail(context, flag_key, default).try_map(
575 |val| val.as_int(),
576 default,
577 eval::Error::WrongType,
578 )
579 }
580
581 pub fn json_variation_detail(
588 &self,
589 context: &Context,
590 flag_key: &str,
591 default: serde_json::Value,
592 ) -> Detail<serde_json::Value> {
593 self.variation_detail(context, flag_key, default.clone())
594 .try_map(|val| val.as_json(), default, eval::Error::WrongType)
595 }
596
597 #[cfg(any(feature = "crypto-aws-lc-rs", feature = "crypto-openssl"))]
598 pub fn secure_mode_hash(&self, context: &Context) -> Result<String, String> {
603 #[cfg(feature = "crypto-aws-lc-rs")]
604 {
605 let key =
606 aws_lc_rs::hmac::Key::new(aws_lc_rs::hmac::HMAC_SHA256, self.sdk_key.as_bytes());
607 let tag = aws_lc_rs::hmac::sign(&key, context.canonical_key().as_bytes());
608
609 Ok(data_encoding::HEXLOWER.encode(tag.as_ref()))
610 }
611 #[cfg(feature = "crypto-openssl")]
612 {
613 use openssl::hash::MessageDigest;
614 use openssl::pkey::PKey;
615 use openssl::sign::Signer;
616
617 let key = PKey::hmac(self.sdk_key.as_bytes())
618 .map_err(|e| format!("Failed to create HMAC key: {e}"))?;
619 let mut signer = Signer::new(MessageDigest::sha256(), &key)
620 .map_err(|e| format!("Failed to create signer: {e}"))?;
621 signer
622 .update(context.canonical_key().as_bytes())
623 .map_err(|e| format!("Failed to update signer: {e}"))?;
624 let hmac = signer
625 .sign_to_vec()
626 .map_err(|e| format!("Failed to sign: {e}"))?;
627
628 Ok(data_encoding::HEXLOWER.encode(&hmac))
629 }
630 }
631
632 pub fn all_flags_detail(
643 &self,
644 context: &Context,
645 flag_state_config: FlagDetailConfig,
646 ) -> FlagDetail {
647 if self.offline {
648 warn!(
649 "all_flags_detail() called, but client is in offline mode. Returning empty state"
650 );
651 return FlagDetail::new(false);
652 }
653
654 if !self.initialized() {
655 warn!("all_flags_detail() called before client has finished initializing! Feature store unavailable - returning empty state");
656 return FlagDetail::new(false);
657 }
658
659 let data_store = self.data_store.read();
660
661 let mut flag_detail = FlagDetail::new(true);
662 flag_detail.populate(&*data_store, context, flag_state_config);
663
664 flag_detail
665 }
666
667 pub fn variation_detail<T: Into<FlagValue> + Clone>(
673 &self,
674 context: &Context,
675 flag_key: &str,
676 default: T,
677 ) -> Detail<FlagValue> {
678 let (detail, _) =
679 self.variation_internal(context, flag_key, default, &self.events_with_reasons);
680 detail
681 }
682
683 pub fn variation<T: Into<FlagValue> + Clone>(
694 &self,
695 context: &Context,
696 flag_key: &str,
697 default: T,
698 ) -> FlagValue {
699 let (detail, _) = self.variation_internal(context, flag_key, default, &self.events_default);
700 detail.value.unwrap()
701 }
702
703 pub fn migration_variation(
708 &self,
709 context: &Context,
710 flag_key: &str,
711 default_stage: Stage,
712 ) -> (Stage, Arc<Mutex<MigrationOpTracker>>) {
713 let (detail, flag) =
714 self.variation_internal(context, flag_key, default_stage, &self.events_default);
715
716 let migration_detail =
717 detail.try_map(|v| v.try_into().ok(), default_stage, eval::Error::WrongType);
718 let tracker = MigrationOpTracker::new(
719 flag_key.into(),
720 flag,
721 context.clone(),
722 migration_detail.clone(),
723 default_stage,
724 );
725
726 (
727 migration_detail.value.unwrap_or(default_stage),
728 Arc::new(Mutex::new(tracker)),
729 )
730 }
731
732 pub fn track_event(&self, context: Context, key: impl Into<String>) {
742 let _ = self.track(context, key, None, serde_json::Value::Null);
743 }
744
745 pub fn track_data(
758 &self,
759 context: Context,
760 key: impl Into<String>,
761 data: impl Serialize,
762 ) -> serde_json::Result<()> {
763 self.track(context, key, None, data)
764 }
765
766 pub fn track_metric(
777 &self,
778 context: Context,
779 key: impl Into<String>,
780 value: f64,
781 data: impl Serialize,
782 ) {
783 let _ = self.track(context, key, Some(value), data);
784 }
785
786 fn track(
787 &self,
788 context: Context,
789 key: impl Into<String>,
790 metric_value: Option<f64>,
791 data: impl Serialize,
792 ) -> serde_json::Result<()> {
793 if !self.events_default.disabled {
794 let event =
795 self.events_default
796 .event_factory
797 .new_custom(context, key, metric_value, data)?;
798
799 self.send_internal(event);
800 }
801
802 Ok(())
803 }
804
805 pub fn track_migration_op(&self, tracker: Arc<Mutex<MigrationOpTracker>>) {
812 if self.events_default.disabled {
813 return;
814 }
815
816 match tracker.lock() {
817 Ok(tracker) => {
818 let event = tracker.build();
819 match event {
820 Ok(event) => {
821 self.send_internal(
822 self.events_default.event_factory.new_migration_op(event),
823 );
824 }
825 Err(e) => error!("Failed to build migration event, no event will be sent: {e}"),
826 }
827 }
828 Err(e) => error!("Failed to lock migration tracker, no event will be sent: {e}"),
829 }
830 }
831
832 fn variation_internal<T: Into<FlagValue> + Clone>(
833 &self,
834 context: &Context,
835 flag_key: &str,
836 default: T,
837 events_scope: &EventsScope,
838 ) -> (Detail<FlagValue>, Option<eval::Flag>) {
839 if self.offline {
840 return (
841 Detail::err_default(eval::Error::ClientNotReady, default.into()),
842 None,
843 );
844 }
845
846 let (flag, result) = match self.initialized() {
847 false => (
848 None,
849 Detail::err_default(eval::Error::ClientNotReady, default.clone().into()),
850 ),
851 true => {
852 let data_store = self.data_store.read();
853 match data_store.flag(flag_key) {
854 Some(flag) => {
855 let result = eval::evaluate(
856 data_store.to_store(),
857 &flag,
858 context,
859 Some(&*events_scope.prerequisite_event_recorder),
860 )
861 .map(|v| v.clone())
862 .or(default.clone().into());
863
864 (Some(flag), result)
865 }
866 None => (
867 None,
868 Detail::err_default(eval::Error::FlagNotFound, default.clone().into()),
869 ),
870 }
871 }
872 };
873
874 if !events_scope.disabled {
875 let event = match &flag {
876 Some(f) => events_scope.event_factory.new_eval_event(
877 flag_key,
878 context.clone(),
879 f,
880 result.clone(),
881 default.into(),
882 None,
883 ),
884 None => events_scope.event_factory.new_unknown_flag_event(
885 flag_key,
886 context.clone(),
887 result.clone(),
888 default.into(),
889 ),
890 };
891 self.send_internal(event);
892 }
893
894 (result, flag)
895 }
896
897 fn send_internal(&self, event: InputEvent) {
898 self.event_processor.send(event);
899 }
900}
901
902#[cfg(test)]
903mod tests {
904 use assert_json_diff::assert_json_eq;
905 use crossbeam_channel::Receiver;
906 use eval::{ContextBuilder, MultiContextBuilder};
907 use futures::FutureExt;
908 use launchdarkly_server_sdk_evaluation::{Flag, Reason, Segment};
909 use maplit::hashmap;
910 use std::collections::HashMap;
911 use tokio::time::Instant;
912
913 use crate::data_source::MockDataSource;
914 use crate::data_source_builders::MockDataSourceBuilder;
915 use crate::evaluation::FlagFilter;
916 use crate::events::create_event_sender;
917 use crate::events::event::{OutputEvent, VariationKey};
918 use crate::events::processor_builders::EventProcessorBuilder;
919 use crate::stores::persistent_store::tests::InMemoryPersistentDataStore;
920 use crate::stores::store_types::{PatchTarget, StorageItem};
921 use crate::test_common::{
922 self, basic_flag, basic_flag_with_prereq, basic_flag_with_prereqs_and_visibility,
923 basic_flag_with_visibility, basic_int_flag, basic_migration_flag, basic_off_flag,
924 };
925 use crate::test_data::TestData;
926 use crate::{
927 AllData, ConfigBuilder, MigratorBuilder, NullEventProcessorBuilder, Operation, Origin,
928 PersistentDataStore, PersistentDataStoreBuilder, PersistentDataStoreFactory,
929 SerializedItem,
930 };
931 use test_case::test_case;
932
933 use super::*;
934
935 fn is_send_and_sync<T: Send + Sync>() {}
936
937 #[test]
938 fn ensure_client_is_send_and_sync() {
939 is_send_and_sync::<Client>()
940 }
941
942 #[tokio::test]
943 async fn client_asynchronously_initializes_within_timeout() {
944 let (client, _event_rx) = make_mocked_client_with_delay(1000, false, false);
945 client.start_with_default_executor();
946
947 let now = Instant::now();
948 let initialized = client
949 .wait_for_initialization(Duration::from_millis(1500))
950 .await;
951 let elapsed_time = now.elapsed();
952 assert!(elapsed_time.as_millis() > 500);
954 assert_eq!(initialized, Some(true));
955 }
956
957 #[tokio::test]
958 async fn client_asynchronously_initializes_slower_than_timeout() {
959 let (client, _event_rx) = make_mocked_client_with_delay(2000, false, false);
960 client.start_with_default_executor();
961
962 let now = Instant::now();
963 let initialized = client
964 .wait_for_initialization(Duration::from_millis(500))
965 .await;
966 let elapsed_time = now.elapsed();
967 assert!(elapsed_time.as_millis() < 750);
969 assert!(initialized.is_none());
970 }
971
972 #[tokio::test]
973 async fn client_initializes_immediately_in_offline_mode() {
974 let (client, _event_rx) = make_mocked_client_with_delay(1000, true, false);
975 client.start_with_default_executor();
976
977 assert!(client.initialized());
978
979 let now = Instant::now();
980 let initialized = client
981 .wait_for_initialization(Duration::from_millis(2000))
982 .await;
983 let elapsed_time = now.elapsed();
984 assert_eq!(initialized, Some(true));
985 assert!(elapsed_time.as_millis() < 500)
986 }
987
988 #[tokio::test]
989 async fn client_initializes_immediately_in_daemon_mode() {
990 let (client, _event_rx) = make_mocked_client_with_delay(1000, false, true);
991 client.start_with_default_executor();
992
993 assert!(client.initialized());
994
995 let now = Instant::now();
996 let initialized = client
997 .wait_for_initialization(Duration::from_millis(2000))
998 .await;
999 let elapsed_time = now.elapsed();
1000 assert_eq!(initialized, Some(true));
1001 assert!(elapsed_time.as_millis() < 500)
1002 }
1003
1004 #[test_case(basic_flag("myFlag"), false.into(), true.into())]
1005 #[test_case(basic_int_flag("myFlag"), 0.into(), test_common::FLOAT_TO_INT_MAX.into())]
1006 fn client_updates_changes_evaluation_results(
1007 flag: eval::Flag,
1008 default: FlagValue,
1009 expected: FlagValue,
1010 ) {
1011 let context = ContextBuilder::new("foo")
1012 .build()
1013 .expect("Failed to create context");
1014
1015 let (client, _event_rx) = make_mocked_client();
1016
1017 let result = client.variation_detail(&context, "myFlag", default.clone());
1018 assert_eq!(result.value.unwrap(), default);
1019
1020 client.start_with_default_executor();
1021 client
1022 .data_store
1023 .write()
1024 .upsert(
1025 &flag.key,
1026 PatchTarget::Flag(StorageItem::Item(flag.clone())),
1027 )
1028 .expect("patch should apply");
1029
1030 let result = client.variation_detail(&context, "myFlag", default);
1031 assert_eq!(result.value.unwrap(), expected);
1032 assert!(matches!(
1033 result.reason,
1034 Reason::Fallthrough {
1035 in_experiment: false
1036 }
1037 ));
1038 }
1039
1040 #[test]
1041 fn all_flags_detail_is_invalid_when_offline() {
1042 let (client, _event_rx) = make_mocked_offline_client();
1043 client.start_with_default_executor();
1044
1045 let context = ContextBuilder::new("bob")
1046 .build()
1047 .expect("Failed to create context");
1048
1049 let all_flags = client.all_flags_detail(&context, FlagDetailConfig::new());
1050 assert_json_eq!(all_flags, json!({"$valid": false, "$flagsState" : {}}));
1051 }
1052
1053 #[test]
1054 fn all_flags_detail_is_invalid_when_not_initialized() {
1055 let (client, _event_rx) = make_mocked_client();
1056
1057 let context = ContextBuilder::new("bob")
1058 .build()
1059 .expect("Failed to create context");
1060
1061 let all_flags = client.all_flags_detail(&context, FlagDetailConfig::new());
1062 assert_json_eq!(all_flags, json!({"$valid": false, "$flagsState" : {}}));
1063 }
1064
1065 #[tokio::test]
1066 async fn all_flags_detail_returns_flag_states() {
1067 let td = TestData::new();
1068 td.use_preconfigured_flag(basic_flag("myFlag1"));
1069 td.use_preconfigured_flag(basic_flag("myFlag2"));
1070 let (client, _event_rx) = make_client_with_test_data(&td);
1071 client.start_with_default_executor();
1072
1073 let context = ContextBuilder::new("bob")
1074 .build()
1075 .expect("Failed to create context");
1076
1077 let all_flags = client.all_flags_detail(&context, FlagDetailConfig::new());
1078
1079 client.close();
1080
1081 assert_json_eq!(
1082 all_flags,
1083 json!({
1084 "myFlag1": true,
1085 "myFlag2": true,
1086 "$flagsState": {
1087 "myFlag1": {
1088 "version": 1,
1089 "variation": 1
1090 },
1091 "myFlag2": {
1092 "version": 1,
1093 "variation": 1
1094 },
1095 },
1096 "$valid": true
1097 })
1098 );
1099 }
1100
1101 #[tokio::test]
1102 async fn all_flags_detail_returns_prerequisite_relations() {
1103 let td = TestData::new();
1104 td.use_preconfigured_flag(basic_flag("prereq1"));
1105 td.use_preconfigured_flag(basic_flag("prereq2"));
1106 td.use_preconfigured_flag(basic_flag_with_prereqs_and_visibility(
1107 "toplevel",
1108 &["prereq1", "prereq2"],
1109 false,
1110 false,
1111 ));
1112 let (client, _event_rx) = make_client_with_test_data(&td);
1113 client.start_with_default_executor();
1114
1115 let context = ContextBuilder::new("bob")
1116 .build()
1117 .expect("Failed to create context");
1118
1119 let all_flags = client.all_flags_detail(&context, FlagDetailConfig::new());
1120
1121 client.close();
1122
1123 assert_json_eq!(
1124 all_flags,
1125 json!({
1126 "prereq1": true,
1127 "prereq2": true,
1128 "toplevel": true,
1129 "$flagsState": {
1130 "toplevel": {
1131 "version": 1,
1132 "variation": 1,
1133 "prerequisites": ["prereq1", "prereq2"]
1134 },
1135 "prereq1": {
1136 "version": 1,
1137 "variation": 1
1138 },
1139 "prereq2": {
1140 "version": 1,
1141 "variation": 1
1142 },
1143 },
1144 "$valid": true
1145 })
1146 );
1147 }
1148
1149 #[tokio::test]
1150 async fn all_flags_detail_returns_prerequisite_relations_when_not_visible_to_clients() {
1151 let td = TestData::new();
1152 td.use_preconfigured_flag(basic_flag_with_visibility("prereq1", false, false));
1153 td.use_preconfigured_flag(basic_flag_with_visibility("prereq2", false, false));
1154 td.use_preconfigured_flag(basic_flag_with_prereqs_and_visibility(
1155 "toplevel",
1156 &["prereq1", "prereq2"],
1157 true,
1158 false,
1159 ));
1160 let (client, _event_rx) = make_client_with_test_data(&td);
1161 client.start_with_default_executor();
1162
1163 let context = ContextBuilder::new("bob")
1164 .build()
1165 .expect("Failed to create context");
1166
1167 let mut config = FlagDetailConfig::new();
1168 config.flag_filter(FlagFilter::CLIENT);
1169
1170 let all_flags = client.all_flags_detail(&context, config);
1171
1172 client.close();
1173
1174 assert_json_eq!(
1175 all_flags,
1176 json!({
1177 "toplevel": true,
1178 "$flagsState": {
1179 "toplevel": {
1180 "version": 1,
1181 "variation": 1,
1182 "prerequisites": ["prereq1", "prereq2"]
1183 },
1184 },
1185 "$valid": true
1186 })
1187 );
1188 }
1189
1190 #[tokio::test]
1191 async fn variation_tracks_events_correctly() {
1192 let td = TestData::new();
1193 td.use_preconfigured_flag(basic_flag("myFlag"));
1194 let (client, event_rx) = make_client_with_test_data(&td);
1195 client.start_with_default_executor();
1196
1197 let context = ContextBuilder::new("bob")
1198 .build()
1199 .expect("Failed to create context");
1200
1201 let flag_value = client.variation(&context, "myFlag", FlagValue::Bool(false));
1202
1203 assert!(flag_value.as_bool().unwrap());
1204 client.flush();
1205 client.close();
1206
1207 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1208 assert_eq!(events.len(), 2);
1209 assert_eq!(events[0].kind(), "index");
1210 assert_eq!(events[1].kind(), "summary");
1211
1212 if let OutputEvent::Summary(event_summary) = events[1].clone() {
1213 let variation_key = VariationKey {
1214 version: Some(1),
1215 variation: Some(1),
1216 };
1217 let feature = event_summary.features.get("myFlag");
1218 assert!(feature.is_some());
1219
1220 let feature = feature.unwrap();
1221 assert!(feature.counters.contains_key(&variation_key));
1222 } else {
1223 panic!("Event should be a summary type");
1224 }
1225 }
1226
1227 #[test]
1228 fn variation_handles_offline_mode() {
1229 let (client, event_rx) = make_mocked_offline_client();
1230 client.start_with_default_executor();
1231
1232 let context = ContextBuilder::new("bob")
1233 .build()
1234 .expect("Failed to create context");
1235 let flag_value = client.variation(&context, "myFlag", FlagValue::Bool(false));
1236
1237 assert!(!flag_value.as_bool().unwrap());
1238 client.flush();
1239 client.close();
1240
1241 assert_eq!(event_rx.iter().count(), 0);
1242 }
1243
1244 #[test]
1245 fn variation_handles_unknown_flags() {
1246 let (client, event_rx) = make_mocked_client();
1247 client.start_with_default_executor();
1248 let context = ContextBuilder::new("bob")
1249 .build()
1250 .expect("Failed to create context");
1251
1252 let flag_value = client.variation(&context, "non-existent-flag", FlagValue::Bool(false));
1253
1254 assert!(!flag_value.as_bool().unwrap());
1255 client.flush();
1256 client.close();
1257
1258 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1259 assert_eq!(events.len(), 2);
1260 assert_eq!(events[0].kind(), "index");
1261 assert_eq!(events[1].kind(), "summary");
1262
1263 if let OutputEvent::Summary(event_summary) = events[1].clone() {
1264 let variation_key = VariationKey {
1265 version: None,
1266 variation: None,
1267 };
1268
1269 let feature = event_summary.features.get("non-existent-flag");
1270 assert!(feature.is_some());
1271
1272 let feature = feature.unwrap();
1273 assert!(feature.counters.contains_key(&variation_key));
1274 } else {
1275 panic!("Event should be a summary type");
1276 }
1277 }
1278
1279 #[tokio::test]
1280 async fn variation_detail_handles_debug_events_correctly() {
1281 let td = TestData::new();
1282 let mut flag = basic_flag("myFlag");
1283 flag.debug_events_until_date = Some(64_060_606_800_000); td.use_preconfigured_flag(flag);
1285 let (client, event_rx) = make_client_with_test_data(&td);
1286 client.start_with_default_executor();
1287
1288 let context = ContextBuilder::new("bob")
1289 .build()
1290 .expect("Failed to create context");
1291
1292 let detail = client.variation_detail(&context, "myFlag", FlagValue::Bool(false));
1293
1294 assert!(detail.value.unwrap().as_bool().unwrap());
1295 assert!(matches!(
1296 detail.reason,
1297 Reason::Fallthrough {
1298 in_experiment: false
1299 }
1300 ));
1301 client.flush();
1302 client.close();
1303
1304 let events = event_rx.try_iter().collect::<Vec<OutputEvent>>();
1305 assert_eq!(events.len(), 3);
1306 assert_eq!(events[0].kind(), "index");
1307 assert_eq!(events[1].kind(), "debug");
1308 assert_eq!(events[2].kind(), "summary");
1309
1310 if let OutputEvent::Summary(event_summary) = events[2].clone() {
1311 let variation_key = VariationKey {
1312 version: Some(1),
1313 variation: Some(1),
1314 };
1315
1316 let feature = event_summary.features.get("myFlag");
1317 assert!(feature.is_some());
1318
1319 let feature = feature.unwrap();
1320 assert!(feature.counters.contains_key(&variation_key));
1321 } else {
1322 panic!("Event should be a summary type");
1323 }
1324 }
1325
1326 #[tokio::test]
1327 async fn variation_detail_tracks_events_correctly() {
1328 let td = TestData::new();
1329 td.use_preconfigured_flag(basic_flag("myFlag"));
1330 let (client, event_rx) = make_client_with_test_data(&td);
1331 client.start_with_default_executor();
1332
1333 let context = ContextBuilder::new("bob")
1334 .build()
1335 .expect("Failed to create context");
1336
1337 let detail = client.variation_detail(&context, "myFlag", FlagValue::Bool(false));
1338
1339 assert!(detail.value.unwrap().as_bool().unwrap());
1340 assert!(matches!(
1341 detail.reason,
1342 Reason::Fallthrough {
1343 in_experiment: false
1344 }
1345 ));
1346 client.flush();
1347 client.close();
1348
1349 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1350 assert_eq!(events.len(), 2);
1351 assert_eq!(events[0].kind(), "index");
1352 assert_eq!(events[1].kind(), "summary");
1353
1354 if let OutputEvent::Summary(event_summary) = events[1].clone() {
1355 let variation_key = VariationKey {
1356 version: Some(1),
1357 variation: Some(1),
1358 };
1359
1360 let feature = event_summary.features.get("myFlag");
1361 assert!(feature.is_some());
1362
1363 let feature = feature.unwrap();
1364 assert!(feature.counters.contains_key(&variation_key));
1365 } else {
1366 panic!("Event should be a summary type");
1367 }
1368 }
1369
1370 #[test]
1371 fn variation_detail_handles_offline_mode() {
1372 let (client, event_rx) = make_mocked_offline_client();
1373 client.start_with_default_executor();
1374
1375 let context = ContextBuilder::new("bob")
1376 .build()
1377 .expect("Failed to create context");
1378
1379 let detail = client.variation_detail(&context, "myFlag", FlagValue::Bool(false));
1380
1381 assert!(!detail.value.unwrap().as_bool().unwrap());
1382 assert!(matches!(
1383 detail.reason,
1384 Reason::Error {
1385 error: eval::Error::ClientNotReady
1386 }
1387 ));
1388 client.flush();
1389 client.close();
1390
1391 assert_eq!(event_rx.iter().count(), 0);
1392 }
1393
1394 struct InMemoryPersistentDataStoreFactory {
1395 data: AllData<Flag, Segment>,
1396 initialized: bool,
1397 }
1398
1399 impl PersistentDataStoreFactory for InMemoryPersistentDataStoreFactory {
1400 fn create_persistent_data_store(
1401 &self,
1402 ) -> Result<Box<dyn PersistentDataStore + 'static>, std::io::Error> {
1403 let serialized_data =
1404 AllData::<SerializedItem, SerializedItem>::try_from(self.data.clone())?;
1405 Ok(Box::new(InMemoryPersistentDataStore {
1406 data: serialized_data,
1407 initialized: self.initialized,
1408 }))
1409 }
1410 }
1411
1412 #[test]
1413 fn variation_detail_handles_daemon_mode() {
1414 testing_logger::setup();
1415 let factory = InMemoryPersistentDataStoreFactory {
1416 data: AllData {
1417 flags: hashmap!["flag".into() => basic_flag("flag")],
1418 segments: HashMap::new(),
1419 },
1420 initialized: true,
1421 };
1422 let builder = PersistentDataStoreBuilder::new(Arc::new(factory));
1423
1424 let config = ConfigBuilder::new("sdk-key")
1425 .daemon_mode(true)
1426 .data_store(&builder)
1427 .event_processor(&NullEventProcessorBuilder::new())
1428 .build()
1429 .expect("config should build");
1430
1431 let client = Client::build(config).expect("Should be built.");
1432
1433 client.start_with_default_executor();
1434
1435 let context = ContextBuilder::new("bob")
1436 .build()
1437 .expect("Failed to create context");
1438
1439 let detail = client.variation_detail(&context, "flag", FlagValue::Bool(false));
1440
1441 assert!(detail.value.unwrap().as_bool().unwrap());
1442 assert!(matches!(
1443 detail.reason,
1444 Reason::Fallthrough {
1445 in_experiment: false
1446 }
1447 ));
1448 client.flush();
1449 client.close();
1450
1451 testing_logger::validate(|captured_logs| {
1452 assert_eq!(captured_logs.len(), 1);
1453 assert_eq!(
1454 captured_logs[0].body,
1455 "Started LaunchDarkly Client in daemon mode"
1456 );
1457 });
1458 }
1459
1460 #[test]
1461 fn daemon_mode_is_quiet_if_store_is_not_initialized() {
1462 testing_logger::setup();
1463
1464 let factory = InMemoryPersistentDataStoreFactory {
1465 data: AllData {
1466 flags: HashMap::new(),
1467 segments: HashMap::new(),
1468 },
1469 initialized: false,
1470 };
1471 let builder = PersistentDataStoreBuilder::new(Arc::new(factory));
1472
1473 let config = ConfigBuilder::new("sdk-key")
1474 .daemon_mode(true)
1475 .data_store(&builder)
1476 .event_processor(&NullEventProcessorBuilder::new())
1477 .build()
1478 .expect("config should build");
1479
1480 let client = Client::build(config).expect("Should be built.");
1481
1482 client.start_with_default_executor();
1483
1484 let context = ContextBuilder::new("bob")
1485 .build()
1486 .expect("Failed to create context");
1487
1488 client.variation_detail(&context, "flag", FlagValue::Bool(false));
1489
1490 testing_logger::validate(|captured_logs| {
1491 assert_eq!(captured_logs.len(), 1);
1492 assert_eq!(
1493 captured_logs[0].body,
1494 "Started LaunchDarkly Client in daemon mode"
1495 );
1496 });
1497 }
1498
1499 #[tokio::test]
1500 async fn variation_handles_off_flag_without_variation() {
1501 let td = TestData::new();
1502 td.use_preconfigured_flag(basic_off_flag("myFlag"));
1503 let (client, event_rx) = make_client_with_test_data(&td);
1504 client.start_with_default_executor();
1505
1506 let context = ContextBuilder::new("bob")
1507 .build()
1508 .expect("Failed to create context");
1509
1510 let result = client.variation(&context, "myFlag", FlagValue::Bool(false));
1511
1512 assert!(!result.as_bool().unwrap());
1513 client.flush();
1514 client.close();
1515
1516 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1517 assert_eq!(events.len(), 2);
1518 assert_eq!(events[0].kind(), "index");
1519 assert_eq!(events[1].kind(), "summary");
1520
1521 if let OutputEvent::Summary(event_summary) = events[1].clone() {
1522 let variation_key = VariationKey {
1523 version: Some(1),
1524 variation: None,
1525 };
1526 let feature = event_summary.features.get("myFlag");
1527 assert!(feature.is_some());
1528
1529 let feature = feature.unwrap();
1530 assert!(feature.counters.contains_key(&variation_key));
1531 } else {
1532 panic!("Event should be a summary type");
1533 }
1534 }
1535
1536 #[tokio::test]
1537 async fn variation_detail_tracks_prereq_events_correctly() {
1538 let td = TestData::new();
1539 let mut prereq_flag = basic_flag("prereqFlag");
1540 prereq_flag.track_events = true;
1541 td.use_preconfigured_flag(prereq_flag);
1542
1543 let mut main_flag = basic_flag_with_prereq("myFlag", "prereqFlag");
1544 main_flag.track_events = true;
1545 td.use_preconfigured_flag(main_flag);
1546
1547 let (client, event_rx) = make_client_with_test_data(&td);
1548 client.start_with_default_executor();
1549
1550 let context = ContextBuilder::new("bob")
1551 .build()
1552 .expect("Failed to create context");
1553
1554 let detail = client.variation_detail(&context, "myFlag", FlagValue::Bool(false));
1555
1556 assert!(detail.value.unwrap().as_bool().unwrap());
1557 assert!(matches!(
1558 detail.reason,
1559 Reason::Fallthrough {
1560 in_experiment: false
1561 }
1562 ));
1563 client.flush();
1564 client.close();
1565
1566 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1567 assert_eq!(events.len(), 4);
1568 assert_eq!(events[0].kind(), "index");
1569 assert_eq!(events[1].kind(), "feature");
1570 assert_eq!(events[2].kind(), "feature");
1571 assert_eq!(events[3].kind(), "summary");
1572
1573 if let OutputEvent::Summary(event_summary) = events[3].clone() {
1574 let variation_key = VariationKey {
1575 version: Some(1),
1576 variation: Some(1),
1577 };
1578 let feature = event_summary.features.get("myFlag");
1579 assert!(feature.is_some());
1580
1581 let feature = feature.unwrap();
1582 assert!(feature.counters.contains_key(&variation_key));
1583
1584 let variation_key = VariationKey {
1585 version: Some(1),
1586 variation: Some(1),
1587 };
1588 let feature = event_summary.features.get("prereqFlag");
1589 assert!(feature.is_some());
1590
1591 let feature = feature.unwrap();
1592 assert!(feature.counters.contains_key(&variation_key));
1593 }
1594 }
1595
1596 #[tokio::test]
1597 async fn variation_handles_failed_prereqs_correctly() {
1598 let td = TestData::new();
1599 let mut prereq_flag = basic_off_flag("prereqFlag");
1600 prereq_flag.track_events = true;
1601 td.use_preconfigured_flag(prereq_flag);
1602
1603 let mut main_flag = basic_flag_with_prereq("myFlag", "prereqFlag");
1604 main_flag.track_events = true;
1605 td.use_preconfigured_flag(main_flag);
1606
1607 let (client, event_rx) = make_client_with_test_data(&td);
1608 client.start_with_default_executor();
1609
1610 let context = ContextBuilder::new("bob")
1611 .build()
1612 .expect("Failed to create context");
1613
1614 let detail = client.variation(&context, "myFlag", FlagValue::Bool(false));
1615
1616 assert!(!detail.as_bool().unwrap());
1617 client.flush();
1618 client.close();
1619
1620 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1621 assert_eq!(events.len(), 4);
1622 assert_eq!(events[0].kind(), "index");
1623 assert_eq!(events[1].kind(), "feature");
1624 assert_eq!(events[2].kind(), "feature");
1625 assert_eq!(events[3].kind(), "summary");
1626
1627 if let OutputEvent::Summary(event_summary) = events[3].clone() {
1628 let variation_key = VariationKey {
1629 version: Some(1),
1630 variation: Some(0),
1631 };
1632 let feature = event_summary.features.get("myFlag");
1633 assert!(feature.is_some());
1634
1635 let feature = feature.unwrap();
1636 assert!(feature.counters.contains_key(&variation_key));
1637
1638 let variation_key = VariationKey {
1639 version: Some(1),
1640 variation: None,
1641 };
1642 let feature = event_summary.features.get("prereqFlag");
1643 assert!(feature.is_some());
1644
1645 let feature = feature.unwrap();
1646 assert!(feature.counters.contains_key(&variation_key));
1647 }
1648 }
1649
1650 #[test]
1651 fn variation_detail_handles_flag_not_found() {
1652 let (client, event_rx) = make_mocked_client();
1653 client.start_with_default_executor();
1654
1655 let context = ContextBuilder::new("bob")
1656 .build()
1657 .expect("Failed to create context");
1658 let detail = client.variation_detail(&context, "non-existent-flag", FlagValue::Bool(false));
1659
1660 assert!(!detail.value.unwrap().as_bool().unwrap());
1661 assert!(matches!(
1662 detail.reason,
1663 Reason::Error {
1664 error: eval::Error::FlagNotFound
1665 }
1666 ));
1667 client.flush();
1668 client.close();
1669
1670 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1671 assert_eq!(events.len(), 2);
1672 assert_eq!(events[0].kind(), "index");
1673 assert_eq!(events[1].kind(), "summary");
1674
1675 if let OutputEvent::Summary(event_summary) = events[1].clone() {
1676 let variation_key = VariationKey {
1677 version: None,
1678 variation: None,
1679 };
1680 let feature = event_summary.features.get("non-existent-flag");
1681 assert!(feature.is_some());
1682
1683 let feature = feature.unwrap();
1684 assert!(feature.counters.contains_key(&variation_key));
1685 } else {
1686 panic!("Event should be a summary type");
1687 }
1688 }
1689
1690 #[tokio::test]
1691 async fn variation_detail_handles_client_not_ready() {
1692 let (client, event_rx) = make_mocked_client_with_delay(u64::MAX, false, false);
1693 client.start_with_default_executor();
1694 let context = ContextBuilder::new("bob")
1695 .build()
1696 .expect("Failed to create context");
1697
1698 let detail = client.variation_detail(&context, "non-existent-flag", FlagValue::Bool(false));
1699
1700 assert!(!detail.value.unwrap().as_bool().unwrap());
1701 assert!(matches!(
1702 detail.reason,
1703 Reason::Error {
1704 error: eval::Error::ClientNotReady
1705 }
1706 ));
1707 client.flush();
1708 client.close();
1709
1710 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1711 assert_eq!(events.len(), 2);
1712 assert_eq!(events[0].kind(), "index");
1713 assert_eq!(events[1].kind(), "summary");
1714
1715 if let OutputEvent::Summary(event_summary) = events[1].clone() {
1716 let variation_key = VariationKey {
1717 version: None,
1718 variation: None,
1719 };
1720 let feature = event_summary.features.get("non-existent-flag");
1721 assert!(feature.is_some());
1722
1723 let feature = feature.unwrap();
1724 assert!(feature.counters.contains_key(&variation_key));
1725 } else {
1726 panic!("Event should be a summary type");
1727 }
1728 }
1729
1730 #[test]
1731 fn identify_sends_identify_event() {
1732 let (client, event_rx) = make_mocked_client();
1733 client.start_with_default_executor();
1734
1735 let context = ContextBuilder::new("bob")
1736 .build()
1737 .expect("Failed to create context");
1738
1739 client.identify(context);
1740 client.flush();
1741 client.close();
1742
1743 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1744 assert_eq!(events.len(), 1);
1745 assert_eq!(events[0].kind(), "identify");
1746 }
1747
1748 #[test]
1749 fn identify_sends_sends_nothing_in_offline_mode() {
1750 let (client, event_rx) = make_mocked_offline_client();
1751 client.start_with_default_executor();
1752
1753 let context = ContextBuilder::new("bob")
1754 .build()
1755 .expect("Failed to create context");
1756
1757 client.identify(context);
1758 client.flush();
1759 client.close();
1760
1761 assert_eq!(event_rx.iter().count(), 0);
1762 }
1763
1764 #[test]
1765 #[cfg(any(feature = "crypto-aws-lc-rs", feature = "crypto-openssl"))]
1766 fn secure_mode_hash() {
1767 let config = ConfigBuilder::new("secret")
1768 .offline(true)
1769 .build()
1770 .expect("config should build");
1771 let client = Client::build(config).expect("Should be built.");
1772 let context = ContextBuilder::new("Message")
1773 .build()
1774 .expect("Failed to create context");
1775
1776 assert_eq!(
1777 client
1778 .secure_mode_hash(&context)
1779 .expect("Hash should be computed"),
1780 "aa747c502a898200f9e4fa21bac68136f886a0e27aec70ba06daf2e2a5cb5597"
1781 );
1782 }
1783
1784 #[test]
1785 #[cfg(any(feature = "crypto-aws-lc-rs", feature = "crypto-openssl"))]
1786 fn secure_mode_hash_with_multi_kind() {
1787 let config = ConfigBuilder::new("secret")
1788 .offline(true)
1789 .build()
1790 .expect("config should build");
1791 let client = Client::build(config).expect("Should be built.");
1792
1793 let org = ContextBuilder::new("org-key|1")
1794 .kind("org")
1795 .build()
1796 .expect("Failed to create context");
1797 let user = ContextBuilder::new("user-key:2")
1798 .build()
1799 .expect("Failed to create context");
1800
1801 let context = MultiContextBuilder::new()
1802 .add_context(org)
1803 .add_context(user)
1804 .build()
1805 .expect("failed to build multi-context");
1806
1807 assert_eq!(
1808 client
1809 .secure_mode_hash(&context)
1810 .expect("Hash should be computed"),
1811 "5687e6383b920582ed50c2a96c98a115f1b6aad85a60579d761d9b8797415163"
1812 );
1813 }
1814
1815 #[derive(Serialize)]
1816 struct MyCustomData {
1817 pub answer: u32,
1818 }
1819
1820 #[test]
1821 fn track_sends_track_and_index_events() -> serde_json::Result<()> {
1822 let (client, event_rx) = make_mocked_client();
1823 client.start_with_default_executor();
1824
1825 let context = ContextBuilder::new("bob")
1826 .build()
1827 .expect("Failed to create context");
1828
1829 client.track_event(context.clone(), "event-with-null");
1830 client.track_data(context.clone(), "event-with-string", "string-data")?;
1831 client.track_data(context.clone(), "event-with-json", json!({"answer": 42}))?;
1832 client.track_data(
1833 context.clone(),
1834 "event-with-struct",
1835 MyCustomData { answer: 42 },
1836 )?;
1837 client.track_metric(context, "event-with-metric", 42.0, serde_json::Value::Null);
1838
1839 client.flush();
1840 client.close();
1841
1842 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
1843 assert_eq!(events.len(), 6);
1844
1845 let mut events_by_type: HashMap<&str, usize> = HashMap::new();
1846 for event in events {
1847 if let Some(count) = events_by_type.get_mut(event.kind()) {
1848 *count += 1;
1849 } else {
1850 events_by_type.insert(event.kind(), 1);
1851 }
1852 }
1853 assert!(matches!(events_by_type.get("index"), Some(1)));
1854 assert!(matches!(events_by_type.get("custom"), Some(5)));
1855
1856 Ok(())
1857 }
1858
1859 #[test]
1860 fn track_sends_nothing_in_offline_mode() -> serde_json::Result<()> {
1861 let (client, event_rx) = make_mocked_offline_client();
1862 client.start_with_default_executor();
1863
1864 let context = ContextBuilder::new("bob")
1865 .build()
1866 .expect("Failed to create context");
1867
1868 client.track_event(context.clone(), "event-with-null");
1869 client.track_data(context.clone(), "event-with-string", "string-data")?;
1870 client.track_data(context.clone(), "event-with-json", json!({"answer": 42}))?;
1871 client.track_data(
1872 context.clone(),
1873 "event-with-struct",
1874 MyCustomData { answer: 42 },
1875 )?;
1876 client.track_metric(context, "event-with-metric", 42.0, serde_json::Value::Null);
1877
1878 client.flush();
1879 client.close();
1880
1881 assert_eq!(event_rx.iter().count(), 0);
1882
1883 Ok(())
1884 }
1885
1886 #[test]
1887 fn migration_handles_flag_not_found() {
1888 let (client, _event_rx) = make_mocked_client();
1889 client.start_with_default_executor();
1890
1891 let context = ContextBuilder::new("bob")
1892 .build()
1893 .expect("Failed to create context");
1894
1895 let (stage, _tracker) =
1896 client.migration_variation(&context, "non-existent-flag-key", Stage::Off);
1897
1898 assert_eq!(stage, Stage::Off);
1899 }
1900
1901 #[tokio::test]
1902 async fn migration_uses_non_migration_flag() {
1903 let td = TestData::new();
1904 td.use_preconfigured_flag(basic_flag("boolean-flag"));
1905 let (client, _event_rx) = make_client_with_test_data(&td);
1906 client.start_with_default_executor();
1907
1908 let context = ContextBuilder::new("bob")
1909 .build()
1910 .expect("Failed to create context");
1911
1912 let (stage, _tracker) = client.migration_variation(&context, "boolean-flag", Stage::Off);
1913
1914 assert_eq!(stage, Stage::Off);
1915 }
1916
1917 #[test_case(Stage::Off)]
1918 #[test_case(Stage::DualWrite)]
1919 #[test_case(Stage::Shadow)]
1920 #[test_case(Stage::Live)]
1921 #[test_case(Stage::Rampdown)]
1922 #[test_case(Stage::Complete)]
1923 #[tokio::test]
1924 async fn migration_can_determine_correct_stage_from_flag(stage: Stage) {
1925 let td = TestData::new();
1926 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
1927 let (client, _event_rx) = make_client_with_test_data(&td);
1928 client.start_with_default_executor();
1929
1930 let context = ContextBuilder::new("bob")
1931 .build()
1932 .expect("Failed to create context");
1933
1934 let (evaluated_stage, _tracker) =
1935 client.migration_variation(&context, "stage-flag", Stage::Off);
1936
1937 assert_eq!(evaluated_stage, stage);
1938 }
1939
1940 #[tokio::test]
1941 async fn migration_tracks_invoked_correctly() {
1942 migration_tracks_invoked_correctly_driver(Stage::Off, Operation::Read, vec![Origin::Old])
1943 .await;
1944 migration_tracks_invoked_correctly_driver(
1945 Stage::DualWrite,
1946 Operation::Read,
1947 vec![Origin::Old],
1948 )
1949 .await;
1950 migration_tracks_invoked_correctly_driver(
1951 Stage::Shadow,
1952 Operation::Read,
1953 vec![Origin::Old, Origin::New],
1954 )
1955 .await;
1956 migration_tracks_invoked_correctly_driver(
1957 Stage::Live,
1958 Operation::Read,
1959 vec![Origin::Old, Origin::New],
1960 )
1961 .await;
1962 migration_tracks_invoked_correctly_driver(
1963 Stage::Rampdown,
1964 Operation::Read,
1965 vec![Origin::New],
1966 )
1967 .await;
1968 migration_tracks_invoked_correctly_driver(
1969 Stage::Complete,
1970 Operation::Read,
1971 vec![Origin::New],
1972 )
1973 .await;
1974 migration_tracks_invoked_correctly_driver(Stage::Off, Operation::Write, vec![Origin::Old])
1975 .await;
1976 migration_tracks_invoked_correctly_driver(
1977 Stage::DualWrite,
1978 Operation::Write,
1979 vec![Origin::Old, Origin::New],
1980 )
1981 .await;
1982 migration_tracks_invoked_correctly_driver(
1983 Stage::Shadow,
1984 Operation::Write,
1985 vec![Origin::Old, Origin::New],
1986 )
1987 .await;
1988 migration_tracks_invoked_correctly_driver(
1989 Stage::Live,
1990 Operation::Write,
1991 vec![Origin::Old, Origin::New],
1992 )
1993 .await;
1994 migration_tracks_invoked_correctly_driver(
1995 Stage::Rampdown,
1996 Operation::Write,
1997 vec![Origin::Old, Origin::New],
1998 )
1999 .await;
2000 migration_tracks_invoked_correctly_driver(
2001 Stage::Complete,
2002 Operation::Write,
2003 vec![Origin::New],
2004 )
2005 .await;
2006 }
2007
2008 async fn migration_tracks_invoked_correctly_driver(
2009 stage: Stage,
2010 operation: Operation,
2011 origins: Vec<Origin>,
2012 ) {
2013 let td = TestData::new();
2014 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
2015 let (client, event_rx) = make_client_with_test_data(&td);
2016 let client = Arc::new(client);
2017 client.start_with_default_executor();
2018
2019 let mut migrator = MigratorBuilder::new(client.clone())
2020 .read(
2021 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2022 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2023 Some(|_, _| true),
2024 )
2025 .write(
2026 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2027 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2028 )
2029 .build()
2030 .expect("migrator should build");
2031
2032 let context = ContextBuilder::new("bob")
2033 .build()
2034 .expect("Failed to create context");
2035
2036 if let Operation::Read = operation {
2037 migrator
2038 .read(
2039 &context,
2040 "stage-flag".into(),
2041 Stage::Off,
2042 serde_json::Value::Null,
2043 )
2044 .await;
2045 } else {
2046 migrator
2047 .write(
2048 &context,
2049 "stage-flag".into(),
2050 Stage::Off,
2051 serde_json::Value::Null,
2052 )
2053 .await;
2054 }
2055
2056 client.flush();
2057 client.close();
2058
2059 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2060 assert_eq!(events.len(), 3);
2061 match &events[1] {
2062 OutputEvent::MigrationOp(event) => {
2063 assert!(event.invoked.len() == origins.len());
2064 assert!(event.invoked.iter().all(|i| origins.contains(i)));
2065 }
2066 _ => panic!("Expected migration event"),
2067 }
2068 }
2069
2070 #[tokio::test]
2071 async fn migration_tracks_latency() {
2072 migration_tracks_latency_driver(Stage::Off, Operation::Read, vec![Origin::Old]).await;
2073 migration_tracks_latency_driver(Stage::DualWrite, Operation::Read, vec![Origin::Old]).await;
2074 migration_tracks_latency_driver(
2075 Stage::Shadow,
2076 Operation::Read,
2077 vec![Origin::Old, Origin::New],
2078 )
2079 .await;
2080 migration_tracks_latency_driver(
2081 Stage::Live,
2082 Operation::Read,
2083 vec![Origin::Old, Origin::New],
2084 )
2085 .await;
2086 migration_tracks_latency_driver(Stage::Rampdown, Operation::Read, vec![Origin::New]).await;
2087 migration_tracks_latency_driver(Stage::Complete, Operation::Read, vec![Origin::New]).await;
2088 migration_tracks_latency_driver(Stage::Off, Operation::Write, vec![Origin::Old]).await;
2089 migration_tracks_latency_driver(
2090 Stage::DualWrite,
2091 Operation::Write,
2092 vec![Origin::Old, Origin::New],
2093 )
2094 .await;
2095 migration_tracks_latency_driver(
2096 Stage::Shadow,
2097 Operation::Write,
2098 vec![Origin::Old, Origin::New],
2099 )
2100 .await;
2101 migration_tracks_latency_driver(
2102 Stage::Live,
2103 Operation::Write,
2104 vec![Origin::Old, Origin::New],
2105 )
2106 .await;
2107 migration_tracks_latency_driver(
2108 Stage::Rampdown,
2109 Operation::Write,
2110 vec![Origin::Old, Origin::New],
2111 )
2112 .await;
2113 migration_tracks_latency_driver(Stage::Complete, Operation::Write, vec![Origin::New]).await;
2114 }
2115
2116 async fn migration_tracks_latency_driver(
2117 stage: Stage,
2118 operation: Operation,
2119 origins: Vec<Origin>,
2120 ) {
2121 let td = TestData::new();
2122 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
2123 let (client, event_rx) = make_client_with_test_data(&td);
2124 let client = Arc::new(client);
2125 client.start_with_default_executor();
2126
2127 let mut migrator = MigratorBuilder::new(client.clone())
2128 .track_latency(true)
2129 .read(
2130 |_| {
2131 async move {
2132 async_std::task::sleep(Duration::from_millis(100)).await;
2133 Ok(serde_json::Value::Null)
2134 }
2135 .boxed()
2136 },
2137 |_| {
2138 async move {
2139 async_std::task::sleep(Duration::from_millis(100)).await;
2140 Ok(serde_json::Value::Null)
2141 }
2142 .boxed()
2143 },
2144 Some(|_, _| true),
2145 )
2146 .write(
2147 |_| {
2148 async move {
2149 async_std::task::sleep(Duration::from_millis(100)).await;
2150 Ok(serde_json::Value::Null)
2151 }
2152 .boxed()
2153 },
2154 |_| {
2155 async move {
2156 async_std::task::sleep(Duration::from_millis(100)).await;
2157 Ok(serde_json::Value::Null)
2158 }
2159 .boxed()
2160 },
2161 )
2162 .build()
2163 .expect("migrator should build");
2164
2165 let context = ContextBuilder::new("bob")
2166 .build()
2167 .expect("Failed to create context");
2168
2169 if let Operation::Read = operation {
2170 migrator
2171 .read(
2172 &context,
2173 "stage-flag".into(),
2174 Stage::Off,
2175 serde_json::Value::Null,
2176 )
2177 .await;
2178 } else {
2179 migrator
2180 .write(
2181 &context,
2182 "stage-flag".into(),
2183 Stage::Off,
2184 serde_json::Value::Null,
2185 )
2186 .await;
2187 }
2188
2189 client.flush();
2190 client.close();
2191
2192 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2193 assert_eq!(events.len(), 3);
2194 match &events[1] {
2195 OutputEvent::MigrationOp(event) => {
2196 assert!(event.latency.len() == origins.len());
2197 assert!(event
2198 .latency
2199 .values()
2200 .all(|l| l > &Duration::from_millis(100)));
2201 }
2202 _ => panic!("Expected migration event"),
2203 }
2204 }
2205
2206 #[tokio::test]
2207 async fn migration_tracks_read_errors() {
2208 migration_tracks_read_errors_driver(Stage::Off, vec![Origin::Old]).await;
2209 migration_tracks_read_errors_driver(Stage::DualWrite, vec![Origin::Old]).await;
2210 migration_tracks_read_errors_driver(Stage::Shadow, vec![Origin::Old, Origin::New]).await;
2211 migration_tracks_read_errors_driver(Stage::Live, vec![Origin::Old, Origin::New]).await;
2212 migration_tracks_read_errors_driver(Stage::Rampdown, vec![Origin::New]).await;
2213 migration_tracks_read_errors_driver(Stage::Complete, vec![Origin::New]).await;
2214 }
2215
2216 async fn migration_tracks_read_errors_driver(stage: Stage, origins: Vec<Origin>) {
2217 let td = TestData::new();
2218 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
2219 let (client, event_rx) = make_client_with_test_data(&td);
2220 let client = Arc::new(client);
2221 client.start_with_default_executor();
2222
2223 let mut migrator = MigratorBuilder::new(client.clone())
2224 .track_latency(true)
2225 .read(
2226 |_| async move { Err("fail".into()) }.boxed(),
2227 |_| async move { Err("fail".into()) }.boxed(),
2228 Some(|_: &String, _: &String| true),
2229 )
2230 .write(
2231 |_| async move { Err("fail".into()) }.boxed(),
2232 |_| async move { Err("fail".into()) }.boxed(),
2233 )
2234 .build()
2235 .expect("migrator should build");
2236
2237 let context = ContextBuilder::new("bob")
2238 .build()
2239 .expect("Failed to create context");
2240
2241 migrator
2242 .read(
2243 &context,
2244 "stage-flag".into(),
2245 Stage::Off,
2246 serde_json::Value::Null,
2247 )
2248 .await;
2249 client.flush();
2250 client.close();
2251
2252 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2253 assert_eq!(events.len(), 3);
2254 match &events[1] {
2255 OutputEvent::MigrationOp(event) => {
2256 assert!(event.errors.len() == origins.len());
2257 assert!(event.errors.iter().all(|i| origins.contains(i)));
2258 }
2259 _ => panic!("Expected migration event"),
2260 }
2261 }
2262
2263 #[tokio::test]
2264 async fn migration_tracks_authoritative_write_errors() {
2265 migration_tracks_authoritative_write_errors_driver(Stage::Off, vec![Origin::Old]).await;
2266 migration_tracks_authoritative_write_errors_driver(Stage::DualWrite, vec![Origin::Old])
2267 .await;
2268 migration_tracks_authoritative_write_errors_driver(Stage::Shadow, vec![Origin::Old]).await;
2269 migration_tracks_authoritative_write_errors_driver(Stage::Live, vec![Origin::New]).await;
2270 migration_tracks_authoritative_write_errors_driver(Stage::Rampdown, vec![Origin::New])
2271 .await;
2272 migration_tracks_authoritative_write_errors_driver(Stage::Complete, vec![Origin::New])
2273 .await;
2274 }
2275
2276 async fn migration_tracks_authoritative_write_errors_driver(
2277 stage: Stage,
2278 origins: Vec<Origin>,
2279 ) {
2280 let td = TestData::new();
2281 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
2282 let (client, event_rx) = make_client_with_test_data(&td);
2283 let client = Arc::new(client);
2284 client.start_with_default_executor();
2285
2286 let mut migrator = MigratorBuilder::new(client.clone())
2287 .track_latency(true)
2288 .read(
2289 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2290 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2291 None,
2292 )
2293 .write(
2294 |_| async move { Err("fail".into()) }.boxed(),
2295 |_| async move { Err("fail".into()) }.boxed(),
2296 )
2297 .build()
2298 .expect("migrator should build");
2299
2300 let context = ContextBuilder::new("bob")
2301 .build()
2302 .expect("Failed to create context");
2303
2304 migrator
2305 .write(
2306 &context,
2307 "stage-flag".into(),
2308 Stage::Off,
2309 serde_json::Value::Null,
2310 )
2311 .await;
2312
2313 client.flush();
2314 client.close();
2315
2316 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2317 assert_eq!(events.len(), 3);
2318 match &events[1] {
2319 OutputEvent::MigrationOp(event) => {
2320 assert!(event.errors.len() == origins.len());
2321 assert!(event.errors.iter().all(|i| origins.contains(i)));
2322 }
2323 _ => panic!("Expected migration event"),
2324 }
2325 }
2326
2327 #[tokio::test]
2328 async fn migration_tracks_nonauthoritative_write_errors() {
2329 migration_tracks_nonauthoritative_write_errors_driver(
2330 Stage::DualWrite,
2331 false,
2332 true,
2333 vec![Origin::New],
2334 )
2335 .await;
2336 migration_tracks_nonauthoritative_write_errors_driver(
2337 Stage::Shadow,
2338 false,
2339 true,
2340 vec![Origin::New],
2341 )
2342 .await;
2343 migration_tracks_nonauthoritative_write_errors_driver(
2344 Stage::Live,
2345 true,
2346 false,
2347 vec![Origin::Old],
2348 )
2349 .await;
2350 migration_tracks_nonauthoritative_write_errors_driver(
2351 Stage::Rampdown,
2352 true,
2353 false,
2354 vec![Origin::Old],
2355 )
2356 .await;
2357 }
2358
2359 async fn migration_tracks_nonauthoritative_write_errors_driver(
2360 stage: Stage,
2361 fail_old: bool,
2362 fail_new: bool,
2363 origins: Vec<Origin>,
2364 ) {
2365 let td = TestData::new();
2366 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
2367 let (client, event_rx) = make_client_with_test_data(&td);
2368 let client = Arc::new(client);
2369 client.start_with_default_executor();
2370
2371 let mut migrator = MigratorBuilder::new(client.clone())
2372 .track_latency(true)
2373 .read(
2374 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2375 |_| async move { Ok(serde_json::Value::Null) }.boxed(),
2376 None,
2377 )
2378 .write(
2379 move |_| {
2380 async move {
2381 if fail_old {
2382 Err("fail".into())
2383 } else {
2384 Ok(serde_json::Value::Null)
2385 }
2386 }
2387 .boxed()
2388 },
2389 move |_| {
2390 async move {
2391 if fail_new {
2392 Err("fail".into())
2393 } else {
2394 Ok(serde_json::Value::Null)
2395 }
2396 }
2397 .boxed()
2398 },
2399 )
2400 .build()
2401 .expect("migrator should build");
2402
2403 let context = ContextBuilder::new("bob")
2404 .build()
2405 .expect("Failed to create context");
2406
2407 migrator
2408 .write(
2409 &context,
2410 "stage-flag".into(),
2411 Stage::Off,
2412 serde_json::Value::Null,
2413 )
2414 .await;
2415
2416 client.flush();
2417 client.close();
2418
2419 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2420 assert_eq!(events.len(), 3);
2421 match &events[1] {
2422 OutputEvent::MigrationOp(event) => {
2423 assert!(event.errors.len() == origins.len());
2424 assert!(event.errors.iter().all(|i| origins.contains(i)));
2425 }
2426 _ => panic!("Expected migration event"),
2427 }
2428 }
2429
2430 #[tokio::test]
2431 async fn migration_tracks_consistency() {
2432 migration_tracks_consistency_driver(Stage::Shadow, "same", "same", true).await;
2433 migration_tracks_consistency_driver(Stage::Shadow, "same", "different", false).await;
2434 migration_tracks_consistency_driver(Stage::Live, "same", "same", true).await;
2435 migration_tracks_consistency_driver(Stage::Live, "same", "different", false).await;
2436 }
2437
2438 async fn migration_tracks_consistency_driver(
2439 stage: Stage,
2440 old_return: &'static str,
2441 new_return: &'static str,
2442 expected_consistency: bool,
2443 ) {
2444 let td = TestData::new();
2445 td.use_preconfigured_flag(basic_migration_flag("stage-flag", stage));
2446 let (client, event_rx) = make_client_with_test_data(&td);
2447 let client = Arc::new(client);
2448 client.start_with_default_executor();
2449
2450 let mut migrator = MigratorBuilder::new(client.clone())
2451 .track_latency(true)
2452 .read(
2453 |_| {
2454 async move {
2455 async_std::task::sleep(Duration::from_millis(100)).await;
2456 Ok(serde_json::Value::String(old_return.to_string()))
2457 }
2458 .boxed()
2459 },
2460 |_| {
2461 async move {
2462 async_std::task::sleep(Duration::from_millis(100)).await;
2463 Ok(serde_json::Value::String(new_return.to_string()))
2464 }
2465 .boxed()
2466 },
2467 Some(|lhs, rhs| lhs == rhs),
2468 )
2469 .write(
2470 |_| {
2471 async move {
2472 async_std::task::sleep(Duration::from_millis(100)).await;
2473 Ok(serde_json::Value::Null)
2474 }
2475 .boxed()
2476 },
2477 |_| {
2478 async move {
2479 async_std::task::sleep(Duration::from_millis(100)).await;
2480 Ok(serde_json::Value::Null)
2481 }
2482 .boxed()
2483 },
2484 )
2485 .build()
2486 .expect("migrator should build");
2487
2488 let context = ContextBuilder::new("bob")
2489 .build()
2490 .expect("Failed to create context");
2491
2492 migrator
2493 .read(
2494 &context,
2495 "stage-flag".into(),
2496 Stage::Off,
2497 serde_json::Value::Null,
2498 )
2499 .await;
2500
2501 client.flush();
2502 client.close();
2503
2504 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2505 assert_eq!(events.len(), 3);
2506 match &events[1] {
2507 OutputEvent::MigrationOp(event) => {
2508 assert!(event.consistency_check == Some(expected_consistency))
2509 }
2510 _ => panic!("Expected migration event"),
2511 }
2512 }
2513
2514 #[tokio::test]
2515 async fn client_flush_blocking_completes_successfully() {
2516 let (client, event_rx) = make_mocked_client();
2517 client.start_with_default_executor();
2518 client.wait_for_initialization(Duration::from_secs(1)).await;
2519
2520 let context = ContextBuilder::new("user-key")
2521 .build()
2522 .expect("Failed to create context");
2523
2524 client.identify(context);
2525
2526 let result = client.flush_blocking(Duration::from_secs(5)).await;
2527 assert!(result, "flush_blocking should complete successfully");
2528
2529 client.close();
2530
2531 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2532 assert!(!events.is_empty(), "Should have received identify event");
2533 }
2534
2535 #[tokio::test]
2536 async fn client_flush_blocking_with_zero_timeout() {
2537 let (client, event_rx) = make_mocked_client();
2538 client.start_with_default_executor();
2539 client.wait_for_initialization(Duration::from_secs(1)).await;
2540
2541 let context = ContextBuilder::new("user-key")
2542 .build()
2543 .expect("Failed to create context");
2544
2545 client.identify(context);
2546
2547 let result = client.flush_blocking(Duration::ZERO).await;
2548 assert!(
2549 result,
2550 "flush_blocking with zero timeout should complete successfully"
2551 );
2552
2553 client.close();
2554
2555 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2556 assert!(!events.is_empty(), "Should have received identify event");
2557 }
2558
2559 #[tokio::test]
2560 async fn client_flush_blocking_with_no_events() {
2561 let (client, _event_rx) = make_mocked_client();
2562 client.start_with_default_executor();
2563 client.wait_for_initialization(Duration::from_secs(1)).await;
2564
2565 let result = client.flush_blocking(Duration::from_secs(1)).await;
2566 assert!(
2567 result,
2568 "flush_blocking with no events should complete immediately"
2569 );
2570
2571 client.close();
2572 }
2573
2574 #[tokio::test]
2575 async fn client_flush_blocking_multiple_concurrent_calls() {
2576 let (client, event_rx) = make_mocked_client();
2577 client.start_with_default_executor();
2578 client.wait_for_initialization(Duration::from_secs(1)).await;
2579
2580 let context = ContextBuilder::new("user-key")
2581 .build()
2582 .expect("Failed to create context");
2583
2584 client.identify(context);
2585
2586 let (result1, result2, result3) = tokio::join!(
2588 client.flush_blocking(Duration::from_secs(5)),
2589 client.flush_blocking(Duration::from_secs(5)),
2590 client.flush_blocking(Duration::from_secs(5)),
2591 );
2592
2593 assert!(result1, "First flush_blocking should succeed");
2594 assert!(result2, "Second flush_blocking should succeed");
2595 assert!(result3, "Third flush_blocking should succeed");
2596
2597 client.close();
2598
2599 let events = event_rx.iter().collect::<Vec<OutputEvent>>();
2600 assert!(!events.is_empty(), "Should have received identify event");
2601 }
2602
2603 fn make_mocked_client_with_delay(
2604 delay: u64,
2605 offline: bool,
2606 daemon_mode: bool,
2607 ) -> (Client, Receiver<OutputEvent>) {
2608 let updates = Arc::new(MockDataSource::new_with_init_delay(delay));
2609 let (event_sender, event_rx) = create_event_sender();
2610
2611 let config = ConfigBuilder::new("sdk-key")
2612 .offline(offline)
2613 .daemon_mode(daemon_mode)
2614 .data_source(MockDataSourceBuilder::new().data_source(updates))
2615 .event_processor(
2616 EventProcessorBuilder::<launchdarkly_sdk_transport::HyperTransport>::new()
2617 .event_sender(Arc::new(event_sender)),
2618 )
2619 .build()
2620 .expect("config should build");
2621
2622 let client = Client::build(config).expect("Should be built.");
2623
2624 (client, event_rx)
2625 }
2626
2627 fn make_mocked_offline_client() -> (Client, Receiver<OutputEvent>) {
2628 make_mocked_client_with_delay(0, true, false)
2629 }
2630
2631 fn make_mocked_client() -> (Client, Receiver<OutputEvent>) {
2632 make_mocked_client_with_delay(0, false, false)
2633 }
2634
2635 fn make_client_with_test_data(td: &TestData) -> (Client, Receiver<OutputEvent>) {
2636 let (event_sender, event_rx) = create_event_sender();
2637 let config = ConfigBuilder::new("sdk-key")
2638 .data_source(td)
2639 .event_processor(
2640 EventProcessorBuilder::<launchdarkly_sdk_transport::HyperTransport>::new()
2641 .event_sender(Arc::new(event_sender)),
2642 )
2643 .build()
2644 .expect("config should build");
2645 let client = Client::build(config).expect("Should be built.");
2646 (client, event_rx)
2647 }
2648
2649 #[test]
2650 fn client_builds_successfully() {
2651 let config = ConfigBuilder::new("sdk-key")
2652 .offline(true)
2653 .build()
2654 .expect("config should build");
2655
2656 let client = Client::build(config).expect("client should build successfully");
2657
2658 assert!(
2659 !client.started.load(Ordering::SeqCst),
2660 "client should not be started yet"
2661 );
2662 assert!(client.offline, "client should be in offline mode");
2663 assert_eq!(client.sdk_key, "sdk-key", "sdk_key should match");
2664 }
2665}