1use 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#[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 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
86pub 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 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 #[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 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 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 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 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 fn update_ta_install_dates(
405 &self,
406 db: &Database,
407 writer: &mut Writer,
408 state: &mut MutexGuard<InternalMutableState>,
409 ) -> Result<()> {
410 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 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 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 self.evolve_experiments(db, &mut writer, &mut state, &new_experiments)?
491 }
492 None => vec![],
493 };
494
495 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)] fn get_installation_date(&self, db: &Database, writer: &mut Writer) -> Result<DateTime<Utc>> {
506 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 (Some(persisted), Some(current), Some(date)) if persisted == *current => date,
540 (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 (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 (_, _, Some(date)) => date,
556 _ => 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 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 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 let store = db.get_store(StoreId::Meta);
599 if store.get::<String, _>(&writer, DB_KEY_NIMBUS_ID)?.is_some() {
600 events = reset_telemetry_identifiers(db, &mut writer)?;
602
603 db.clear_event_count_data(&mut writer)?;
605
606 store.delete(&mut writer, DB_KEY_NIMBUS_ID)?;
609 self.end_initialize(db, writer, &mut state)?;
610 }
611
612 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 writer.commit()?;
629 Ok(uuid)
630 }
631
632 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 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 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 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 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 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 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 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 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 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");