Skip to main content

launchdarkly_server_sdk/
client.rs

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/// Error type used to represent failures when building a [Client] instance.
54#[non_exhaustive]
55#[derive(Debug, Error)]
56pub enum BuildError {
57    /// Error used when a configuration setting is invalid. This typically indicates an invalid URL.
58    #[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/// Error type used to represent failures when starting the [Client].
87#[non_exhaustive]
88#[derive(Debug, Error)]
89pub enum StartError {
90    /// Error used when spawning a background there fails.
91    #[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
119/// A client for the LaunchDarkly API.
120///
121/// In order to create a client instance, first create a config using [crate::ConfigBuilder].
122///
123/// # Examples
124///
125/// Creating a client, with default configuration.
126/// ```
127/// # use launchdarkly_server_sdk::{Client, ConfigBuilder, BuildError};
128/// # fn main() -> Result<(), BuildError> {
129///     let ld_client = Client::build(ConfigBuilder::new("sdk-key").build()?)?;
130/// #   Ok(())
131/// # }
132/// ```
133///
134/// Creating an instance which connects to a relay proxy.
135/// ```
136/// # use launchdarkly_server_sdk::{Client, ConfigBuilder, ServiceEndpointsBuilder, BuildError};
137/// # fn main() -> Result<(), BuildError> {
138///     let ld_client = Client::build(ConfigBuilder::new("sdk-key")
139///         .service_endpoints(ServiceEndpointsBuilder::new()
140///             .relay_proxy("http://my-relay-hostname:8080")
141///         ).build()?
142///     )?;
143/// #   Ok(())
144/// # }
145/// ```
146///
147/// Each builder type includes usage examples for the builder.
148pub 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    /// Create a new instance of a [Client] based on the provided [Config] parameter.
170    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    /// Starts a client in the current thread, which must have a default tokio runtime.
234    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        // These clones are going to move into the closure, we
244        // do not want to move or reference `self`, because
245        // then lifetimes will get involved.
246        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    /// Creates a new tokio runtime and then starts the client. Tasks from the client will
266    /// be executed on created runtime.
267    /// If your application already has a tokio runtime, then you can use
268    /// [crate::Client::start_with_default_executor] and the client will dispatch tasks to
269    /// your existing runtime.
270    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    /// This is an async method that will resolve once initialization is complete or the specified
286    /// timeout has occurred.
287    ///
288    /// If the timeout is triggered, this method will return `None`. Otherwise, the method will
289    /// return a boolean indicating whether or not the SDK has successfully initialized.
290    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 the client is not initialized, then we need to wait for it to be initialized.
305        // Because we are using atomic types, and not a lock, then there is still the possibility
306        // that the value will change between the read and when we wait. We use a semaphore to wait,
307        // and we do not forget the permit, therefore if the permit has been added, then we will get
308        // it very quickly and reduce blocking.
309        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    /// This function synchronously returns if the SDK is initialized.
316    /// In the case of unrecoverable errors in establishing a connection it is possible for the
317    /// SDK to never become initialized.
318    pub fn initialized(&self) -> bool {
319        self.offline
320            || self.daemon_mode
321            || ClientInitState::Initialized == self.init_state.load(Ordering::SeqCst)
322    }
323
324    /// Close shuts down the LaunchDarkly client. After calling this, the LaunchDarkly client
325    /// should no longer be used. The method will block until all pending analytics events (if any)
326    /// been sent.
327    pub fn close(&self) {
328        self.event_processor.close();
329
330        // If the system is in offline mode or daemon mode, no receiver will be listening to this
331        // broadcast channel, so sending on it would always result in an error.
332        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        // Potentially take the runtime we created when starting the client and do nothing with it
339        // so it drops, closing out all spawned tasks.
340        self.runtime.write().take();
341    }
342
343    /// Flush tells the client that all pending analytics events (if any) should be delivered as
344    /// soon as possible. Flushing is asynchronous, so this method will return before it is
345    /// complete. However, if you call [Client::close], events are guaranteed to be sent before
346    /// that method returns.
347    ///
348    /// For more information, see the Reference Guide:
349    /// <https://docs.launchdarkly.com/sdk/features/flush#rust>.
350    pub fn flush(&self) {
351        self.event_processor.flush();
352    }
353
354    /// Flush tells the client that all pending analytics events should be delivered as
355    /// soon as possible, and blocks until delivery is complete or the timeout expires.
356    ///
357    /// This method is particularly useful in short-lived execution environments like AWS Lambda
358    /// where you need to ensure events are sent before the function terminates.
359    ///
360    /// This method triggers a flush of events currently buffered and waits for that specific
361    /// flush to complete. Note that if periodic flushes or other flush operations are in-flight
362    /// when this is called, those may still be completing after this method returns.
363    ///
364    /// # Arguments
365    ///
366    /// * `timeout` - Maximum time to wait for flush to complete. Use `Duration::ZERO` to wait indefinitely.
367    ///
368    /// # Returns
369    ///
370    /// Returns `true` if flush completed successfully, `false` if timeout occurred.
371    ///
372    /// # Examples
373    ///
374    /// ```no_run
375    /// # use launchdarkly_server_sdk::{Client, ConfigBuilder};
376    /// # use std::time::Duration;
377    /// # async fn example() {
378    /// # let client = Client::build(ConfigBuilder::new("sdk-key").build().unwrap()).unwrap();
379    /// // Wait up to 5 seconds for flush to complete
380    /// let success = client.flush_blocking(Duration::from_secs(5)).await;
381    /// if !success {
382    ///     eprintln!("Warning: flush timed out");
383    /// }
384    /// # }
385    /// ```
386    ///
387    /// For more information, see the Reference Guide:
388    /// <https://docs.launchdarkly.com/sdk/features/flush#rust>.
389    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            // Wait indefinitely
397            flush_future.await.unwrap_or(false)
398        } else {
399            // Apply timeout at async level too
400            match tokio::time::timeout(timeout, flush_future).await {
401                Ok(Ok(result)) => result,
402                Ok(Err(_)) => false, // spawn_blocking panicked
403                Err(_) => false,     // Timeout
404            }
405        }
406    }
407
408    /// Identify reports details about a context.
409    ///
410    /// For more information, see the Reference Guide:
411    /// <https://docs.launchdarkly.com/sdk/features/identify#rust>
412    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    /// Returns the value of a boolean feature flag for a given context.
421    ///
422    /// Returns `default` if there is an error, if the flag doesn't exist, or the feature is turned
423    /// off and has no off variation.
424    ///
425    /// For more information, see the Reference Guide:
426    /// <https://docs.launchdarkly.com/sdk/features/evaluating#rust>.
427    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    /// Returns the value of a string feature flag for a given context.
438    ///
439    /// Returns `default` if there is an error, if the flag doesn't exist, or the feature is turned
440    /// off and has no off variation.
441    ///
442    /// For more information, see the Reference Guide:
443    /// <https://docs.launchdarkly.com/sdk/features/evaluating#rust>.
444    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    /// Returns the value of a float feature flag for a given context.
455    ///
456    /// Returns `default` if there is an error, if the flag doesn't exist, or the feature is turned
457    /// off and has no off variation.
458    ///
459    /// For more information, see the Reference Guide:
460    /// <https://docs.launchdarkly.com/sdk/features/evaluating#rust>.
461    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    /// Returns the value of a integer feature flag for a given context.
472    ///
473    /// Returns `default` if there is an error, if the flag doesn't exist, or the feature is turned
474    /// off and has no off variation.
475    ///
476    /// For more information, see the Reference Guide:
477    /// <https://docs.launchdarkly.com/sdk/features/evaluating#rust>.
478    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    /// Returns the value of a feature flag for the given context, allowing the value to be
489    /// of any JSON type.
490    ///
491    /// The value is returned as an [serde_json::Value].
492    ///
493    /// Returns `default` if there is an error, if the flag doesn't exist, or the feature is turned off.
494    ///
495    /// For more information, see the Reference Guide:
496    /// <https://docs.launchdarkly.com/sdk/features/evaluating#rust>.
497    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    /// This method is the same as [Client::bool_variation], but also returns further information
509    /// about how the value was calculated. The "reason" data will also be included in analytics
510    /// events.
511    ///
512    /// For more information, see the Reference Guide:
513    /// <https://docs.launchdarkly.com/sdk/features/evaluation-reasons#rust>.
514    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    /// This method is the same as [Client::str_variation], but also returns further information
528    /// about how the value was calculated. The "reason" data will also be included in analytics
529    /// events.
530    ///
531    /// For more information, see the Reference Guide:
532    /// <https://docs.launchdarkly.com/sdk/features/evaluation-reasons#rust>.
533    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    /// This method is the same as [Client::float_variation], but also returns further information
544    /// about how the value was calculated. The "reason" data will also be included in analytics
545    /// events.
546    ///
547    /// For more information, see the Reference Guide:
548    /// <https://docs.launchdarkly.com/sdk/features/evaluation-reasons#rust>.
549    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    /// This method is the same as [Client::int_variation], but also returns further information
563    /// about how the value was calculated. The "reason" data will also be included in analytics
564    /// events.
565    ///
566    /// For more information, see the Reference Guide:
567    /// <https://docs.launchdarkly.com/sdk/features/evaluation-reasons#rust>.
568    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    /// This method is the same as [Client::json_variation], but also returns further information
582    /// about how the value was calculated. The "reason" data will also be included in analytics
583    /// events.
584    ///
585    /// For more information, see the Reference Guide:
586    /// <https://docs.launchdarkly.com/sdk/features/evaluation-reasons#rust>.
587    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    /// Generates the secure mode hash value for a context.
599    ///
600    /// For more information, see the Reference Guide:
601    /// <https://docs.launchdarkly.com/sdk/features/secure-mode#rust>.
602    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    /// Returns an object that encapsulates the state of all feature flags for a given context. This
633    /// includes the flag values, and also metadata that can be used on the front end.
634    ///
635    /// The most common use case for this method is to bootstrap a set of client-side feature flags
636    /// from a back-end service.
637    ///
638    /// You may pass any configuration of [FlagDetailConfig] to control what data is included.
639    ///
640    /// For more information, see the Reference Guide:
641    /// <https://docs.launchdarkly.com/sdk/features/all-flags#rust>
642    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    /// This method is the same as [Client::variation], but also returns further information about
668    /// how the value was calculated. The "reason" data will also be included in analytics events.
669    ///
670    /// For more information, see the Reference Guide:
671    /// <https://docs.launchdarkly.com/sdk/features/evaluation-reasons#rust>.
672    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    /// This is a generic function which returns the value of a feature flag for a given context.
684    ///
685    /// This method is an alternatively to the type specified methods (e.g.
686    /// [Client::bool_variation], [Client::int_variation], etc.).
687    ///
688    /// Returns `default` if there is an error, if the flag doesn't exist, or the feature is turned
689    /// off and has no off variation.
690    ///
691    /// For more information, see the Reference Guide:
692    /// <https://docs.launchdarkly.com/sdk/features/evaluating#rust>.
693    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    /// This method returns the migration stage of the migration feature flag for the given
704    /// evaluation context.
705    ///
706    /// This method returns the default stage if there is an error or the flag does not exist.
707    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    /// Reports that a context has performed an event.
733    ///
734    /// The `key` parameter is defined by the application and will be shown in analytics reports;
735    /// it normally corresponds to the event name of a metric that you have created through the
736    /// LaunchDarkly dashboard. If you want to associate additional data with this event, use
737    /// [Client::track_data] or [Client::track_metric].
738    ///
739    /// For more information, see the Reference Guide:
740    /// <https://docs.launchdarkly.com/sdk/features/events#rust>.
741    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    /// Reports that a context has performed an event, and associates it with custom data.
746    ///
747    /// The `key` parameter is defined by the application and will be shown in analytics reports;
748    /// it normally corresponds to the event name of a metric that you have created through the
749    /// LaunchDarkly dashboard.
750    ///
751    /// `data` parameter is any type that implements [Serialize]. If no such value is needed, use
752    /// [serde_json::Value::Null] (or call [Client::track_event] instead). To send a numeric value
753    /// for experimentation, use [Client::track_metric].
754    ///
755    /// For more information, see the Reference Guide:
756    /// <https://docs.launchdarkly.com/sdk/features/events#rust>.
757    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    /// Reports that a context has performed an event, and associates it with a numeric value. This
767    /// value is used by the LaunchDarkly experimentation feature in numeric custom metrics, and
768    /// will also be returned as part of the custom event for Data Export.
769    ///
770    /// The `key` parameter is defined by the application and will be shown in analytics reports;
771    /// it normally corresponds to the event name of a metric that you have created through the
772    /// LaunchDarkly dashboard.
773    ///
774    /// For more information, see the Reference Guide:
775    /// <https://docs.launchdarkly.com/sdk/features/events#rust>.
776    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    /// Tracks the results of a migrations operation. This event includes measurements which can be
806    /// used to enhance the observability of a migration within the LaunchDarkly UI.
807    ///
808    /// This event should be generated through [crate::MigrationOpTracker]. If you are using the
809    /// [crate::Migrator] to handle migrations, this event will be created and emitted
810    /// automatically.
811    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        // Give ourself a good margin for thread scheduling.
953        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        // Give ourself a good margin for thread scheduling.
968        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); // Jan. 1st, 4000
1284        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        // Make multiple concurrent flush_blocking calls
2587        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}