Created
September 21, 2026 13:01
-
-
Save badboy/7ddc415b3da2d52004ef576becfa3627 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| --- ../glean-ro/glean-core/src/database/mod.rs 2026-09-21 10:10:30 | |
| +++ glean-core/src/database/rkv.rs 2026-09-21 14:54:27 | |
| @@ -14,6 +14,7 @@ | |
| use std::sync::RwLock; | |
| use std::time::{Duration, Instant}; | |
| +use crate::metrics::dual_labeled_counter::RECORD_SEPARATOR; | |
| use crate::ErrorKind; | |
| use malloc_size_of::MallocSizeOf; | |
| @@ -356,6 +357,18 @@ | |
| data.insert(metric_id, metric); | |
| } | |
| } | |
| + } | |
| + | |
| + pub fn iter_store<F>( | |
| + &self, | |
| + lifetime: Lifetime, | |
| + storage_name: &str, | |
| + transaction_fn: F, | |
| + ) -> Result<()> | |
| + where | |
| + F: FnMut(&[u8], &[&str], &Metric), | |
| + { | |
| + self.iter_store_from(lifetime, storage_name, None, transaction_fn) | |
| } | |
| /// Iterates with the provided transaction function | |
| @@ -386,8 +399,9 @@ | |
| storage_name: &str, | |
| metric_key: Option<&str>, | |
| mut transaction_fn: F, | |
| - ) where | |
| - F: FnMut(&[u8], &Metric), | |
| + ) -> Result<()> | |
| + where | |
| + F: FnMut(&[u8], &[&str], &Metric), | |
| { | |
| let iter_start = Self::get_storage_key(storage_name, metric_key); | |
| let len = iter_start.len(); | |
| @@ -401,18 +415,30 @@ | |
| .expect("Can't read ping lifetime data"); | |
| for (key, value) in data.iter() { | |
| if key.starts_with(&iter_start) { | |
| - let key = &key[len..]; | |
| - transaction_fn(key.as_bytes(), value); | |
| + let metric_id = &key[len..]; | |
| + let (metric_id, labels_str) = metric_id | |
| + .split_once(|c| ['/', RECORD_SEPARATOR].contains(&c)) | |
| + .unwrap_or((metric_id, "")); | |
| + | |
| + let labels: &[&str] = if labels_str.is_empty() { | |
| + &[] | |
| + } else if labels_str.contains(RECORD_SEPARATOR) { | |
| + let (key, category) = labels_str.split_once(RECORD_SEPARATOR).unwrap(); | |
| + &[key, category] | |
| + } else { | |
| + &[labels_str] | |
| + }; | |
| + transaction_fn(metric_id.as_bytes(), labels, value); | |
| } | |
| } | |
| - return; | |
| + return Ok(()); | |
| } | |
| } | |
| - let reader = unwrap_or!(self.rkv.read(), return); | |
| + let reader = unwrap_or!(self.rkv.read(), return Ok(())); | |
| let mut iter = unwrap_or!( | |
| self.get_store(lifetime).iter_from(&reader, &iter_start), | |
| - return | |
| + return Ok(()) | |
| ); | |
| while let Some(Ok((metric_id, value))) = iter.next() { | |
| @@ -420,13 +446,38 @@ | |
| break; | |
| } | |
| - let metric_id = &metric_id[len..]; | |
| + if metric_key.is_some() && metric_id.len() > iter_start.len() { | |
| + // It's either the key we're looking for or a labeled metric (after `/` or the record separator) | |
| + let next_char = metric_id[iter_start.len()] as char; | |
| + if !['/', RECORD_SEPARATOR].contains(&next_char) { | |
| + break; | |
| + } | |
| + } | |
| + | |
| let metric: Metric = match value { | |
| rkv::Value::Blob(blob) => unwrap_or!(bincode::deserialize(blob), continue), | |
| _ => continue, | |
| }; | |
| - transaction_fn(metric_id, &metric); | |
| + | |
| + let metric_id = &metric_id[len..]; | |
| + let metric_id = String::from_utf8_lossy(metric_id); | |
| + let (metric_id, labels_str) = metric_id | |
| + .split_once(|c| ['/', RECORD_SEPARATOR].contains(&c)) | |
| + .unwrap_or((&metric_id, "")); | |
| + | |
| + let labels: &[&str] = if labels_str.is_empty() { | |
| + &[] | |
| + } else if labels_str.contains(RECORD_SEPARATOR) { | |
| + let (key, category) = labels_str.split_once(RECORD_SEPARATOR).unwrap(); | |
| + &[key, category] | |
| + } else { | |
| + &[labels_str] | |
| + }; | |
| + | |
| + transaction_fn(metric_id.as_bytes(), labels, &metric); | |
| } | |
| + | |
| + Ok(()) | |
| } | |
| /// Determines if the storage has the given metric. | |
| @@ -466,6 +517,39 @@ | |
| .get(&reader, &key) | |
| .unwrap_or(None) | |
| .is_some() | |
| + } | |
| + | |
| + /// Get a single metric by name from storage | |
| + pub fn get_metric( | |
| + &self, | |
| + glean: &Glean, | |
| + data: &CommonMetricDataInternal, | |
| + storage_name: &str, | |
| + ) -> Option<Metric> { | |
| + let metric_identifier = &data.identifier(glean, false); | |
| + let key = Self::get_storage_key(storage_name, Some(metric_identifier)); | |
| + let lifetime = data.inner.lifetime; | |
| + | |
| + // Lifetime::Ping data is not persisted to disk if | |
| + // Glean has `delay_ping_lifetime_io` set to true | |
| + if lifetime == Lifetime::Ping { | |
| + if let Some(ping_lifetime_data) = &self.ping_lifetime_data { | |
| + let data = ping_lifetime_data.read().unwrap(); | |
| + let metric = data.get(&key)?; | |
| + return Some(metric.clone()); | |
| + } | |
| + } | |
| + | |
| + let reader = unwrap_or!(self.rkv.read(), return None); | |
| + let value = self | |
| + .get_store(lifetime) | |
| + .get(&reader, &key) | |
| + .unwrap_or(None)?; | |
| + if let rkv::Value::Blob(blob) = value { | |
| + bincode::deserialize(blob).ok() | |
| + } else { | |
| + None | |
| + } | |
| } | |
| /// Writes to the specified storage with the provided transaction function. | |
| @@ -486,7 +570,7 @@ | |
| /// Records a metric in the underlying storage system. | |
| pub fn record(&self, glean: &Glean, data: &CommonMetricDataInternal, value: &Metric) { | |
| - let name = data.identifier(glean); | |
| + let name = data.identifier(glean, true); | |
| for ping_name in data.storage_names() { | |
| if glean.is_ping_enabled(ping_name) { | |
| if let Err(e) = | |
| @@ -554,7 +638,7 @@ | |
| where | |
| F: FnMut(Option<Metric>) -> Metric, | |
| { | |
| - let name = data.identifier(glean); | |
| + let name = data.identifier(glean, true); | |
| for ping_name in data.storage_names() { | |
| if glean.is_ping_enabled(ping_name) { | |
| if let Err(e) = self.record_per_lifetime_with( | |
| @@ -658,22 +742,24 @@ | |
| /// | |
| /// This function will **not** panic on database errors. | |
| pub fn clear_ping_lifetime_storage(&self, storage_name: &str) -> Result<()> { | |
| + let storage_name = Self::get_storage_key(storage_name, None); | |
| + | |
| // Lifetime::Ping data will be saved to `ping_lifetime_data` | |
| // in case `delay_ping_lifetime_io` is set to true | |
| if let Some(ping_lifetime_data) = &self.ping_lifetime_data { | |
| ping_lifetime_data | |
| .write() | |
| .expect("Can't access ping lifetime data as writable") | |
| - .retain(|metric_id, _| !metric_id.starts_with(storage_name)); | |
| + .retain(|metric_id, _| !metric_id.starts_with(&storage_name)); | |
| } | |
| self.write_with_store(Lifetime::Ping, |mut writer, store| { | |
| let mut metrics = Vec::new(); | |
| { | |
| - let mut iter = store.iter_from(&writer, storage_name)?; | |
| + let mut iter = store.iter_from(&writer, &storage_name)?; | |
| while let Some(Ok((metric_id, _))) = iter.next() { | |
| if let Ok(metric_id) = std::str::from_utf8(metric_id) { | |
| - if !metric_id.starts_with(storage_name) { | |
| + if !metric_id.starts_with(&storage_name) { | |
| break; | |
| } | |
| metrics.push(metric_id.to_owned()); | |
| @@ -1027,7 +1113,8 @@ | |
| // Verify that the data is correctly recorded. | |
| let mut found_metrics = 0; | |
| - let mut snapshotter = |metric_id: &[u8], metric: &Metric| { | |
| + let mut snapshotter = |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| found_metrics += 1; | |
| let metric_id = String::from_utf8_lossy(metric_id).into_owned(); | |
| assert_eq!(test_metric_id, metric_id); | |
| @@ -1037,7 +1124,7 @@ | |
| } | |
| }; | |
| - db.iter_store_from(Lifetime::Ping, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::Ping, test_storage, None, &mut snapshotter); | |
| assert_eq!(1, found_metrics, "We only expect 1 Lifetime.Ping metric."); | |
| } | |
| @@ -1061,7 +1148,8 @@ | |
| // Verify that the data is correctly recorded. | |
| let mut found_metrics = 0; | |
| - let mut snapshotter = |metric_id: &[u8], metric: &Metric| { | |
| + let mut snapshotter = |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| found_metrics += 1; | |
| let metric_id = String::from_utf8_lossy(metric_id).into_owned(); | |
| assert_eq!(test_metric_id, metric_id); | |
| @@ -1071,7 +1159,7 @@ | |
| } | |
| }; | |
| - db.iter_store_from(Lifetime::Application, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::Application, test_storage, None, &mut snapshotter); | |
| assert_eq!( | |
| 1, found_metrics, | |
| "We only expect 1 Lifetime.Application metric." | |
| @@ -1098,7 +1186,8 @@ | |
| // Verify that the data is correctly recorded. | |
| let mut found_metrics = 0; | |
| - let mut snapshotter = |metric_id: &[u8], metric: &Metric| { | |
| + let mut snapshotter = |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| found_metrics += 1; | |
| let metric_id = String::from_utf8_lossy(metric_id).into_owned(); | |
| assert_eq!(test_metric_id, metric_id); | |
| @@ -1108,7 +1197,7 @@ | |
| } | |
| }; | |
| - db.iter_store_from(Lifetime::User, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::User, test_storage, None, &mut snapshotter); | |
| assert_eq!(1, found_metrics, "We only expect 1 Lifetime.User metric."); | |
| } | |
| @@ -1145,7 +1234,8 @@ | |
| // Take a snapshot for the data, all the lifetimes. | |
| { | |
| let mut snapshot: HashMap<String, String> = HashMap::new(); | |
| - let mut snapshotter = |metric_id: &[u8], metric: &Metric| { | |
| + let mut snapshotter = |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| let metric_id = String::from_utf8_lossy(metric_id).into_owned(); | |
| match metric { | |
| Metric::String(s) => snapshot.insert(metric_id, s.to_string()), | |
| @@ -1153,9 +1243,9 @@ | |
| }; | |
| }; | |
| - db.iter_store_from(Lifetime::User, test_storage, None, &mut snapshotter); | |
| - db.iter_store_from(Lifetime::Ping, test_storage, None, &mut snapshotter); | |
| - db.iter_store_from(Lifetime::Application, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::User, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::Ping, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::Application, test_storage, None, &mut snapshotter); | |
| assert_eq!(3, snapshot.len(), "We expect all lifetimes to be present."); | |
| assert!(snapshot.contains_key("telemetry_test.test_name_user")); | |
| @@ -1169,7 +1259,8 @@ | |
| // Take a snapshot again and check that we're only clearing the Ping lifetime. | |
| { | |
| let mut snapshot: HashMap<String, String> = HashMap::new(); | |
| - let mut snapshotter = |metric_id: &[u8], metric: &Metric| { | |
| + let mut snapshotter = |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| let metric_id = String::from_utf8_lossy(metric_id).into_owned(); | |
| match metric { | |
| Metric::String(s) => snapshot.insert(metric_id, s.to_string()), | |
| @@ -1177,9 +1268,9 @@ | |
| }; | |
| }; | |
| - db.iter_store_from(Lifetime::User, test_storage, None, &mut snapshotter); | |
| - db.iter_store_from(Lifetime::Ping, test_storage, None, &mut snapshotter); | |
| - db.iter_store_from(Lifetime::Application, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::User, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::Ping, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(Lifetime::Application, test_storage, None, &mut snapshotter); | |
| assert_eq!(2, snapshot.len(), "We only expect 2 metrics to be left."); | |
| assert!(snapshot.contains_key("telemetry_test.test_name_user")); | |
| @@ -1224,7 +1315,8 @@ | |
| // Verify that "telemetry_test.single_metric_retain" is still around for all lifetimes. | |
| for lifetime in lifetimes.iter() { | |
| let mut found_metrics = 0; | |
| - let mut snapshotter = |metric_id: &[u8], metric: &Metric| { | |
| + let mut snapshotter = |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| found_metrics += 1; | |
| let metric_id = String::from_utf8_lossy(metric_id).into_owned(); | |
| assert_eq!(format!("{}_retain", metric_id_pattern), metric_id); | |
| @@ -1235,7 +1327,7 @@ | |
| }; | |
| // Check the User lifetime. | |
| - db.iter_store_from(*lifetime, test_storage, None, &mut snapshotter); | |
| + _ = db.iter_store_from(*lifetime, test_storage, None, &mut snapshotter); | |
| assert_eq!( | |
| 1, found_metrics, | |
| "We only expect 1 metric for this lifetime." | |
| @@ -1500,17 +1592,18 @@ | |
| let test_storage = "test-storage"; | |
| let test_data = CommonMetricDataInternal::new("category", "name", test_storage); | |
| - let test_metric_id = test_data.identifier(&glean); | |
| + let test_metric_id = test_data.identifier(&glean, true); | |
| // Attempt to record metric with the record and record_with functions, | |
| // this should work since upload is enabled. | |
| let db = Database::new(dir.path(), true, 0, Duration::ZERO).unwrap(); | |
| db.record(&glean, &test_data, &Metric::String("record".to_owned())); | |
| - db.iter_store_from( | |
| + _ = db.iter_store_from( | |
| Lifetime::Ping, | |
| test_storage, | |
| None, | |
| - &mut |metric_id: &[u8], metric: &Metric| { | |
| + &mut |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| assert_eq!( | |
| String::from_utf8_lossy(metric_id).into_owned(), | |
| test_metric_id | |
| @@ -1525,11 +1618,12 @@ | |
| db.record_with(&glean, &test_data, |_| { | |
| Metric::String("record_with".to_owned()) | |
| }); | |
| - db.iter_store_from( | |
| + _ = db.iter_store_from( | |
| Lifetime::Ping, | |
| test_storage, | |
| None, | |
| - &mut |metric_id: &[u8], metric: &Metric| { | |
| + &mut |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| assert_eq!( | |
| String::from_utf8_lossy(metric_id).into_owned(), | |
| test_metric_id | |
| @@ -1547,11 +1641,12 @@ | |
| // Attempt to record metric with the record and record_with functions, | |
| // this should work since upload is now **disabled**. | |
| db.record(&glean, &test_data, &Metric::String("record_nop".to_owned())); | |
| - db.iter_store_from( | |
| + _ = db.iter_store_from( | |
| Lifetime::Ping, | |
| test_storage, | |
| None, | |
| - &mut |metric_id: &[u8], metric: &Metric| { | |
| + &mut |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| assert_eq!( | |
| String::from_utf8_lossy(metric_id).into_owned(), | |
| test_metric_id | |
| @@ -1565,11 +1660,12 @@ | |
| db.record_with(&glean, &test_data, |_| { | |
| Metric::String("record_with_nop".to_owned()) | |
| }); | |
| - db.iter_store_from( | |
| + _ = db.iter_store_from( | |
| Lifetime::Ping, | |
| test_storage, | |
| None, | |
| - &mut |metric_id: &[u8], metric: &Metric| { | |
| + &mut |metric_id: &[u8], labels: &[&str], metric: &Metric| { | |
| + assert!(labels.is_empty()); | |
| assert_eq!( | |
| String::from_utf8_lossy(metric_id).into_owned(), | |
| test_metric_id |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment