Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

[Full changelog](https://github.com/mozilla/glean/compare/v70.2.0...main)

* General
* Set up RKV storage to store submitted pings in memory when enabled ([#3661](https://github.com/mozilla/glean/pull/3661))
* iOS
* Re-enable the SQLite backend by default ([#3650](https://github.com/mozilla/glean/pull/3650))

Expand Down
73 changes: 73 additions & 0 deletions glean-core/src/database/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ pub mod migration;
#[cfg(feature = "sqlite")]
pub mod sqlite;

use crate::{JsonValue, Result};
use chrono::{DateTime, Utc};
#[cfg(feature = "sqlite")]
pub use conn_ext::ConnExt;

Expand All @@ -21,3 +23,74 @@ pub use sqlite::Database;

#[cfg(not(feature = "sqlite"))]
pub use rkv::Database;

/// A trait defining the methods for a database to handle storing submitted pings.
pub(crate) trait StoredSubmittedPingHandler {
/// Gets all pings in the `submitted_pings` table.
fn get_all_submitted_pings(&self) -> Vec<crate::SubmittedPing>;

/// Returns all submitted pings in the `submitted_pings` table that match a supplied ping name.
///
/// # Arguments
///
/// * `ping` - The name of the pings to return.
fn get_submitted_pings_by_name(&self, ping: &str) -> Vec<crate::SubmittedPing>;

/// Marks a particular ping as uploaded.
///
/// # Arguments
///
/// * `document_id` - The ping to mark as uploaded.
/// * `date_uploaded` - The UTC date/time the ping was uploaded.
///
/// # Returns
///
/// A `usize` representing the number of rows updated.
fn mark_ping_as_uploaded(&self, document_id: &str, date_uploaded: DateTime<Utc>) -> usize;

/// Marks a particular ping as upload failed.
///
/// # Arguments
///
/// * `document_id` - The ping to mark as upload failed.
///
/// # Returns
///
/// A `usize` representing the number of rows updated.
fn mark_ping_as_upload_failed(&self, document_id: &str) -> usize;

/// Stores a submitted ping into the `submitted_pings` table.
///
/// # Arguments
///
/// * `document_id` - The unique identifier for the ping.
/// * `ping` - The name of the ping.
/// * `date_submitted` - The UTC date/time the ping was submitted.
/// * `date_uploaded` - An optional UTC date/time the ping was uploaded.
/// * `payload` - A JSON representation of the content of the ping.
///
/// # Returns
///
/// An empty `Result`.
fn store_submitted_ping(
&self,
document_id: &str,
ping: &str,
date_submitted: DateTime<Utc>,
date_uploaded: Option<DateTime<Utc>>,
upload_failed: Option<DateTime<Utc>>,
payload: JsonValue,
) -> Result<()>;

/// Remove stored submitted pings that are older than `before_time` (or 30 days if not specified)
///
/// # Arguments
///
/// * `before_time` - An optional date – when supplied uses that date as the oldest date_submitted we should keep.
/// Defaults to 30 days if `None` is supplied.
///
/// # Returns
///
/// An empty `Result`.
fn cleanup_submitted_pings(&self, before_time: Option<DateTime<Utc>>) -> Result<()>;
}
135 changes: 130 additions & 5 deletions glean-core/src/database/rkv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,21 +2,21 @@
// License, v. 2.0. If a copy of the MPL was not distributed with this
// file, You can obtain one at https://mozilla.org/MPL/2.0/.

use crate::metrics::dual_labeled_counter::RECORD_SEPARATOR;
use crate::{ErrorKind, JsonValue, SubmittedPing};
use chrono::{DateTime, Utc};
use std::cell::{Cell, RefCell};
use std::collections::btree_map::Entry;
use std::collections::BTreeMap;
use std::collections::{BTreeMap, HashMap};
use std::fs;
use std::io;
use std::num::NonZeroU64;
use std::path::Path;
use std::str;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::RwLock;
use std::sync::{Mutex, RwLock};
use std::time::{Duration, Instant};

use crate::metrics::dual_labeled_counter::RECORD_SEPARATOR;
use crate::ErrorKind;

use malloc_size_of::MallocSizeOf;
use rkv::{StoreError, StoreOptions};

Expand Down Expand Up @@ -95,6 +95,7 @@ pub fn rkv_new(path: &Path) -> std::result::Result<(Rkv, RkvLoadState), rkv::Sto
}

use crate::common_metric_data::CommonMetricDataInternal;
use crate::database::StoredSubmittedPingHandler;
use crate::metrics::Metric;
use crate::Glean;
use crate::Lifetime;
Expand Down Expand Up @@ -147,6 +148,10 @@ pub struct Database {
/// Times an Rkv write-commit took.
/// Re-applied as samples in a timing distribution later.
pub(crate) write_timings: RefCell<Vec<i64>>,

/// An in-memory store for submitted pings. This could get pretty big, but
/// since it's a development-only thing that's alright.
submitted_pings_store: Mutex<HashMap<String, SubmittedPing>>,
}

impl MallocSizeOf for Database {
Expand Down Expand Up @@ -264,6 +269,7 @@ impl Database {
file_size,
rkv_load_state,
write_timings,
submitted_pings_store: Mutex::new(HashMap::new()),
Comment thread
jeddai marked this conversation as resolved.
};

db.load_ping_lifetime_data();
Expand Down Expand Up @@ -1009,6 +1015,125 @@ impl Database {
}
}

impl StoredSubmittedPingHandler for Database {
fn get_all_submitted_pings(&self) -> Vec<SubmittedPing> {
let lock = self
.submitted_pings_store
.lock()
.expect("Unable to lock submitted pings store");
let mut res = lock.values().cloned().collect::<Vec<SubmittedPing>>();
res.sort_by(|a, b| {
if a.submitted_date() < b.submitted_date() {
std::cmp::Ordering::Greater
} else if a.submitted_date() > b.submitted_date() {
std::cmp::Ordering::Less
} else {
std::cmp::Ordering::Equal
}
});
res
}

fn get_submitted_pings_by_name(&self, ping: &str) -> Vec<SubmittedPing> {
let lock = self
.submitted_pings_store
.lock()
.expect("Unable to lock submitted pings store");
let mut res = lock
.values()
.filter(|v| v.ping == ping)
.cloned()
.collect::<Vec<SubmittedPing>>();
res.sort_by(|a, b| {
if a.submitted_date() < b.submitted_date() {
std::cmp::Ordering::Greater
} else if a.submitted_date() > b.submitted_date() {
std::cmp::Ordering::Less
} else {
std::cmp::Ordering::Equal
}
});
res
}

fn mark_ping_as_uploaded(&self, document_id: &str, date_uploaded: DateTime<Utc>) -> usize {
let mut lock = self
.submitted_pings_store
.lock()
.expect("Unable to lock submitted pings store");
let entry = lock.get_mut(document_id);
if let Some(p) = entry {
p.uploaded_date = Some(date_uploaded.to_rfc3339());
1
} else {
0
}
}

fn mark_ping_as_upload_failed(&self, document_id: &str) -> usize {
let mut lock = self
.submitted_pings_store
.lock()
.expect("Unable to lock submitted pings store");
let entry = lock.get_mut(document_id);
if let Some(p) = entry {
p.upload_failed = Some(Utc::now().to_rfc3339());
1
} else {
0
}
}

fn store_submitted_ping(
&self,
document_id: &str,
ping: &str,
date_submitted: DateTime<Utc>,
date_uploaded: Option<DateTime<Utc>>,
upload_failed: Option<DateTime<Utc>>,
payload: JsonValue,
) -> Result<()> {
let mut lock = self
.submitted_pings_store
.lock()
.expect("Unable to lock submitted pings store");
lock.insert(
document_id.to_string(),
SubmittedPing {
document_id: document_id.to_string(),
ping: ping.to_string(),
submitted_date: date_submitted.to_rfc3339(),
uploaded_date: date_uploaded.map(|d| d.to_rfc3339()),
upload_failed: upload_failed.map(|d| d.to_rfc3339()),
payload: Some(payload),
},
);
Ok(())
}

fn cleanup_submitted_pings(&self, before_time: Option<DateTime<Utc>>) -> Result<()> {
Comment thread
jeddai marked this conversation as resolved.
let mut lock = self
.submitted_pings_store
.lock()
.expect("Unable to lock submitted pings store");
let days_30 = Duration::from_secs(30 * 24 * 60 * 60);
let mut time = before_time.unwrap_or_else(|| Utc::now() - days_30);
if let Some(t) = before_time {
time = t;
}
let mut keys = vec![];
for (k, v) in lock.iter_mut() {
if v.submitted_date() <= time {
keys.push(k.clone());
}
}
for k in keys {
lock.remove(&k);
}
Ok(())
}
}

#[cfg(test)]
mod test {
use super::*;
Expand Down
Loading