diff --git a/aw-datastore/examples/export_memory.rs b/aw-datastore/examples/export_memory.rs index fa8be614..73648dd3 100644 --- a/aw-datastore/examples/export_memory.rs +++ b/aw-datastore/examples/export_memory.rs @@ -1,16 +1,42 @@ //! 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. +//! Run the resulting executable under a memory profiler with one of: +//! +//! - `stream 100000` — JSON export, incremental (`write_export`) +//! - `materialized 100000` — JSON export, pre-#677 in-memory baseline +//! - `csv 100000` — CSV export, row-at-a-time (`write_events_csv`) +//! - `csv-union 100000` — CSV export where every event carries a distinct data +//! key, exercising the `MAX_CSV_DATA_COLUMNS` fallback +//! +//! 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}; +/// Discards written bytes while counting them, so a single run reports both the +/// serialization cost and the size of the produced document. +#[derive(Default)] +struct CountingWriter(u64); + +impl std::io::Write for CountingWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0 += buf.len() as u64; + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + fn main() { let mode = std::env::args().nth(1).unwrap_or_else(|| "stream".into()); - assert!(matches!(mode.as_str(), "stream" | "materialized")); + assert!(matches!( + mode.as_str(), + "stream" | "materialized" | "csv" | "csv-union" + )); let count: u32 = std::env::args() .nth(2) .unwrap_or_else(|| "100000".into()) @@ -34,28 +60,56 @@ fn main() { }, ) .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 { + let mut out = CountingWriter::default(); + ds.write_events_csv(&conn, "synthetic", None, None, None, &mut out) + .unwrap(); + println!( + "{mode}: {count} synthetic events, {:?}, {} bytes of CSV", + start.elapsed(), + out.0 + ); + } + "stream" => { + ds.write_export(&conn, None, std::io::sink()).unwrap(); + println!("{mode}: {count} synthetic events, {:?}", start.elapsed()); + } + "materialized" => { + 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()); } - let export = BucketsExport { buckets }; - let body = serde_json::to_string(&export).unwrap(); - black_box((&export, &body)); + _ => unreachable!(), } - println!("{mode}: {count} synthetic events, {:?}", start.elapsed()); } diff --git a/aw-datastore/src/datastore.rs b/aw-datastore/src/datastore.rs index aa414036..023570d7 100644 --- a/aw-datastore/src/datastore.rs +++ b/aw-datastore/src/datastore.rs @@ -227,7 +227,12 @@ fn _migrate_v5_to_v6(conn: &Connection) { // bucket's time span as a cheap selectivity estimate; this affects performance, // never which rows qualify. Limited queries retain their ordered starttime scan // so LIMIT can stop early without sorting all matching rows. -fn prefer_endtime_index(bucket: &Bucket, start: i64, end: i64, limit: Option) -> bool { +pub(crate) fn prefer_endtime_index( + bucket: &Bucket, + start: i64, + end: i64, + limit: Option, +) -> bool { if limit.is_some() { return false; } @@ -689,6 +694,29 @@ impl DatastoreInstance { })) } + /// Stream one bucket's events as RFC-4180 CSV, one row at a time. + /// + /// Uses the same filters, clipping, and corrupt-row policy as `get_events`. + pub fn write_events_csv( + &self, + conn: &Connection, + bucket_id: &str, + start: Option>, + end: Option>, + limit: Option, + writer: impl std::io::Write, + ) -> Result<(), DatastoreError> { + crate::export::write_events_csv( + conn, + &self.buckets_cache, + bucket_id, + start, + end, + limit, + writer, + ) + } + pub fn insert_events( &mut self, conn: &Connection, diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index f106ceaf..4ba6e12b 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -1,12 +1,15 @@ -use std::{collections::HashMap, io::Write}; +use std::collections::{HashMap, HashSet}; +use std::io::Write; -use aw_models::Bucket; +use aw_models::{Bucket, Event}; +use chrono::{DateTime, Utc}; use rusqlite::Connection; use serde::{ ser::{Error, SerializeMap, SerializeSeq}, Serialize, Serializer, }; +use crate::datastore::{parse_event_row, prefer_endtime_index}; use crate::DatastoreError; struct EventRows<'a> { @@ -122,6 +125,278 @@ pub(crate) fn write_export( .map_err(|err| DatastoreError::InternalError(format!("Failed to write export: {err}"))) } +fn csv_io_err(err: std::io::Error) -> DatastoreError { + DatastoreError::InternalError(format!("Failed to write CSV export: {err}")) +} + +/// Prefix spreadsheet-formula starters so Excel/Sheets will not execute them. +/// Characters spreadsheets ignore before evaluating a formula starter +/// (Excel and LibreOffice trim leading whitespace, so " =1+1" executes). +const FORMULA_LEADING_WHITESPACE: [char; 4] = [' ', '\t', '\r', '\n']; + +/// Prefix spreadsheet-formula starters so Excel/Sheets will not execute them. +/// Looks past leading whitespace so values like " =1+1" are neutralized too. +fn neutralize_formula(s: &str) -> String { + let trimmed = s.trim_start_matches(FORMULA_LEADING_WHITESPACE); + match trimmed.chars().next() { + Some('=' | '+' | '-' | '@') => format!("'{s}"), + _ => s.to_owned(), + } +} + +/// RFC-4180 field escaping, with formula neutralization applied first. +fn csv_escape(s: &str) -> String { + let s = neutralize_formula(s); + if s.contains([',', '"', '\n', '\r']) { + format!("\"{}\"", s.replace('"', "\"\"")) + } else { + s + } +} + +/// Exact fractional-second duration, matching the JSON nanosecond contract +/// without going through `num_milliseconds()` (which truncates sub-ms). +/// +/// Composed from whole seconds and the signed subsecond remainder rather than +/// `num_nanoseconds()`, which returns `None` (and would export a silent +/// `0.000000000`) for magnitudes outside the `i64` nanosecond range. +fn duration_csv(duration: &chrono::Duration) -> String { + // chrono stores a negative duration as a negative whole-second part plus a + // non-negative subsecond remainder: -1.5s is secs=-2, nanos=500_000_000. + // The sign therefore has to be taken off before splitting into seconds and + // nanos — `num_seconds()` alone reports -2 for -1.5s, and pairing it with + // the remainder would render "-2.500000000". `duration` comes from + // `Duration::nanoseconds(endtime - starttime)`, so negating cannot overflow. + let (sign, magnitude) = if *duration < chrono::Duration::zero() { + ("-", -*duration) + } else { + ("", *duration) + }; + format!( + "{sign}{}.{:09}", + magnitude.num_seconds(), + magnitude.subsec_nanos() + ) +} + +/// SQLite binds `LIMIT` as a signed 64-bit integer, and a negative limit means +/// "unbounded". Converting a `u64` limit with `as` would wrap values above +/// `i64::MAX` to a negative number, turning the client's cap into an unbounded +/// export; saturate instead. +fn sql_limit(limit_opt: Option) -> i64 { + match limit_opt { + Some(l) => i64::try_from(l).unwrap_or(i64::MAX), + None => -1, + } +} + +fn event_field_value(event: &Event, key: &str) -> String { + match event.data.get(key) { + Some(serde_json::Value::String(s)) => s.clone(), + Some(v) => v.to_string(), + None => String::new(), + } +} + +fn write_csv_record( + writer: &mut impl Write, + fields: impl IntoIterator, +) -> Result<(), DatastoreError> { + let mut first = true; + for field in fields { + if !first { + writer.write_all(b",").map_err(csv_io_err)?; + } + first = false; + writer + .write_all(csv_escape(&field).as_bytes()) + .map_err(csv_io_err)?; + } + // RFC-4180 puts each record on its own line delimited by CRLF. + writer.write_all(b"\r\n").map_err(csv_io_err) +} + +/// Maximum number of `data` keys exported as individual CSV columns. +/// +/// A bucket whose events each carry a distinct key would otherwise emit +/// `events × keys` cells — quadratic in the event count, and enough to fill the +/// staging file. Above this cap the export keeps a single `data` column holding +/// each event's full data object as JSON: bounded, and still lossless. +const MAX_CSV_DATA_COLUMNS: usize = 32; + +/// Resolve the data columns for an export: the union of keys, or a single +/// `data` JSON column when that union exceeds `MAX_CSV_DATA_COLUMNS`. +fn csv_data_columns(key_order: Vec) -> (Vec, bool) { + if key_order.len() > MAX_CSV_DATA_COLUMNS { + (vec!["data".to_string()], true) + } else { + (key_order, false) + } +} + +fn write_csv_header(writer: &mut impl Write, data_keys: &[String]) -> Result<(), DatastoreError> { + let mut fields = vec![ + "id".to_string(), + "timestamp".to_string(), + "duration".to_string(), + ]; + fields.extend(data_keys.iter().cloned()); + write_csv_record(writer, fields) +} + +fn write_csv_event( + writer: &mut impl Write, + event: &Event, + data_keys: &[String], + data_json_column: bool, +) -> Result<(), DatastoreError> { + let mut fields = vec![ + event.id.map(|i| i.to_string()).unwrap_or_default(), + event.timestamp.to_rfc3339(), + duration_csv(&event.duration), + ]; + if data_json_column { + fields.push(serde_json::Value::Object(event.data.clone()).to_string()); + } else { + for key in data_keys { + fields.push(event_field_value(event, key)); + } + } + write_csv_record(writer, fields) +} + +/// Stream events for one bucket as RFC-4180 CSV, writing one row at a time. +/// +/// Columns: `id`, `timestamp`, `duration`, then the union of data keys across +/// all matched events, collected in a pre-pass so heterogeneous event data is +/// not truncated to the first event's schema. If the union exceeds +/// `MAX_CSV_DATA_COLUMNS`, a single `data` column holds each event's JSON data +/// object instead (bounded output, no keys dropped). +/// Query filters, clipping, and corrupt-row skipping match `get_events`. +pub(crate) fn write_events_csv( + conn: &Connection, + buckets: &HashMap, + bucket_id: &str, + starttime_opt: Option>, + endtime_opt: Option>, + limit_opt: Option, + mut writer: impl Write, +) -> Result<(), DatastoreError> { + let bucket = match buckets.get(bucket_id) { + Some(bucket) => bucket, + None => return Err(DatastoreError::NoSuchBucket(bucket_id.to_owned())), + }; + + let starttime_filter_ns: i64 = match starttime_opt { + Some(dt) => dt.timestamp_nanos_opt().unwrap(), + None => 0, + }; + let endtime_filter_ns: i64 = match endtime_opt { + Some(dt) => dt.timestamp_nanos_opt().unwrap(), + None => i64::MAX, + }; + if starttime_filter_ns > endtime_filter_ns { + warn!("Starttime in event query was lower than endtime!"); + write_csv_header(&mut writer, &[])?; + return writer.flush().map_err(csv_io_err); + } + let limit = sql_limit(limit_opt); + + let sql = if prefer_endtime_index(bucket, starttime_filter_ns, endtime_filter_ns, limit_opt) { + "SELECT id, starttime, endtime, data + FROM events INDEXED BY events_bucketrow_endtime_starttime_index + WHERE bucketrow = ?1 AND endtime >= ?2 AND starttime <= ?3 + ORDER BY starttime DESC, endtime ASC, id ASC LIMIT ?4" + } else { + "SELECT id, starttime, endtime, data + FROM events INDEXED BY events_bucketrow_starttime_endtime_index + WHERE bucketrow = ?1 AND endtime >= ?2 AND starttime <= ?3 + ORDER BY starttime DESC, endtime ASC, id ASC LIMIT ?4" + }; + // First pass: collect the union of data keys across matched rows so the + // header includes keys the first event may lack. Rows are still streamed + // in the second pass, and the key set is bounded by + // `MAX_CSV_DATA_COLUMNS`: collection stops as soon as the union exceeds + // the cap, because the fallback decision is already made at that point. + let keys_sql = sql.replace("SELECT id, starttime, endtime, data", "SELECT data"); + let mut key_order: Vec = Vec::new(); + let mut seen: HashSet = HashSet::new(); + let mut keys_stmt = conn.prepare_cached(&keys_sql).map_err(|err| { + DatastoreError::InternalError(format!("Failed to prepare CSV export key pass: {err}")) + })?; + let mut key_rows = keys_stmt + .query(rusqlite::params![ + bucket.bid.unwrap(), + starttime_filter_ns, + endtime_filter_ns, + limit, + ]) + .map_err(|err| { + DatastoreError::InternalError(format!("Failed to query CSV export key pass: {err}")) + })?; + let mut key_overflow = false; + while let Some(row) = key_rows.next().map_err(|err| { + DatastoreError::InternalError(format!("Failed to read CSV export key row: {err}")) + })? { + let data: Option = row.get(0).map_err(|err| { + DatastoreError::InternalError(format!("Failed to read CSV export key data: {err}")) + })?; + let data = match data { + Some(d) => d, + None => continue, + }; + if let Ok(value) = serde_json::from_str::(&data) { + if let Some(obj) = value.as_object() { + for key in obj.keys() { + if seen.insert(key.clone()) { + key_order.push(key.clone()); + if key_order.len() > MAX_CSV_DATA_COLUMNS { + key_overflow = true; + break; + } + } + } + } + } + if key_overflow { + break; + } + } + drop(key_rows); + drop(keys_stmt); + + let mut stmt = conn.prepare_cached(sql).map_err(|err| { + DatastoreError::InternalError(format!("Failed to prepare CSV export SQL: {err}")) + })?; + let mut rows = stmt + .query(rusqlite::params![ + bucket.bid.unwrap(), + starttime_filter_ns, + endtime_filter_ns, + limit, + ]) + .map_err(|err| { + DatastoreError::InternalError(format!("Failed to query CSV export SQL: {err}")) + })?; + + let clip = Some((starttime_filter_ns, endtime_filter_ns)); + let (data_keys, data_json_column) = csv_data_columns(key_order); + write_csv_header(&mut writer, &data_keys)?; + while let Some(row) = rows.next().map_err(|err| { + DatastoreError::InternalError(format!("Failed to read CSV export row: {err}")) + })? { + let event = match parse_event_row(row, clip) { + Ok(event) => event, + Err(err) => { + warn!("Corrupt event in bucket {}: {}", bucket_id, err); + continue; + } + }; + write_csv_event(&mut writer, &event, &data_keys, data_json_column)?; + } + writer.flush().map_err(csv_io_err) +} + #[cfg(test)] mod tests { use super::*; @@ -217,4 +492,231 @@ mod tests { Err(DatastoreError::InternalError(_)) )); } + + #[test] + fn csv_duration_preserves_sub_millisecond_nanos() { + assert_eq!( + duration_csv(&Duration::nanoseconds(1_500_000)), + "0.001500000" + ); + assert_eq!(duration_csv(&Duration::milliseconds(1)), "0.001000000"); + assert_eq!(duration_csv(&Duration::seconds(0)), "0.000000000"); + } + + #[test] + fn csv_duration_handles_negative_durations() { + // -1.5s is stored by chrono as secs=-2, nanos=500_000_000; taking the + // magnitude before splitting is what keeps this at -1.5 and not -2.5. + assert_eq!( + duration_csv(&Duration::milliseconds(-1_500)), + "-1.500000000" + ); + assert_eq!( + duration_csv(&Duration::nanoseconds(-1_500_000)), + "-0.001500000" + ); + assert_eq!(duration_csv(&Duration::seconds(-1)), "-1.000000000"); + } + + #[test] + fn csv_limit_saturates_instead_of_wrapping() { + assert_eq!(sql_limit(None), -1); + assert_eq!(sql_limit(Some(10)), 10); + assert_eq!(sql_limit(Some(i64::MAX as u64)), i64::MAX); + // Would have wrapped to -1 (SQLite: "no limit") under `as`. + assert_eq!(sql_limit(Some(u64::MAX)), i64::MAX); + } + + #[test] + fn csv_escape_neutralizes_formula_prefixes_and_quotes_rfc4180() { + assert_eq!(csv_escape("plain"), "plain"); + assert_eq!(csv_escape("=1+1"), "'=1+1"); + assert_eq!(csv_escape("+cmd"), "'+cmd"); + assert_eq!(csv_escape("-1+1"), "'-1+1"); + assert_eq!(csv_escape("@SUM(A1)"), "'@SUM(A1)"); + assert_eq!( + csv_escape("A \"quoted\" title"), + "\"A \"\"quoted\"\" title\"" + ); + assert_eq!(csv_escape("=1,2"), "\"'=1,2\""); + } + + #[test] + fn csv_escape_neutralizes_whitespace_prefixed_formulas() { + assert_eq!(csv_escape(" =1+1"), "' =1+1"); + assert_eq!(csv_escape("\t+cmd"), "'\t+cmd"); + assert_eq!(csv_escape(" plain"), " plain"); + } + + #[test] + fn streamed_csv_columns_use_union_of_keys_across_events() { + let (conn, mut ds) = setup(); + let events = [ + serde_json::json!({"app": "firefox", "title": "t"}), + serde_json::json!({"app": "firefox", "url": "u"}), + ] + .into_iter() + .map(|data| { + Event::new( + DateTime::from_timestamp(0, 0).unwrap(), + Duration::seconds(1), + serde_json::from_value(data).unwrap(), + ) + }) + .collect(); + ds.insert_events(&conn, "empty", events).unwrap(); + + let mut output = Vec::new(); + ds.write_events_csv(&conn, "empty", None, None, None, &mut output) + .unwrap(); + let csv = String::from_utf8(output).unwrap(); + let header = csv.lines().next().unwrap(); + assert!(header.ends_with(",app,title,url"), "header: {header}"); + assert!(csv.contains(",u"), "second row lost url: {csv}"); + } + + #[test] + fn csv_falls_back_to_a_json_data_column_for_wide_schemas() { + let (conn, mut ds) = setup(); + // One distinct key per event: the union grows with the event count. + let events: Vec = (0..MAX_CSV_DATA_COLUMNS + 1) + .map(|i| { + let mut data = serde_json::Map::new(); + data.insert(format!("k{i}"), serde_json::json!(i)); + Event::new( + DateTime::from_timestamp(i as i64, 0).unwrap(), + Duration::seconds(1), + data, + ) + }) + .collect(); + ds.insert_events(&conn, "empty", events).unwrap(); + + let mut output = Vec::new(); + ds.write_events_csv(&conn, "empty", None, None, None, &mut output) + .unwrap(); + let csv = String::from_utf8(output).unwrap(); + // Bounded header: one `data` column instead of one column per key. + assert_eq!(csv.lines().next().unwrap(), "id,timestamp,duration,data"); + assert_eq!(csv.lines().count(), MAX_CSV_DATA_COLUMNS + 2, "{csv}"); + // No keys are dropped: every event's data object survives as JSON. + assert!(csv.contains("k0"), "first key lost: {csv}"); + assert!(csv.contains("k32"), "last key lost: {csv}"); + } + + #[test] + fn streamed_csv_terminates_every_record_with_crlf() { + let (conn, mut ds) = setup(); + let event = Event::new( + DateTime::from_timestamp(0, 0).unwrap(), + Duration::seconds(1), + serde_json::from_value(serde_json::json!({"app": "firefox"})).unwrap(), + ); + ds.insert_events(&conn, "empty", vec![event]).unwrap(); + + let mut output = Vec::new(); + ds.write_events_csv(&conn, "empty", None, None, None, &mut output) + .unwrap(); + let csv = String::from_utf8(output).unwrap(); + // RFC-4180: header + one event row, each CRLF-terminated. + assert!(csv.ends_with("\r\n"), "{csv:?}"); + assert_eq!(csv.matches("\r\n").count(), 2, "{csv:?}"); + } + + #[test] + fn streamed_csv_stops_key_pass_at_the_cap_without_truncating_rows() { + let (conn, mut ds) = setup(); + // Far more distinct keys than the cap: the key pre-pass must stop + // collecting once the fallback is decided instead of holding every key. + let count = 1_000; + let events: Vec = (0..count) + .map(|i| { + let mut data = serde_json::Map::new(); + data.insert(format!("k{i}"), serde_json::json!(i)); + Event::new( + DateTime::from_timestamp(i as i64, 0).unwrap(), + Duration::seconds(1), + data, + ) + }) + .collect(); + ds.insert_events(&conn, "empty", events).unwrap(); + + let mut output = Vec::new(); + ds.write_events_csv(&conn, "empty", None, None, None, &mut output) + .unwrap(); + let csv = String::from_utf8(output).unwrap(); + assert_eq!(csv.lines().next().unwrap(), "id,timestamp,duration,data"); + // Leaving the key pre-pass early must not truncate the data pass. + assert_eq!(csv.lines().count(), count as usize + 1, "{csv}"); + } + + #[test] + fn csv_flush_errors_propagate() { + struct FlushFailingWriter; + impl Write for FlushFailingWriter { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + Ok(bytes.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Err(std::io::Error::new( + std::io::ErrorKind::WriteZero, + "disk full", + )) + } + } + let (conn, mut ds) = setup(); + let event = Event::new( + DateTime::from_timestamp(0, 0).unwrap(), + Duration::nanoseconds(1_500_000), + serde_json::from_value(serde_json::json!({"app": "firefox"})).unwrap(), + ); + ds.insert_events(&conn, "empty", vec![event]).unwrap(); + assert!(matches!( + ds.write_events_csv(&conn, "empty", None, None, None, FlushFailingWriter), + Err(DatastoreError::InternalError(msg)) if msg.contains("disk full") + )); + } + + #[test] + fn streamed_csv_preserves_precision_neutralizes_formulas_and_rejects_missing() { + let (conn, mut ds) = setup(); + let event = Event::new( + DateTime::from_timestamp(0, 0).unwrap(), + Duration::nanoseconds(1_500_000), + serde_json::from_value(serde_json::json!({ + "app": "firefox", + "title": "=cmd|calc" + })) + .unwrap(), + ); + ds.insert_events(&conn, "empty", vec![event]).unwrap(); + + let mut output = Vec::new(); + ds.write_events_csv(&conn, "empty", None, None, None, &mut output) + .unwrap(); + let csv = String::from_utf8(output).unwrap(); + assert!(csv.starts_with("id,timestamp,duration,"), "{csv}"); + assert!(csv.contains("0.001500000"), "duration: {csv}"); + assert!(csv.contains("'=cmd|calc"), "formula: {csv}"); + + let mut missing = Vec::new(); + assert!(matches!( + ds.write_events_csv(&conn, "missing", None, None, None, &mut missing), + Err(DatastoreError::NoSuchBucket(_)) + )); + assert!(missing.is_empty()); + } + + #[test] + fn streamed_csv_matches_get_events_row_count() { + let (conn, mut ds) = setup(); + let events = ds.get_events(&conn, "populated", None, None, None).unwrap(); + let mut output = Vec::new(); + ds.write_events_csv(&conn, "populated", None, None, None, &mut output) + .unwrap(); + let csv = String::from_utf8(output).unwrap(); + assert_eq!(csv.lines().count(), events.len() + 1, "{csv}"); + assert!(csv.contains("\"quotes \"\" and unicode ☀\""), "{csv}"); + } } diff --git a/aw-datastore/src/lib.rs b/aw-datastore/src/lib.rs index 1dd60961..522216f2 100644 --- a/aw-datastore/src/lib.rs +++ b/aw-datastore/src/lib.rs @@ -25,6 +25,7 @@ mod worker; pub use self::datastore::DatastoreInstance; pub use self::datastore::NEWEST_DB_VERSION; + pub use self::worker::Datastore; #[derive(Clone)] diff --git a/aw-datastore/src/worker.rs b/aw-datastore/src/worker.rs index d0c0959f..93b0300e 100644 --- a/aw-datastore/src/worker.rs +++ b/aw-datastore/src/worker.rs @@ -109,6 +109,7 @@ impl fmt::Debug for Datastore { #[derive(Debug)] pub enum Response { Export(File, Option), + ExportCsv(File), Empty(), Bucket(Bucket), BucketMap(HashMap), @@ -123,6 +124,13 @@ pub enum Response { #[derive(Debug)] pub enum Command { Export(Option, File), + ExportCsv( + String, + Option>, + Option>, + Option, + File, + ), CreateBucket(Bucket), DeleteBucket(String), GetBucket(String), @@ -413,6 +421,15 @@ impl DatastoreWorker { drop(writer); Ok(Response::Export(file, name)) } + Command::ExportCsv(bucket_id, start, end, limit, mut file) => { + let mut writer = BufWriter::new(&mut file); + ds.write_events_csv(tx, &bucket_id, start, end, limit, &mut writer)?; + writer.flush().map_err(|err| { + DatastoreError::InternalError(format!("Failed to flush CSV export: {err}")) + })?; + drop(writer); + Ok(Response::ExportCsv(file)) + } Command::CreateBucket(bucket) => match ds.create_bucket(tx, bucket) { Ok(_) => { self.commit = true; @@ -695,6 +712,31 @@ impl Datastore { } } + /// Stream one bucket's events as CSV into `file`, one SQL row at a time. + /// + /// Same snapshot rules as [`Datastore::export_to_file`]: the worker writes + /// including uncommitted events, and the file is returned only after + /// serialization and flush succeed. + pub fn export_csv_to_file( + &self, + bucket_id: &str, + start: Option>, + end: Option>, + limit: Option, + file: File, + ) -> Result { + match self.request(Command::ExportCsv( + bucket_id.to_owned(), + start, + end, + limit, + file, + ))? { + Response::ExportCsv(file) => Ok(file), + _ => panic!("Invalid response"), + } + } + pub fn insert_events( &self, bucket_id: &str, diff --git a/aw-server/src/endpoints/bucket.rs b/aw-server/src/endpoints/bucket.rs index 10e67e6d..ef9c0161 100644 --- a/aw-server/src/endpoints/bucket.rs +++ b/aw-server/src/endpoints/bucket.rs @@ -13,7 +13,7 @@ use aw_models::Event; use rocket::http::Status; use rocket::State; -use crate::endpoints::util::BucketsExportRocket; +use crate::endpoints::util::{BucketEventsCsvRocket, BucketsExportRocket}; use crate::endpoints::{HttpErrorJson, ServerState}; #[get("/")] @@ -223,6 +223,49 @@ pub fn bucket_export( BucketsExportRocket::new(&state.datastore, Some(bucket_id)) } +/// Stream events for a single bucket as a CSV file. +/// +/// Mirrors `GET //export` (JSON) but returns `text/csv`. +/// Sends HTTP headers before serialization begins so large buckets don't +/// look like a hung connection on Android WebView or slow networks. +/// Accepts the same `start`, `end`, and `limit` query params as the JSON +/// events endpoint. +#[get("//export/csv?&&")] +pub fn bucket_events_get_csv( + bucket_id: &str, + start: Option, + end: Option, + limit: Option, + state: &State, +) -> Result { + let starttime: Option> = match start { + Some(dt_str) => match DateTime::parse_from_rfc3339(&dt_str) { + Ok(dt) => Some(dt.with_timezone(&Utc)), + Err(e) => { + let err_msg = format!( + "Failed to parse starttime, datetime needs to be in rfc3339 format: {e}" + ); + warn!("{}", err_msg); + return Err(HttpErrorJson::new(Status::BadRequest, err_msg)); + } + }, + None => None, + }; + let endtime: Option> = match end { + Some(dt_str) => match DateTime::parse_from_rfc3339(&dt_str) { + Ok(dt) => Some(dt.with_timezone(&Utc)), + Err(e) => { + let err_msg = + format!("Failed to parse endtime, datetime needs to be in rfc3339 format: {e}"); + warn!("{}", err_msg); + return Err(HttpErrorJson::new(Status::BadRequest, err_msg)); + } + }, + None => None, + }; + BucketEventsCsvRocket::new(&state.datastore, bucket_id, starttime, endtime, limit) +} + #[delete("/")] pub fn bucket_delete(bucket_id: &str, state: &State) -> Result<(), HttpErrorJson> { let datastore = &state.datastore; diff --git a/aw-server/src/endpoints/mod.rs b/aw-server/src/endpoints/mod.rs index 276677e4..d0944e1f 100644 --- a/aw-server/src/endpoints/mod.rs +++ b/aw-server/src/endpoints/mod.rs @@ -179,6 +179,7 @@ pub fn build_rocket(server_state: ServerState, config: AWConfig) -> rocket::Rock bucket::buckets_get, bucket::bucket_get, bucket::bucket_events_get, + bucket::bucket_events_get_csv, bucket::bucket_events_create, bucket::bucket_events_heartbeat, bucket::bucket_event_count, diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index a96c583f..4eff492b 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -2,6 +2,7 @@ use std::fs::File; use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; use std::thread; +use chrono::{DateTime, Utc}; use rocket::http::ContentType; use rocket::http::Header; use rocket::http::Status; @@ -45,6 +46,35 @@ pub struct BucketsExportRocket { filename: String, } +/// Make a client-supplied bucket id safe to interpolate into a response header. +/// +/// Bucket ids are not restricted at creation, and Rocket percent-decodes path +/// segments, so an id containing CR/LF would otherwise split the +/// `Content-Disposition` header (response splitting / header injection). +fn sanitize_header_value(value: &str) -> String { + value + .chars() + .map(|c| match c { + c if c.is_control() => '_', + '"' | '\\' | ';' => '_', + c => c, + }) + .collect() +} + +/// Build a `Content-Disposition` value for an attachment download. +/// +/// The filename is quoted: bucket ids may legally contain spaces, and an +/// unquoted `filename=my bucket.csv` is a malformed parameter that clients may +/// drop. [`sanitize_header_value`] has already removed the characters that could +/// break out of the quoted string. +fn content_disposition(filename: &str) -> String { + format!( + "attachment; filename=\"{}\"", + sanitize_header_value(filename) + ) +} + fn export_filename( datastore: &aw_datastore::Datastore, bucket_id: Option<&str>, @@ -60,8 +90,8 @@ fn export_filename( } }; Ok(match name { - Some(id) => format!("attachment; filename=aw-bucket-export_{id}.json"), - None => "attachment; filename=aw-buckets-export.json".into(), + Some(id) => content_disposition(&format!("aw-bucket-export_{id}.json")), + None => content_disposition("aw-buckets-export.json"), }) } @@ -159,6 +189,117 @@ impl<'r> Responder<'r, 'static> for BucketsExportRocket { } } +// ── CSV streaming export ────────────────────────────────────────────────────── + +fn spawn_csv_export_stream( + datastore: aw_datastore::Datastore, + bucket_id: String, + start: Option>, + end: Option>, + limit: Option, + staging: File, + writer: PipeWriter, +) { + thread::spawn(move || { + // Serialize on the datastore worker, one SQL row at a time, into the + // staging file (created during preflight, before the 200 was + // committed). The full event set is never materialized in memory, and + // a flush failure (e.g. full staging filesystem) surfaces as an error + // instead of a silently truncated CSV. + let mut staging = match datastore.export_csv_to_file(&bucket_id, start, end, limit, staging) + { + Ok(file) => file, + Err(err) => { + error!("CSV export serialization failed: {err:?}"); + return; + } + }; + if let Err(err) = staging.seek(SeekFrom::Start(0)) { + error!("CSV staging rewind failed: {err}"); + return; + } + let mut writer = pipe_writer_to_file(writer); + if let Err(err) = copy(&mut staging, &mut writer) { + error!("CSV export copy failed: {err}"); + } + }); +} + +pub struct BucketEventsCsvRocket { + datastore: aw_datastore::Datastore, + bucket_id: String, + start: Option>, + end: Option>, + limit: Option, + filename: String, + staging: File, +} + +impl BucketEventsCsvRocket { + pub fn new( + datastore: &aw_datastore::Datastore, + bucket_id: &str, + start: Option>, + end: Option>, + limit: Option, + ) -> Result { + // Resolve 404/500 before headers commit. get_bucket catches a missing + // bucket; LIMIT 1 forces the same SQL the full export will run so a + // down worker or a prepare/read failure still returns JSON instead of + // a 200 with an empty CSV. The staging file is also created here: a + // full staging filesystem is still reportable as JSON at this point. + // Mid-stream failures after 200 cannot change the status without + // delaying headers until serialization finishes — that hung-connection + // behavior is what this endpoint exists to avoid (same tradeoff as + // JSON export / #721). + datastore.get_bucket(bucket_id)?; + datastore.get_events(bucket_id, start, end, Some(1))?; + let staging = tempfile::tempfile().map_err(|err| { + HttpErrorJson::new( + Status::InternalServerError, + format!("Failed to create CSV staging file: {err}"), + ) + })?; + let filename = content_disposition(&format!("aw-events-export-{bucket_id}.csv")); + Ok(Self { + datastore: datastore.clone(), + bucket_id: bucket_id.to_owned(), + start, + end, + limit, + filename, + staging, + }) + } +} + +impl<'r> Responder<'r, 'static> for BucketEventsCsvRocket { + fn respond_to(self, _: &Request) -> response::Result<'static> { + let Self { + datastore, + bucket_id, + start, + end, + limit, + filename, + staging, + } = self; + let (reader, writer) = pipe().map_err(|err| { + error!("Failed to open CSV export pipe: {err}"); + Status::InternalServerError + })?; + spawn_csv_export_stream(datastore, bucket_id, start, end, limit, staging, writer); + Response::build() + .status(Status::Ok) + .header(Header::new("Content-Disposition", filename)) + .header(ContentType::new("text", "csv")) + .streamed_body(rocket::tokio::fs::File::from_std(pipe_reader_to_file( + reader, + ))) + .ok() + } +} + use aw_datastore::DatastoreError; impl From for HttpErrorJson { @@ -193,3 +334,37 @@ impl From for HttpErrorJson { } } } + +#[cfg(test)] +mod tests { + use super::{content_disposition, sanitize_header_value}; + + #[test] + fn content_disposition_quotes_the_filename() { + // Spaces are legal in bucket ids; an unquoted filename= parameter with a + // space is malformed and clients may drop it. + assert_eq!( + content_disposition("aw-events-export-my bucket.csv"), + "attachment; filename=\"aw-events-export-my bucket.csv\"" + ); + assert_eq!( + content_disposition("aw-buckets-export.json"), + "attachment; filename=\"aw-buckets-export.json\"" + ); + } + + #[test] + fn sanitize_header_value_strips_header_metacharacters() { + assert_eq!( + sanitize_header_value("aw-watcher-window_host"), + "aw-watcher-window_host" + ); + // CR/LF would split the header; quote/backslash/semicolon would end or + // re-parameterize the filename value. + assert_eq!( + sanitize_header_value("evil\r\nX-Injected: 1"), + "evil__X-Injected: 1" + ); + assert_eq!(sanitize_header_value("a\"b\\c;d"), "a_b_c_d"); + } +} diff --git a/aw-server/tests/api.rs b/aw-server/tests/api.rs index 7a7b7914..712027f5 100644 --- a/aw-server/tests/api.rs +++ b/aw-server/tests/api.rs @@ -54,7 +54,7 @@ mod api_tests { assert_eq!(response.content_type(), Some(ContentType::JSON)); assert_eq!( response.headers().get_one("Content-Disposition"), - Some("attachment; filename=aw-bucket-export_live.json") + Some("attachment; filename=\"aw-bucket-export_live.json\"") ); let body: Value = serde_json::from_str(&response.into_string().unwrap()).unwrap(); assert_eq!( @@ -83,7 +83,7 @@ mod api_tests { assert_eq!(response.content_type(), Some(ContentType::JSON)); assert_eq!( response.headers().get_one("Content-Disposition"), - Some("attachment; filename=aw-buckets-export.json") + Some("attachment; filename=\"aw-buckets-export.json\"") ); let body: Value = serde_json::from_str(&response.into_string().unwrap()).unwrap(); assert_eq!(body["buckets"], json!({})); @@ -137,6 +137,70 @@ mod api_tests { ); } + #[test] + fn csv_export_returns_csv_with_correct_headers_and_missing_bucket_errors() { + let server = setup_testserver(); + let datastore = server + .state::() + .unwrap() + .datastore + .clone(); + let bucket: Bucket = serde_json::from_value(json!({ + "id": "testbucket", "type": "test", "client": "test", "hostname": "test" + })) + .unwrap(); + datastore.create_bucket(&bucket).unwrap(); + let mut event = aw_models::Event { + duration: chrono::Duration::nanoseconds(1_500_000), + ..Default::default() + }; + event.data.insert("app".into(), json!("firefox")); + event + .data + .insert("title".into(), json!("A \"quoted\" title")); + event.data.insert("formula".into(), json!("=cmd|calc")); + let inserted = datastore.insert_events("testbucket", &[event]).unwrap(); + + let client = Client::untracked(server).unwrap(); + let response = client + .get("/api/0/buckets/testbucket/export/csv") + .header(Header::new("Host", "127.0.0.1:5600")) + .dispatch(); + assert_eq!(response.status(), Status::Ok); + assert_eq!( + response.content_type(), + Some(ContentType::new("text", "csv")) + ); + assert_eq!( + response.headers().get_one("Content-Disposition"), + Some("attachment; filename=\"aw-events-export-testbucket.csv\"") + ); + let body = response.into_string().unwrap(); + // Header row + assert!(body.starts_with("id,timestamp,duration,"), "header: {body}"); + // Quoted field for title with embedded double-quote + assert!( + body.contains("\"A \"\"quoted\"\" title\""), + "quoting: {body}" + ); + // Sub-millisecond duration is not truncated to 0.001000000 + assert!(body.contains("0.001500000"), "duration: {body}"); + // Spreadsheet formula prefixes are neutralized + assert!(body.contains("'=cmd|calc"), "formula: {body}"); + // Event id present + let event_id = inserted[0].id.unwrap().to_string(); + assert!(body.contains(&event_id), "id in body: {body}"); + + // Missing bucket → 404 with JSON body + let response = client + .get("/api/0/buckets/nosuchbucket/export/csv") + .header(Header::new("Host", "127.0.0.1:5600")) + .dispatch(); + assert_eq!(response.status(), Status::NotFound); + let body: Value = serde_json::from_str(&response.into_string().unwrap()).unwrap(); + assert!(body["message"].as_str().unwrap().contains("does not exist")); + } + #[test] fn test_bucket() { let server = setup_testserver();