nimbus/stateful/
nimbus_client.rs

1/* This Source Code Form is subject to the terms of the Mozilla Public
2 * License, v. 2.0. If a copy of the MPL was not distributed with this
3 * file, You can obtain one at https://mozilla.org/MPL/2.0/. */
4
5use std::collections::HashSet;
6use std::fmt::Debug;
7use std::path::{Path, PathBuf};
8use std::sync::{Arc, Mutex, MutexGuard};
9
10use chrono::{DateTime, NaiveDateTime, Utc};
11use once_cell::sync::OnceCell;
12use remote_settings::RemoteSettingsService;
13use serde_json::Value;
14use uuid::Uuid;
15
16use crate::defaults::Defaults;
17use crate::enrollment::{
18    EnrolledFeature, EnrollmentChangeEvent, EnrollmentChangeEventType, EnrollmentsEvolver,
19    ExperimentEnrollment, PreviousGeckoPrefState,
20};
21use crate::error::{BehaviorError, info, warn};
22use crate::evaluator::{
23    CalculatedAttributes, ExperimentAvailable, TargetingAttributes, get_calculated_attributes,
24    is_experiment_available,
25};
26use crate::json::{JsonObject, PrefValue};
27use crate::metrics::{
28    DatabaseLoadExtraDef, DatabaseMigrationExtraDef, EnrollmentStatusExtraDef,
29    FeatureExposureExtraDef, MalformedFeatureConfigExtraDef, MetricsHandler,
30};
31use crate::schema::parse_experiments;
32use crate::stateful::behavior::EventStore;
33use crate::stateful::client::{NimbusServerSettings, SettingsClient, create_client};
34use crate::stateful::dbcache::DatabaseCache;
35use crate::stateful::enrollment::{
36    enroll_in_firefox_lab, get_experiment_participation, get_rollout_participation,
37    opt_in_with_branch, opt_out, reset_telemetry_identifiers, set_experiment_participation,
38    set_rollout_participation, unenroll_for_pref, unenroll_from_all_firefox_labs,
39    unenroll_from_firefox_lab,
40};
41use crate::stateful::firefox_labs::{
42    FirefoxLabsEnrollResult, FirefoxLabsEnrollStatus, FirefoxLabsMetadata,
43    FirefoxLabsUnenrollResult, FirefoxLabsUnenrollStatus,
44};
45use crate::stateful::gecko_prefs::{
46    GeckoPref, GeckoPrefHandler, GeckoPrefState, GeckoPrefStore, OriginalGeckoPref, PrefBranch,
47    PrefEnrollmentData, PrefUnenrollReason,
48};
49use crate::stateful::matcher::AppContext;
50use crate::stateful::persistence::{Database, StoreId, Writer};
51use crate::stateful::targeting::{RecordedContext, execute_event_queries, validate_event_queries};
52use crate::stateful::updating::{read_and_remove_pending_experiments, write_pending_experiments};
53use crate::strings::fmt_with_map;
54use crate::{
55    AvailableExperiment, AvailableRandomizationUnits, EnrolledExperiment, EnrollmentSlugs,
56    EnrollmentStatus,
57};
58use crate::{Experiment, ExperimentBranch, NimbusError, NimbusTargetingHelper, Result};
59
60const DB_KEY_NIMBUS_ID: &str = "nimbus-id";
61pub const DB_KEY_INSTALLATION_DATE: &str = "installation-date";
62pub const DB_KEY_UPDATE_DATE: &str = "update-date";
63pub const DB_KEY_APP_VERSION: &str = "app-version";
64pub const DB_KEY_FETCH_ENABLED: &str = "fetch-enabled";
65
66// The main `NimbusClient` struct must not expose any methods that make an `&mut self`,
67// in order to be compatible with the uniffi's requirements on objects. This is a helper
68// struct to contain the bits that do actually need to be mutable, so they can be
69// protected by a Mutex.
70#[derive(Default)]
71pub struct InternalMutableState {
72    pub(crate) available_randomization_units: AvailableRandomizationUnits,
73    pub(crate) install_date: Option<DateTime<Utc>>,
74    pub(crate) update_date: Option<DateTime<Utc>>,
75    // Application level targeting attributes
76    pub(crate) targeting_attributes: TargetingAttributes,
77}
78
79impl InternalMutableState {
80    pub(crate) fn update_time_to_now(&mut self, now: DateTime<Utc>) {
81        self.targeting_attributes
82            .update_time_to_now(now, &self.install_date, &self.update_date);
83    }
84}
85
86/// Nimbus is the main struct representing the experiments state
87/// It should hold all the information needed to communicate a specific user's
88/// experimentation status
89pub struct NimbusClient {
90    settings_client: Mutex<Box<dyn SettingsClient + Send>>,
91    pub(crate) mutable_state: Mutex<InternalMutableState>,
92    app_context: AppContext,
93    pub(crate) db: OnceCell<Database>,
94    // Manages an in-memory cache so that we can answer certain requests
95    // without doing (or waiting for) IO.
96    database_cache: DatabaseCache,
97    db_path: PathBuf,
98    coenrolling_feature_ids: Vec<String>,
99    event_store: Arc<Mutex<EventStore>>,
100    recorded_context: Option<Arc<dyn RecordedContext>>,
101    pub(crate) gecko_prefs: Option<Arc<GeckoPrefStore>>,
102    metrics_handler: Arc<dyn MetricsHandler>,
103}
104
105impl NimbusClient {
106    // This constructor *must* not do any kind of I/O since it might be called on the main
107    // thread in the gecko Javascript stack, hence the use of OnceCell for the db.
108    #[allow(clippy::too_many_arguments)]
109    pub fn new<P: Into<PathBuf>>(
110        app_context: AppContext,
111        recorded_context: Option<Arc<dyn RecordedContext>>,
112        coenrolling_feature_ids: Vec<String>,
113        db_path: P,
114        metrics_handler: Arc<dyn MetricsHandler>,
115        gecko_pref_handler: Option<Arc<dyn GeckoPrefHandler>>,
116        remote_settings_info: Option<NimbusServerSettings>,
117    ) -> Result<Self> {
118        let settings_client = Mutex::new(create_client(remote_settings_info)?);
119
120        let targeting_attributes: TargetingAttributes = app_context.clone().into();
121        let mutable_state = Mutex::new(InternalMutableState {
122            available_randomization_units: Default::default(),
123            targeting_attributes,
124            install_date: Default::default(),
125            update_date: Default::default(),
126        });
127
128        let mut prefs = None;
129        if let Some(handler) = gecko_pref_handler {
130            prefs = Some(Arc::new(GeckoPrefStore::new(handler)));
131        }
132
133        info!(
134            "Initialized NimbusClient with: app_context = {:?}; recorded_context = {:?}",
135            app_context,
136            recorded_context
137                .as_ref()
138                .map(|rc| serde_json::Value::Object(rc.to_json()))
139                .unwrap_or(serde_json::Value::Null)
140        );
141
142        Ok(Self {
143            settings_client,
144            mutable_state,
145            app_context,
146            database_cache: Default::default(),
147            db_path: db_path.into(),
148            coenrolling_feature_ids,
149            db: OnceCell::default(),
150            event_store: Arc::default(),
151            recorded_context,
152            gecko_prefs: prefs,
153            metrics_handler,
154        })
155    }
156
157    pub fn with_targeting_attributes(&mut self, targeting_attributes: TargetingAttributes) {
158        let mut state = self.mutable_state.lock().unwrap();
159        state.targeting_attributes = targeting_attributes;
160    }
161
162    pub fn get_targeting_attributes(&self) -> TargetingAttributes {
163        let mut state = self.mutable_state.lock().unwrap();
164        state.update_time_to_now(Utc::now());
165        state.targeting_attributes.clone()
166    }
167
168    pub fn initialize(&self) -> Result<()> {
169        let db = self.db()?;
170        // We're not actually going to write, we just want to exclude concurrent writers.
171        let mut writer = db.write()?;
172
173        let mut state = self.mutable_state.lock().unwrap();
174        self.begin_initialize(db, &mut writer, &mut state)?;
175        self.end_initialize(db, writer, &mut state)?;
176
177        Ok(())
178    }
179
180    // These are tasks which should be in the initialize and apply_pending_experiments
181    // but should happen before the enrollment calculations are done.
182    fn begin_initialize(
183        &self,
184        db: &Database,
185        writer: &mut Writer,
186        state: &mut MutexGuard<InternalMutableState>,
187    ) -> Result<()> {
188        self.read_or_create_nimbus_id(db, writer, state)?;
189        self.update_ta_install_dates(db, writer, state)?;
190        self.event_store
191            .lock()
192            .expect("unable to lock event_store mutex")
193            .read_from_db(db)?;
194
195        if let Some(recorded_context) = &self.recorded_context {
196            let targeting_helper = self.create_targeting_helper_with_context(match serde_json::to_value(
197                &state.targeting_attributes,
198            ) {
199                Ok(v) => v,
200                Err(e) => return Err(NimbusError::JSONError("targeting_helper = nimbus::stateful::nimbus_client::NimbusClient::begin_initialize::serde_json::to_value".into(), e.to_string()))
201            });
202            execute_event_queries(&**recorded_context, targeting_helper.as_ref())?;
203            state
204                .targeting_attributes
205                .set_recorded_context(recorded_context.to_json());
206        }
207
208        if let Some(gecko_prefs) = &self.gecko_prefs {
209            gecko_prefs.initialize()?;
210        }
211
212        Ok(())
213    }
214
215    // These are tasks which should be in the initialize and apply_pending_experiments
216    // but should happen after the enrollment calculations are done.
217    fn end_initialize(
218        &self,
219        db: &Database,
220        writer: Writer,
221        state: &mut MutexGuard<InternalMutableState>,
222    ) -> Result<()> {
223        self.update_ta_active_experiments(db, &writer, state)?;
224        let coenrolling_ids = self
225            .coenrolling_feature_ids
226            .iter()
227            .map(|s| s.as_str())
228            .collect();
229        self.database_cache.commit_and_update(
230            db,
231            writer,
232            &coenrolling_ids,
233            self.gecko_prefs.clone(),
234            true,
235        )?;
236        Ok(())
237    }
238
239    pub fn get_enrollment_by_feature(&self, feature_id: String) -> Result<Option<EnrolledFeature>> {
240        self.database_cache.get_enrollment_by_feature(&feature_id)
241    }
242
243    // Note: the contract for this function is that it never blocks on IO.
244    pub fn get_experiment_branch(&self, slug: String) -> Result<Option<String>> {
245        self.database_cache.get_experiment_branch(&slug)
246    }
247
248    pub fn get_feature_config_variables(&self, feature_id: String) -> Result<Option<String>> {
249        Ok(
250            if let Some(s) = self
251                .database_cache
252                .get_feature_config_variables(&feature_id)?
253            {
254                self.record_feature_activation_if_needed(&feature_id);
255                Some(s)
256            } else {
257                None
258            },
259        )
260    }
261
262    pub fn get_experiment_branches(&self, slug: String) -> Result<Vec<ExperimentBranch>> {
263        self.get_all_experiments()?
264            .into_iter()
265            .find(|e| e.slug == slug)
266            .map(|e| e.branches.into_iter().map(|b| b.into()).collect())
267            .ok_or(NimbusError::NoSuchExperiment(slug))
268    }
269
270    pub fn get_experiment_participation(&self) -> Result<bool> {
271        let db = self.db()?;
272        let reader = db.read()?;
273        get_experiment_participation(db, &reader)
274    }
275
276    pub fn get_rollout_participation(&self) -> Result<bool> {
277        let db = self.db()?;
278        let reader = db.read()?;
279        get_rollout_participation(db, &reader)
280    }
281
282    pub fn set_experiment_participation(
283        &self,
284        user_participating: bool,
285    ) -> Result<Vec<EnrollmentChangeEvent>> {
286        let db = self.db()?;
287        let mut writer = db.write()?;
288        let mut state = self.mutable_state.lock().unwrap();
289        set_experiment_participation(db, &mut writer, user_participating)?;
290
291        let existing_experiments: Vec<Experiment> =
292            db.get_store(StoreId::Experiments).collect_all(&writer)?;
293        let events = self.evolve_experiments(db, &mut writer, &mut state, &existing_experiments)?;
294        let res = self.end_initialize(db, writer, &mut state);
295        self.record_enrollment_status_telemetry(&mut state)?;
296        res?;
297        Ok(events)
298    }
299
300    pub fn set_rollout_participation(
301        &self,
302        user_participating: bool,
303    ) -> Result<Vec<EnrollmentChangeEvent>> {
304        let db = self.db()?;
305        let mut writer = db.write()?;
306        let mut state = self.mutable_state.lock().unwrap();
307        set_rollout_participation(db, &mut writer, user_participating)?;
308
309        let existing_experiments: Vec<Experiment> =
310            db.get_store(StoreId::Experiments).collect_all(&writer)?;
311        let events = self.evolve_experiments(db, &mut writer, &mut state, &existing_experiments)?;
312        let res = self.end_initialize(db, writer, &mut state);
313        self.record_enrollment_status_telemetry(&mut state)?;
314        res?;
315        Ok(events)
316    }
317
318    pub fn get_active_experiments(&self) -> Result<Vec<EnrolledExperiment>> {
319        self.database_cache.get_active_experiments()
320    }
321
322    pub fn get_all_experiments(&self) -> Result<Vec<Experiment>> {
323        let db = self.db()?;
324        let reader = db.read()?;
325        db.get_store(StoreId::Experiments)
326            .collect_all::<Experiment, _>(&reader)
327    }
328
329    pub fn get_available_experiments(&self) -> Result<Vec<AvailableExperiment>> {
330        let th = self.create_targeting_helper(None)?;
331        Ok(self
332            .get_all_experiments()?
333            .into_iter()
334            .filter(|exp| {
335                is_experiment_available(&th, exp, false) == ExperimentAvailable::Available
336            })
337            .map(|exp| exp.into())
338            .collect())
339    }
340
341    pub fn opt_in_with_branch(
342        &self,
343        experiment_slug: String,
344        branch: String,
345    ) -> Result<Vec<EnrollmentChangeEvent>> {
346        let db = self.db()?;
347        let mut writer = db.write()?;
348        let result = opt_in_with_branch(db, &mut writer, &experiment_slug, &branch)?;
349        let mut state = self.mutable_state.lock().unwrap();
350        self.end_initialize(db, writer, &mut state)?;
351        Ok(result)
352    }
353
354    pub fn opt_out(&self, experiment_slug: String) -> Result<Vec<EnrollmentChangeEvent>> {
355        let db = self.db()?;
356        let mut writer = db.write()?;
357        let result = opt_out(
358            db,
359            &mut writer,
360            &experiment_slug,
361            self.gecko_prefs.as_deref(),
362        )?;
363        let mut state = self.mutable_state.lock().unwrap();
364        self.end_initialize(db, writer, &mut state)?;
365        Ok(result)
366    }
367
368    pub fn fetch_experiments(&self) -> Result<()> {
369        if !self.is_fetch_enabled()? {
370            return Ok(());
371        }
372        info!("fetching experiments");
373        let settings_client = self.settings_client.lock().unwrap();
374        let new_experiments = settings_client.fetch_experiments()?;
375        let db = self.db()?;
376        let mut writer = db.write()?;
377        write_pending_experiments(db, &mut writer, new_experiments)?;
378        writer.commit()?;
379        Ok(())
380    }
381
382    pub fn set_fetch_enabled(&self, allow: bool) -> Result<()> {
383        let db = self.db()?;
384        let mut writer = db.write()?;
385        db.get_store(StoreId::Meta)
386            .put(&mut writer, DB_KEY_FETCH_ENABLED, &allow)?;
387        writer.commit()?;
388        Ok(())
389    }
390
391    pub(crate) fn is_fetch_enabled(&self) -> Result<bool> {
392        let db = self.db()?;
393        let reader = db.read()?;
394        let enabled = db
395            .get_store(StoreId::Meta)
396            .get(&reader, DB_KEY_FETCH_ENABLED)?
397            .unwrap_or(true);
398        Ok(enabled)
399    }
400
401    /**
402     * Calculate the days since install and days since update on the targeting_attributes.
403     */
404    fn update_ta_install_dates(
405        &self,
406        db: &Database,
407        writer: &mut Writer,
408        state: &mut MutexGuard<InternalMutableState>,
409    ) -> Result<()> {
410        // Only set install_date and update_date with this method if it hasn't been set already.
411        // This cuts down on deriving the dates at runtime, but also allows us to use
412        // the test methods set_install_date() and set_update_date() to set up
413        // scenarios for test.
414        if state.install_date.is_none() {
415            let installation_date = self.get_installation_date(db, writer)?;
416            state.install_date = Some(installation_date);
417        }
418        if state.update_date.is_none() {
419            let update_date = self.get_update_date(db, writer)?;
420            state.update_date = Some(update_date);
421        }
422        state.update_time_to_now(Utc::now());
423
424        Ok(())
425    }
426
427    /**
428     * Calculates the active_experiments based on current enrollments for the targeting attributes.
429     */
430    fn update_ta_active_experiments(
431        &self,
432        db: &Database,
433        writer: &Writer,
434        state: &mut MutexGuard<InternalMutableState>,
435    ) -> Result<()> {
436        let enrollments_store = db.get_store(StoreId::Enrollments);
437        let prev_enrollments: Vec<ExperimentEnrollment> = enrollments_store.collect_all(writer)?;
438
439        state
440            .targeting_attributes
441            .update_enrollments(&prev_enrollments);
442
443        Ok(())
444    }
445
446    fn evolve_experiments(
447        &self,
448        db: &Database,
449        writer: &mut Writer,
450        state: &mut InternalMutableState,
451        experiments: &[Experiment],
452    ) -> Result<Vec<EnrollmentChangeEvent>> {
453        let mut targeting_helper = NimbusTargetingHelper::with_targeting_attributes(
454            &state.targeting_attributes,
455            self.event_store.clone(),
456            self.gecko_prefs.clone(),
457        );
458        if let Some(ref recorded_context) = self.recorded_context {
459            recorded_context.record();
460        }
461        let coenrolling_feature_ids = self
462            .coenrolling_feature_ids
463            .iter()
464            .map(|s| s.as_str())
465            .collect();
466        let mut evolver = EnrollmentsEvolver::new(
467            &state.available_randomization_units,
468            &mut targeting_helper,
469            &coenrolling_feature_ids,
470        );
471        evolver.evolve_enrollments_in_db(db, writer, experiments, self.gecko_prefs.as_deref())
472    }
473
474    pub fn apply_pending_experiments(&self) -> Result<Vec<EnrollmentChangeEvent>> {
475        info!("updating experiment list");
476        let db = self.db()?;
477        let mut writer = db.write()?;
478
479        // We'll get the pending experiments which were stored for us, either by fetch_experiments
480        // or by set_experiments_locally.
481        let pending_updates = read_and_remove_pending_experiments(db, &mut writer)?;
482        let mut state = self.mutable_state.lock().unwrap();
483        self.begin_initialize(db, &mut writer, &mut state)?;
484
485        let should_record_enrollment_status = pending_updates.is_some();
486        let res = match pending_updates {
487            Some(new_experiments) => {
488                self.update_ta_active_experiments(db, &writer, &mut state)?;
489                // Perform the enrollment calculations if there are pending experiments.
490                self.evolve_experiments(db, &mut writer, &mut state, &new_experiments)?
491            }
492            None => vec![],
493        };
494
495        // Finish up any cleanup, e.g. copying from database in to memory.
496        let end_init_res = self.end_initialize(db, writer, &mut state);
497        if should_record_enrollment_status {
498            self.record_enrollment_status_telemetry(&mut state)?;
499        }
500        end_init_res?;
501        Ok(res)
502    }
503
504    #[allow(deprecated)] // Bug 1960256 - use of deprecated chrono functions.
505    fn get_installation_date(&self, db: &Database, writer: &mut Writer) -> Result<DateTime<Utc>> {
506        // we first check our context
507        if let Some(context_installation_date) = self.app_context.installation_date {
508            let res = DateTime::<Utc>::from_naive_utc_and_offset(
509                NaiveDateTime::from_timestamp_opt(context_installation_date / 1_000, 0).unwrap(),
510                Utc,
511            );
512            info!("[Nimbus] Retrieved date from Context: {}", res);
513            return Ok(res);
514        }
515        let store = db.get_store(StoreId::Meta);
516        let persisted_installation_date: Option<DateTime<Utc>> =
517            store.get(writer, DB_KEY_INSTALLATION_DATE)?;
518        Ok(
519            if let Some(installation_date) = persisted_installation_date {
520                installation_date
521            } else {
522                Utc::now()
523            },
524        )
525    }
526
527    fn get_update_date(&self, db: &Database, writer: &mut Writer) -> Result<DateTime<Utc>> {
528        let store = db.get_store(StoreId::Meta);
529
530        let persisted_app_version: Option<String> = store.get(writer, DB_KEY_APP_VERSION)?;
531        let update_date: Option<DateTime<Utc>> = store.get(writer, DB_KEY_UPDATE_DATE)?;
532        Ok(
533            match (
534                persisted_app_version,
535                &self.app_context.app_version,
536                update_date,
537            ) {
538                // The app been run before, but has not just been updated.
539                (Some(persisted), Some(current), Some(date)) if persisted == *current => date,
540                // The app has been run before, and just been updated.
541                (Some(persisted), Some(current), _) if persisted != *current => {
542                    let now = Utc::now();
543                    store.put(writer, DB_KEY_APP_VERSION, current)?;
544                    store.put(writer, DB_KEY_UPDATE_DATE, &now)?;
545                    now
546                }
547                // The app has just been installed
548                (None, Some(current), _) => {
549                    let now = Utc::now();
550                    store.put(writer, DB_KEY_APP_VERSION, current)?;
551                    store.put(writer, DB_KEY_UPDATE_DATE, &now)?;
552                    now
553                }
554                // The current version is not available, or the persisted date is not available.
555                (_, _, Some(date)) => date,
556                // Either way, this doesn't appear to be a good production environment.
557                _ => Utc::now(),
558            },
559        )
560    }
561
562    pub fn set_experiments_locally(&self, experiments_json: String) -> Result<()> {
563        let new_experiments = parse_experiments(&experiments_json)?;
564        let db = self.db()?;
565        let mut writer = db.write()?;
566        write_pending_experiments(db, &mut writer, new_experiments)?;
567        writer.commit()?;
568        Ok(())
569    }
570
571    /// Reset all enrollments and experiments in the database.
572    ///
573    /// This should only be used in testing.
574    pub fn reset_enrollments(&self) -> Result<()> {
575        let db = self.db()?;
576        let mut writer = db.write()?;
577        let mut state = self.mutable_state.lock().unwrap();
578        db.clear_experiments_and_enrollments(&mut writer)?;
579        self.end_initialize(db, writer, &mut state)?;
580        Ok(())
581    }
582
583    /// Reset internal state in response to application-level telemetry reset.
584    ///
585    /// When the user resets their telemetry state in the consuming application, we need learn
586    /// the new values of any external randomization units, and we need to reset any unique
587    /// identifiers used internally by the SDK. If we don't then we risk accidentally tracking
588    /// across the telemetry reset, since we could use Nimbus metrics to link their pings from
589    /// before and after the reset.
590    ///
591    pub fn reset_telemetry_identifiers(&self) -> Result<Vec<EnrollmentChangeEvent>> {
592        let mut events = vec![];
593        let db = self.db()?;
594        let mut writer = db.write()?;
595        let mut state = self.mutable_state.lock().unwrap();
596        // If we have no `nimbus_id` when we can safely assume that there's
597        // no other experiment state that needs to be reset.
598        let store = db.get_store(StoreId::Meta);
599        if store.get::<String, _>(&writer, DB_KEY_NIMBUS_ID)?.is_some() {
600            // Each enrollment state now opts out because we don't want to leak information between resets.
601            events = reset_telemetry_identifiers(db, &mut writer)?;
602
603            // Remove any stored event counts
604            db.clear_event_count_data(&mut writer)?;
605
606            // The `nimbus_id` itself is a unique identifier.
607            // N.B. we do this last, as a signal that all data has been reset.
608            store.delete(&mut writer, DB_KEY_NIMBUS_ID)?;
609            self.end_initialize(db, writer, &mut state)?;
610        }
611
612        // (No need to commit `writer` if the above check was false, since we didn't change anything)
613        state.available_randomization_units = Default::default();
614        state.targeting_attributes.nimbus_id = None;
615
616        Ok(events)
617    }
618
619    pub fn nimbus_id(&self) -> Result<Uuid> {
620        let db = self.db()?;
621        let mut writer = db.write()?;
622        let mut state = self.mutable_state.lock().unwrap();
623        let uuid = self.read_or_create_nimbus_id(db, &mut writer, &mut state)?;
624
625        // We don't know whether we needed to generate and save the uuid, so
626        // we commit just in case - this is hopefully close to a noop in that
627        // case!
628        writer.commit()?;
629        Ok(uuid)
630    }
631
632    /// Return the nimbus ID from the database, or create a new one and write it
633    /// to the database.
634    ///
635    /// The internal state will be updated with the nimbus ID.
636    fn read_or_create_nimbus_id(
637        &self,
638        db: &Database,
639        writer: &mut Writer,
640        state: &mut MutexGuard<'_, InternalMutableState>,
641    ) -> Result<Uuid> {
642        let store = db.get_store(StoreId::Meta);
643        let nimbus_id = match store.get(writer, DB_KEY_NIMBUS_ID)? {
644            Some(nimbus_id) => nimbus_id,
645            None => {
646                let nimbus_id = Uuid::new_v4();
647                store.put(writer, DB_KEY_NIMBUS_ID, &nimbus_id)?;
648                nimbus_id
649            }
650        };
651
652        state.available_randomization_units.nimbus_id = Some(nimbus_id.to_string());
653        state.targeting_attributes.nimbus_id = Some(nimbus_id.to_string());
654
655        Ok(nimbus_id)
656    }
657
658    // Sets the nimbus ID - TEST ONLY - should not be exposed to real clients.
659    // (Useful for testing so you can have some control over what experiments
660    // are enrolled)
661    pub fn set_nimbus_id(&self, uuid: &Uuid) -> Result<()> {
662        let db = self.db()?;
663        let mut writer = db.write()?;
664        db.get_store(StoreId::Meta)
665            .put(&mut writer, DB_KEY_NIMBUS_ID, uuid)?;
666        writer.commit()?;
667        Ok(())
668    }
669
670    pub(crate) fn db(&self) -> Result<&Database> {
671        self.db
672            .get_or_try_init(|| Database::new(&self.db_path, self.metrics_handler.clone()))
673    }
674
675    fn merge_additional_context(&self, context: Option<JsonObject>) -> Result<Value> {
676        let context = context.map(Value::Object);
677        let targeting = match serde_json::to_value(self.get_targeting_attributes()) {
678            Ok(v) => v,
679            Err(e) => return Err(NimbusError::JSONError("targeting = nimbus::stateful::nimbus_client::NimbusClient::merge_additional_context::serde_json::to_value".into(), e.to_string()))
680        };
681        let context = match context {
682            Some(v) => v.defaults(&targeting)?,
683            None => targeting,
684        };
685
686        Ok(context)
687    }
688
689    pub fn create_targeting_helper(
690        &self,
691        additional_context: Option<JsonObject>,
692    ) -> Result<Arc<NimbusTargetingHelper>> {
693        let context = self.merge_additional_context(additional_context)?;
694        let helper =
695            NimbusTargetingHelper::new(context, self.event_store.clone(), self.gecko_prefs.clone());
696        Ok(Arc::new(helper))
697    }
698
699    pub fn create_targeting_helper_with_context(
700        &self,
701        context: Value,
702    ) -> Arc<NimbusTargetingHelper> {
703        Arc::new(NimbusTargetingHelper::new(
704            context,
705            self.event_store.clone(),
706            self.gecko_prefs.clone(),
707        ))
708    }
709
710    pub fn create_string_helper(
711        &self,
712        additional_context: Option<JsonObject>,
713    ) -> Result<Arc<NimbusStringHelper>> {
714        let context = self.merge_additional_context(additional_context)?;
715        let helper = NimbusStringHelper::new(context.as_object().unwrap().to_owned());
716        Ok(Arc::new(helper))
717    }
718
719    /// Records an event for the purposes of behavioral targeting.
720    ///
721    /// This function is used to record and persist data used for the behavioral
722    /// targeting such as "core-active" user targeting.
723    pub fn record_event(&self, event_id: String, count: i64) -> Result<()> {
724        let mut event_store = self.event_store.lock().unwrap();
725        event_store.record_event(count as u64, &event_id, None)?;
726        event_store.persist_data(self.db()?)?;
727        Ok(())
728    }
729
730    /// Records an event for the purposes of behavioral targeting.
731    ///
732    /// This differs from the `record_event` method in that the event is recorded as if it were
733    /// recorded `seconds_ago` in the past. This makes it very useful for testing.
734    pub fn record_past_event(&self, event_id: String, seconds_ago: i64, count: i64) -> Result<()> {
735        if seconds_ago < 0 {
736            return Err(NimbusError::BehaviorError(BehaviorError::InvalidDuration(
737                "Time duration in the past must be positive".to_string(),
738            )));
739        }
740        let mut event_store = self.event_store.lock().unwrap();
741        event_store.record_past_event(
742            count as u64,
743            &event_id,
744            None,
745            chrono::Duration::seconds(seconds_ago),
746        )?;
747        event_store.persist_data(self.db()?)?;
748        Ok(())
749    }
750
751    /// Advances the event store's concept of `now` artificially.
752    ///
753    /// This works alongside `record_event` and `record_past_event` for testing purposes.
754    pub fn advance_event_time(&self, by_seconds: i64) -> Result<()> {
755        if by_seconds < 0 {
756            return Err(NimbusError::BehaviorError(BehaviorError::InvalidDuration(
757                "Time duration in the future must be positive".to_string(),
758            )));
759        }
760        let mut event_store = self.event_store.lock().unwrap();
761        event_store.advance_datum(chrono::Duration::seconds(by_seconds));
762        Ok(())
763    }
764
765    /// Clear all events in the Nimbus event store.
766    ///
767    /// This should only be used in testing or cases where the previous event store is no longer viable.
768    pub fn clear_events(&self) -> Result<()> {
769        let mut event_store = self.event_store.lock().unwrap();
770        event_store.clear(self.db()?)?;
771        Ok(())
772    }
773
774    pub fn event_store(&self) -> Arc<Mutex<EventStore>> {
775        self.event_store.clone()
776    }
777
778    pub fn dump_state_to_log(&self) -> Result<()> {
779        let experiments = self.get_active_experiments()?;
780        info!("{0: <65}| {1: <30}| {2}", "Slug", "Features", "Branch");
781        for exp in &experiments {
782            info!(
783                "{0: <65}| {1: <30}| {2}",
784                &exp.slug,
785                &exp.feature_ids.join(", "),
786                &exp.branch_slug
787            );
788        }
789        Ok(())
790    }
791
792    /// Given a Gecko pref state and a pref unenroll reason, unenroll from an experiment
793    pub fn unenroll_for_gecko_pref(
794        &self,
795        pref_state: GeckoPrefState,
796        pref_unenroll_reason: PrefUnenrollReason,
797    ) -> Result<Vec<EnrollmentChangeEvent>> {
798        let mut events = Vec::new();
799        if let Some(prefs) = self.gecko_prefs.clone() {
800            {
801                let mut pref_store_state = prefs.get_mutable_pref_state();
802                pref_store_state.update_pref_state(&pref_state);
803            }
804            let enrollments = self
805                .database_cache
806                .get_enrollments_for_pref(&pref_state.gecko_pref.pref)?;
807
808            let db = self.db()?;
809            let mut writer = db.write()?;
810
811            if let Some(enrollments) = enrollments {
812                for experiment_slug in enrollments {
813                    unenroll_for_pref(
814                        db,
815                        &mut writer,
816                        &experiment_slug,
817                        pref_unenroll_reason,
818                        &pref_state.gecko_pref.pref,
819                        self.gecko_prefs.as_deref(),
820                        &mut events,
821                    )?;
822                }
823            } else {
824                warn!(
825                    "Could not find enrollment. Could unenrollment already occurred through another preference?"
826                )
827            }
828
829            let mut state = self.mutable_state.lock().unwrap();
830            self.end_initialize(db, writer, &mut state)?;
831        }
832        Ok(events)
833    }
834
835    pub fn register_previous_gecko_pref_states(
836        &self,
837        gecko_pref_states: &[GeckoPrefState],
838    ) -> Result<()> {
839        let all_prev_gecko_pref_states =
840            super::gecko_prefs::build_prev_gecko_pref_states(gecko_pref_states);
841
842        let db = self.db()?;
843        let mut writer = db.write()?;
844
845        for (experiment_slug, prev_gecko_pref_states) in all_prev_gecko_pref_states {
846            Self::add_prev_gecko_pref_state_for_experiment(
847                db,
848                &mut writer,
849                &experiment_slug,
850                prev_gecko_pref_states,
851            )?;
852        }
853
854        let coenrolling_ids = self
855            .coenrolling_feature_ids
856            .iter()
857            .map(|s| s.as_str())
858            .collect();
859
860        // Registering the Gecko original values does not require a Gecko update,
861        // but the cache does need to be refreshed.
862        self.database_cache.commit_and_update(
863            db,
864            writer,
865            &coenrolling_ids,
866            self.gecko_prefs.clone(),
867            false,
868        )?;
869
870        Ok(())
871    }
872
873    pub(crate) fn add_prev_gecko_pref_state_for_experiment(
874        db: &Database,
875        writer: &mut Writer,
876        experiment_slug: &str,
877        prev_gecko_pref_states: Vec<PreviousGeckoPrefState>,
878    ) -> Result<()> {
879        let enrollments = db.get_store(StoreId::Enrollments);
880
881        if let Ok(Some(existing_enrollment)) =
882            enrollments.get::<ExperimentEnrollment, Writer>(writer, experiment_slug)
883        {
884            // Previous states are only valid on Enrolled experiments
885            let updated_states =
886                existing_enrollment.on_add_gecko_pref_states(prev_gecko_pref_states);
887            enrollments.put(writer, experiment_slug, &updated_states)?;
888        }
889        Ok(())
890    }
891
892    pub fn get_previous_gecko_pref_states(
893        &self,
894        experiment_slug: String,
895    ) -> Result<Option<Vec<PreviousGeckoPrefState>>> {
896        let db = self.db()?;
897        let reader = db.read()?;
898
899        Ok(db
900            .get_store(StoreId::Enrollments)
901            .get::<ExperimentEnrollment, _>(&reader, &experiment_slug)?
902            .and_then(|enrollment| {
903                if let EnrollmentStatus::Enrolled {
904                    prev_gecko_pref_states: prev_gecko_pref_state,
905                    ..
906                } = enrollment.status
907                {
908                    prev_gecko_pref_state
909                } else {
910                    None
911                }
912            }))
913    }
914
915    pub fn set_install_time(&mut self, then: DateTime<Utc>) {
916        let mut state = self.mutable_state.lock().unwrap();
917        state.install_date = Some(then);
918        state.update_time_to_now(Utc::now());
919    }
920
921    pub fn set_update_time(&mut self, then: DateTime<Utc>) {
922        let mut state = self.mutable_state.lock().unwrap();
923        state.update_date = Some(then);
924        state.update_time_to_now(Utc::now());
925    }
926
927    /// This is only called from `get_feature_config_variables` which is itself is cached with
928    /// thread safety in the FeatureHolder.kt and FeatureHolder.swift
929    fn record_feature_activation_if_needed(&self, feature_id: &str) {
930        if let Ok(Some(f)) = self.database_cache.get_enrollment_by_feature(feature_id)
931            && f.branch.is_some()
932            && !self.coenrolling_feature_ids.contains(&f.feature_id)
933        {
934            self.metrics_handler.record_feature_activation(f.into());
935        }
936    }
937
938    pub fn record_feature_exposure(&self, feature_id: String, slug: Option<String>) {
939        let event = if let Some(slug) = slug {
940            if let Ok(Some(branch)) = self.database_cache.get_experiment_branch(&slug) {
941                Some(FeatureExposureExtraDef {
942                    feature_id,
943                    branch: Some(branch),
944                    slug,
945                })
946            } else {
947                None
948            }
949        } else if let Ok(Some(f)) = self.database_cache.get_enrollment_by_feature(&feature_id) {
950            if f.branch.is_some() {
951                Some(f.into())
952            } else {
953                None
954            }
955        } else {
956            None
957        };
958
959        if let Some(event) = event {
960            self.metrics_handler.record_feature_exposure(event);
961        }
962    }
963
964    pub fn record_malformed_feature_config(&self, feature_id: String, part_id: String) {
965        let event = if let Ok(Some(f)) = self.database_cache.get_enrollment_by_feature(&feature_id)
966        {
967            MalformedFeatureConfigExtraDef::from_feature_and_part(f, part_id)
968        } else {
969            MalformedFeatureConfigExtraDef::new(feature_id, part_id)
970        };
971        self.metrics_handler.record_malformed_feature_config(event);
972    }
973
974    fn record_enrollment_status_telemetry(
975        &self,
976        state: &mut MutexGuard<InternalMutableState>,
977    ) -> Result<()> {
978        let targeting_helper = NimbusTargetingHelper::new(
979            state.targeting_attributes.clone(),
980            self.event_store.clone(),
981            self.gecko_prefs.clone(),
982        );
983        let experiments = self.database_cache.get_experiments()?;
984        let experiments = experiments
985            .iter()
986            .filter(|exp| {
987                is_experiment_available(&targeting_helper, exp, true)
988                    == ExperimentAvailable::Available
989            })
990            .map(|exp| &*exp.slug)
991            .collect::<HashSet<&str>>();
992        self.metrics_handler.record_enrollment_statuses(
993            self.database_cache
994                .get_enrollments()?
995                .into_iter()
996                .filter_map(|e| match experiments.contains(&*e.slug) {
997                    true => Some(e.into()),
998                    false => None,
999                })
1000                .collect(),
1001        );
1002        self.metrics_handler.submit_targeting_context();
1003        Ok(())
1004    }
1005
1006    pub fn get_available_firefox_labs(&self) -> Result<Vec<FirefoxLabsMetadata>> {
1007        let mut state = self.mutable_state.lock().unwrap();
1008        state.update_time_to_now(Utc::now());
1009
1010        let targeting_helper = NimbusTargetingHelper::with_targeting_attributes(
1011            &state.targeting_attributes,
1012            self.event_store.clone(),
1013            self.gecko_prefs.clone(),
1014        );
1015
1016        self.database_cache.get_available_firefox_labs_metadata(
1017            &state.available_randomization_units,
1018            &targeting_helper,
1019            &self.coenrolling_feature_ids,
1020        )
1021    }
1022
1023    pub fn enroll_in_firefox_lab(&self, slug: &str) -> Result<FirefoxLabsEnrollResult> {
1024        let feature_conflict = self
1025            .database_cache
1026            .check_for_feature_conflict(slug, &self.coenrolling_feature_ids)?;
1027
1028        let db = self.db()?;
1029        let mut writer = db.write()?;
1030        let result = enroll_in_firefox_lab(db, &mut writer, slug, feature_conflict);
1031        let mut state = self.mutable_state.lock().unwrap();
1032        self.end_initialize(db, writer, &mut state)?;
1033        result
1034    }
1035
1036    pub fn unenroll_from_firefox_lab(&self, slug: &str) -> Result<FirefoxLabsUnenrollResult> {
1037        let db = self.db()?;
1038        let mut writer = db.write()?;
1039        let result = unenroll_from_firefox_lab(db, &mut writer, slug, self.gecko_prefs.as_deref());
1040        let mut state = self.mutable_state.lock().unwrap();
1041        self.end_initialize(db, writer, &mut state)?;
1042        result
1043    }
1044
1045    pub fn unenroll_from_all_firefox_labs(&self) -> Result<Vec<EnrollmentChangeEvent>> {
1046        let db = self.db()?;
1047        let mut writer = db.write()?;
1048        let result = unenroll_from_all_firefox_labs(db, &mut writer, self.gecko_prefs.as_deref());
1049        let mut state = self.mutable_state.lock().unwrap();
1050        self.end_initialize(db, writer, &mut state)?;
1051        result
1052    }
1053
1054    #[cfg(test)]
1055    pub fn get_experiment_enrollment(&self, slug: &str) -> Result<Option<ExperimentEnrollment>> {
1056        self.database_cache.get_experiment_enrollment(slug)
1057    }
1058}
1059
1060pub fn get_active_enrollments<P: AsRef<Path>>(db_path: &P) -> Result<Vec<EnrollmentSlugs>> {
1061    let db = Database::open_single(db_path.as_ref(), StoreId::Enrollments)?;
1062    let reader = db.read()?;
1063    let enrollments: Vec<ExperimentEnrollment> = db.store.collect_all(&reader)?;
1064
1065    Ok(enrollments
1066        .into_iter()
1067        .filter_map(|enrollment| {
1068            if let EnrollmentStatus::Enrolled {
1069                branch: branch_slug,
1070                ..
1071            } = enrollment.status
1072            {
1073                Some(EnrollmentSlugs {
1074                    slug: enrollment.slug,
1075                    branch_slug,
1076                })
1077            } else {
1078                None
1079            }
1080        })
1081        .collect())
1082}
1083
1084pub struct NimbusStringHelper {
1085    context: JsonObject,
1086}
1087
1088impl NimbusStringHelper {
1089    fn new(context: JsonObject) -> Self {
1090        Self { context }
1091    }
1092
1093    pub fn get_uuid(&self, template: String) -> Option<String> {
1094        if template.contains("{uuid}") {
1095            let uuid = Uuid::new_v4();
1096            Some(uuid.to_string())
1097        } else {
1098            None
1099        }
1100    }
1101
1102    pub fn string_format(&self, template: String, uuid: Option<String>) -> String {
1103        match uuid {
1104            Some(uuid) => {
1105                let mut map = self.context.clone();
1106                map.insert("uuid".to_string(), Value::String(uuid));
1107                fmt_with_map(&template, &map)
1108            }
1109            _ => fmt_with_map(&template, &self.context),
1110        }
1111    }
1112}
1113
1114#[cfg(feature = "stateful-uniffi-bindings")]
1115uniffi::include_scaffolding!("nimbus");