Skip to content
Merged
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

61 changes: 61 additions & 0 deletions aw-datastore/examples/export_memory.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
//! Reproducible export-memory comparison using synthetic, in-memory data only.
//! Build with `cargo build -p aw-datastore --example export_memory --release`.
//! Run the resulting executable under a memory profiler with `stream 100000`
//! or `materialized 100000`. No existing database or server is opened.
use aw_datastore::DatastoreInstance;
use aw_models::{Bucket, BucketMetadata, BucketsExport, TryVec};
use chrono::Utc;
use rusqlite::Connection;
use std::{hint::black_box, time::Instant};

fn main() {
let mode = std::env::args().nth(1).unwrap_or_else(|| "stream".into());
assert!(matches!(mode.as_str(), "stream" | "materialized"));
let count: u32 = std::env::args()
.nth(2)
.unwrap_or_else(|| "100000".into())
.parse()
.unwrap();
let conn = Connection::open_in_memory().unwrap();
let mut ds = DatastoreInstance::new(&conn, true).unwrap();
ds.create_bucket(
&conn,
Bucket {
bid: None,
id: "synthetic".into(),
_type: "test".into(),
client: "test".into(),
hostname: "test".into(),
created: Some(Utc::now()),
data: Default::default(),
metadata: BucketMetadata::default(),
events: None,
last_updated: None,
},
)
.unwrap();
let payload = serde_json::json!({"app": "browser", "title": "x".repeat(256)}).to_string();
conn.execute(
"WITH RECURSIVE seq(n) AS (SELECT 1 UNION ALL SELECT n+1 FROM seq WHERE n<?1)
INSERT INTO events(bucketrow,starttime,endtime,data)
SELECT 1,n*1000000000,n*1000000000+500000000,?2 FROM seq",
rusqlite::params![count, payload],
)
.unwrap();
let mut ds = DatastoreInstance::new(&conn, true).unwrap();
let start = Instant::now();
if mode == "stream" {
ds.write_export(&conn, None, std::io::sink()).unwrap();
} else {
let mut buckets = ds.get_buckets();
for (id, bucket) in &mut buckets {
bucket.events = Some(TryVec::new(
ds.get_events(&conn, id, None, None, None).unwrap(),
));
}
let export = BucketsExport { buckets };
let body = serde_json::to_string(&export).unwrap();
black_box((&export, &body));
}
println!("{mode}: {count} synthetic events, {:?}", start.elapsed());
}
20 changes: 19 additions & 1 deletion aw-datastore/src/datastore.rs
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,10 @@ pub struct DatastoreInstance {
///
/// When `clip` is set to `(starttime_filter_ns, endtime_filter_ns)`, the event is
/// clamped to that query range.
fn parse_event_row(row: &rusqlite::Row, clip: Option<(i64, i64)>) -> rusqlite::Result<Event> {
pub(crate) fn parse_event_row(
row: &rusqlite::Row,
clip: Option<(i64, i64)>,
) -> rusqlite::Result<Event> {
let id = row.get(0)?;
let mut starttime_ns: i64 = row.get(1)?;
let mut endtime_ns: i64 = row.get(2)?;
Expand Down Expand Up @@ -491,6 +494,21 @@ impl DatastoreInstance {
self.buckets_cache.clone()
}

/// Serialize an export one event at a time using the caller's connection.
/// Returns the bucket ID for a single-bucket export (for download naming).
pub fn write_export(
&self,
conn: &Connection,
bucket_id: Option<&str>,
writer: impl std::io::Write,
) -> Result<Option<String>, DatastoreError> {
crate::export::write_export(conn, &self.buckets_cache, bucket_id, writer)?;
Ok(bucket_id.map(str::to_owned).or_else(|| {
(self.buckets_cache.len() == 1)
.then(|| self.buckets_cache.keys().next().unwrap().clone())
}))
}

pub fn insert_events(
&mut self,
conn: &Connection,
Expand Down
220 changes: 220 additions & 0 deletions aw-datastore/src/export.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,220 @@
use std::{collections::HashMap, io::Write};

use aw_models::Bucket;
use rusqlite::Connection;
use serde::{
ser::{Error, SerializeMap, SerializeSeq},
Serialize, Serializer,
};

use crate::DatastoreError;

struct EventRows<'a> {
conn: &'a Connection,
bucket: &'a Bucket,
}

impl Serialize for EventRows<'_> {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let mut stmt = self
.conn
.prepare_cached(
"SELECT id, starttime, endtime, data
FROM events INDEXED BY events_bucketrow_starttime_endtime_index
WHERE bucketrow = ?1 AND endtime >= 0 AND starttime <= ?2
ORDER BY starttime DESC, endtime ASC, id ASC",
)
.map_err(S::Error::custom)?;
let mut rows = stmt
.query(rusqlite::params![self.bucket.bid.unwrap(), i64::MAX])
.map_err(S::Error::custom)?;
let mut seq = serializer.serialize_seq(None)?;
while let Some(row) = rows.next().map_err(S::Error::custom)? {
// Match get_events' default range, clipping and corrupt-row policy.
match crate::datastore::parse_event_row(row, Some((0, i64::MAX))) {
Ok(event) => seq.serialize_element(&event)?,
Err(err) => warn!("Corrupt event in bucket {}: {}", self.bucket.id, err),
}
}
seq.end()
}
}

struct ExportBucket<'a> {
conn: &'a Connection,
bucket: &'a Bucket,
}

impl Serialize for ExportBucket<'_> {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
// Reuse Bucket's field names and serialization rules, replacing only
// events. This allocation contains bucket metadata, never event rows.
let metadata = serde_json::to_value(self.bucket).map_err(S::Error::custom)?;
let metadata = metadata
.as_object()
.ok_or_else(|| S::Error::custom("invalid bucket metadata"))?;
let mut map = serializer.serialize_map(Some(metadata.len()))?;
for (key, value) in metadata {
if key != "events" {
map.serialize_entry(key, value)?;
}
}
map.serialize_entry(
"events",
&EventRows {
conn: self.conn,
bucket: self.bucket,
},
)?;
map.end()
}
}

struct ExportBuckets<'a> {
conn: &'a Connection,
buckets: &'a HashMap<String, Bucket>,
selected: Option<&'a str>,
}

impl Serialize for ExportBuckets<'_> {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let mut map = serializer.serialize_map(None)?;
for (id, bucket) in self.buckets {
if self.selected.is_none() || self.selected == Some(id.as_str()) {
map.serialize_entry(
id,
&ExportBucket {
conn: self.conn,
bucket,
},
)?;
}
}
map.end()
}
}

pub(crate) fn write_export(
conn: &Connection,
buckets: &HashMap<String, Bucket>,
selected: Option<&str>,
writer: impl Write,
) -> Result<(), DatastoreError> {
if let Some(id) = selected {
if !buckets.contains_key(id) {
return Err(DatastoreError::NoSuchBucket(id.to_owned()));
}
}
#[derive(Serialize)]
struct Export<'a> {
buckets: ExportBuckets<'a>,
}
serde_json::to_writer(
writer,
&Export {
buckets: ExportBuckets {
conn,
buckets,
selected,
},
},
)
.map_err(|err| DatastoreError::InternalError(format!("Failed to write export: {err}")))
}

#[cfg(test)]
mod tests {
use super::*;
use crate::DatastoreInstance;
use aw_models::{BucketMetadata, BucketsExport, Event, TryVec};
use chrono::{DateTime, Duration};

fn setup() -> (Connection, DatastoreInstance) {
let conn = Connection::open_in_memory().unwrap();
let mut ds = DatastoreInstance::new(&conn, true).unwrap();
for id in ["populated", "empty"] {
ds.create_bucket(
&conn,
Bucket {
bid: None,
id: id.into(),
_type: "test".into(),
client: "test".into(),
hostname: "host".into(),
created: None,
data: Default::default(),
metadata: BucketMetadata::default(),
events: None,
last_updated: None,
},
)
.unwrap();
}
let events = [(-10, 20), (5, 20), (5, 10), (5, 10), (30, 0)]
.into_iter()
.map(|(start, duration)| {
Event::new(
DateTime::from_timestamp(start, 0).unwrap(),
Duration::seconds(duration),
serde_json::from_value(
serde_json::json!({"text": "quotes \" and unicode ☀", "nested": [1, true]}),
)
.unwrap(),
)
})
.collect();
ds.insert_events(&conn, "populated", events).unwrap();
conn.execute("INSERT INTO events(bucketrow,starttime,endtime,data) VALUES(1,6000000000,7000000000,'invalid json')", []).unwrap();
(conn, ds)
}

#[test]
fn streamed_json_matches_materialized_exports() {
let (conn, mut ds) = setup();
for selected in [None, Some("populated"), Some("empty")] {
let mut buckets = ds.get_buckets();
buckets.retain(|id, _| selected.is_none() || selected == Some(id.as_str()));
for (id, bucket) in &mut buckets {
bucket.events = Some(TryVec::new(
ds.get_events(&conn, id, None, None, None).unwrap(),
));
}
let expected = serde_json::to_value(BucketsExport { buckets }).unwrap();
let mut output = Vec::new();
ds.write_export(&conn, selected, &mut output).unwrap();
assert_eq!(
serde_json::from_slice::<serde_json::Value>(&output).unwrap(),
expected
);
}
}

#[test]
fn missing_bucket_and_writer_failures_propagate() {
let (conn, ds) = setup();
let mut output = Vec::new();
assert!(matches!(
ds.write_export(&conn, Some("missing"), &mut output),
Err(DatastoreError::NoSuchBucket(_))
));
assert!(output.is_empty());
struct FailingWriter(usize);
impl Write for FailingWriter {
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
if self.0 == 0 {
return Err(std::io::Error::new(std::io::ErrorKind::Other, "disk full"));
}
let n = bytes.len().min(self.0);
self.0 -= n;
Ok(n)
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
assert!(matches!(
ds.write_export(&conn, None, FailingWriter(500)),
Err(DatastoreError::InternalError(_))
));
}
}
1 change: 1 addition & 0 deletions aw-datastore/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ macro_rules! json_map {
}

mod datastore;
mod export;
mod legacy_import;
mod privacy_filter;
mod worker;
Expand Down
Loading
Loading