From d348d83712c3406420907a953ce0e827c762ad08 Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 23 Sep 2026 11:17:03 +0000 Subject: [PATCH 01/11] feat(server): add streaming CSV export endpoint for bucket events GET /api/0/buckets/{id}/export/csv streams events as CSV, sending HTTP headers before serialization so large buckets don't look like a hung connection on Android WebView or slow networks. Mirrors the existing streaming JSON export (BucketsExportRocket / #721): OS pipe + background thread serializes into a tempfile, then copies to the client. The route is at /export/csv (parallel to /export for JSON) to avoid any Rocket route collision with the single-event endpoint. RFC-4180 CSV: id, timestamp, duration (fractional seconds), then all data-map keys from the first event. Fields containing commas, quotes, or newlines are double-quoted and internal quotes are doubled. Test: csv_export_returns_csv_with_correct_headers_and_missing_bucket_errors Git-Session-Id: cdc6 --- aw-server/src/endpoints/bucket.rs | 45 ++++++++- aw-server/src/endpoints/mod.rs | 1 + aw-server/src/endpoints/util.rs | 155 +++++++++++++++++++++++++++++- aw-server/tests/api.rs | 56 +++++++++++ 4 files changed, 255 insertions(+), 2 deletions(-) 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..58f7284a 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -1,7 +1,8 @@ use std::fs::File; -use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; +use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom, Write}; use std::thread; +use chrono::{DateTime, Utc}; use rocket::http::ContentType; use rocket::http::Header; use rocket::http::Status; @@ -159,6 +160,158 @@ impl<'r> Responder<'r, 'static> for BucketsExportRocket { } } +// ── CSV streaming export ────────────────────────────────────────────────────── + +/// Minimal RFC-4180 CSV field escaper. +fn csv_escape(s: &str) -> String { + if s.contains([',', '"', '\n', '\r']) { + format!("\"{}\"", s.replace('"', "\"\"")) + } else { + s.to_owned() + } +} + +/// Serialize a slice of events as RFC-4180 CSV. +/// +/// Columns: `id`, `timestamp`, `duration`, then all keys from the first +/// event's `data` map (same schema the webui uses for client-side CSV). +fn events_to_csv(events: &[aw_models::Event]) -> String { + let data_keys: Vec = events + .first() + .map(|e| e.data.keys().cloned().collect()) + .unwrap_or_default(); + + let header: Vec<&str> = { + let mut h = vec!["id", "timestamp", "duration"]; + for k in &data_keys { + h.push(k.as_str()); + } + h + }; + + let mut out = header + .iter() + .map(|s| csv_escape(s)) + .collect::>() + .join(","); + out.push('\n'); + + for event in events { + let mut fields = vec![ + event.id.map(|i| i.to_string()).unwrap_or_default(), + event.timestamp.to_rfc3339(), + // duration is serialized as fractional seconds (same as JSON) + format!("{:.9}", event.duration.num_milliseconds() as f64 / 1000.0), + ]; + for key in &data_keys { + let val = match event.data.get(key) { + Some(serde_json::Value::String(s)) => s.clone(), + Some(v) => v.to_string(), + None => String::new(), + }; + fields.push(val); + } + out.push_str( + &fields + .iter() + .map(|s| csv_escape(s)) + .collect::>() + .join(","), + ); + out.push('\n'); + } + out +} + +pub struct BucketEventsCsvRocket { + datastore: aw_datastore::Datastore, + bucket_id: String, + start: Option>, + end: Option>, + limit: Option, + filename: String, +} + +impl BucketEventsCsvRocket { + pub fn new( + datastore: &aw_datastore::Datastore, + bucket_id: &str, + start: Option>, + end: Option>, + limit: Option, + ) -> Result { + // 404-check the bucket before headers are sent so errors return + // a proper JSON body rather than truncating mid-stream. + datastore.get_bucket(bucket_id)?; + let filename = format!("attachment; filename=aw-events-export-{bucket_id}.csv"); + Ok(Self { + datastore: datastore.clone(), + bucket_id: bucket_id.to_owned(), + start, + end, + limit, + filename, + }) + } +} + +impl<'r> Responder<'r, 'static> for BucketEventsCsvRocket { + fn respond_to(self, _: &Request) -> response::Result<'static> { + let Self { + datastore, + bucket_id, + start, + end, + limit, + filename, + } = self; + let (reader, writer) = pipe().map_err(|err| { + error!("Failed to open CSV export pipe: {err}"); + Status::InternalServerError + })?; + // Serialize on the datastore worker into a private tempfile, then copy + // to the client pipe. Headers are sent immediately; the download starts + // before the CSV is fully built. + thread::spawn(move || { + let events = match datastore.get_events(&bucket_id, start, end, limit) { + Ok(e) => e, + Err(err) => { + error!("CSV export: get_events failed: {err:?}"); + return; + } + }; + let csv = events_to_csv(&events); + let mut staging = match tempfile::tempfile() { + Ok(f) => f, + Err(err) => { + error!("Failed to create CSV staging file: {err}"); + return; + } + }; + if let Err(err) = staging.write_all(csv.as_bytes()) { + error!("CSV staging write 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}"); + } + }); + 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 { diff --git a/aw-server/tests/api.rs b/aw-server/tests/api.rs index 7a7b7914..761c97ca 100644 --- a/aw-server/tests/api.rs +++ b/aw-server/tests/api.rs @@ -137,6 +137,62 @@ 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::default(); + event.data.insert("app".into(), json!("firefox")); + event + .data + .insert("title".into(), json!("A \"quoted\" title")); + 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}" + ); + // 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(); From 5789a9fbfec0b7c4ffe4ef79d894e7efe7a319c5 Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 23 Sep 2026 11:41:16 +0000 Subject: [PATCH 02/11] fix(export): stream CSV rows, preserve ns duration, neutralize formulas Address Greptile P1s on ActivityWatch/aw-server-rust#722: - Write CSV from SQL row-by-row on the datastore worker (no full Vec + String) - Format duration from nanoseconds so 0.0015s is not truncated to 0.001 - Neutralize spreadsheet formula prefixes (=, +, -, @) - Preflight LIMIT 1 before 200 so worker/SQL failures still return JSON Git-Session-Id: 79e64905-13b9-5822-b880-3863140ba2d7 --- aw-datastore/src/datastore.rs | 30 +++- aw-datastore/src/export.rs | 238 +++++++++++++++++++++++++++++++- aw-datastore/src/worker.rs | 41 ++++++ aw-server/src/endpoints/util.rs | 136 ++++++------------ aw-server/tests/api.rs | 10 +- 5 files changed, 360 insertions(+), 95 deletions(-) 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..6affd019 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -1,12 +1,14 @@ use std::{collections::HashMap, 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 +124,174 @@ 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. +fn neutralize_formula(s: &str) -> String { + match s.chars().next() { + Some('=' | '+' | '-' | '@' | '\t' | '\r') => 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). +fn duration_csv(duration: &chrono::Duration) -> String { + let ns = duration.num_nanoseconds().unwrap_or(0); + let sign = if ns < 0 { "-" } else { "" }; + let ns = ns.unsigned_abs(); + format!("{sign}{}.{:09}", ns / 1_000_000_000, ns % 1_000_000_000) +} + +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)?; + } + writer.write_all(b"\n").map_err(csv_io_err) +} + +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], +) -> Result<(), DatastoreError> { + let mut fields = vec![ + event.id.map(|i| i.to_string()).unwrap_or_default(), + event.timestamp.to_rfc3339(), + duration_csv(&event.duration), + ]; + 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 all keys from the first +/// valid event's `data` map (same schema the webui uses for client-side CSV). +/// 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!"); + return write_csv_header(&mut writer, &[]); + } + let limit = match limit_opt { + Some(l) => l as i64, + None => -1, + }; + + 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" + }; + 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 mut data_keys: Option> = None; + 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; + } + }; + if data_keys.is_none() { + let keys: Vec = event.data.keys().cloned().collect(); + write_csv_header(&mut writer, &keys)?; + data_keys = Some(keys); + } + write_csv_event(&mut writer, &event, data_keys.as_ref().unwrap())?; + } + if data_keys.is_none() { + write_csv_header(&mut writer, &[])?; + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; @@ -217,4 +387,70 @@ 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_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 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/worker.rs b/aw-datastore/src/worker.rs index d0c0959f..a86cc9c0 100644 --- a/aw-datastore/src/worker.rs +++ b/aw-datastore/src/worker.rs @@ -123,6 +123,13 @@ pub enum Response { #[derive(Debug)] pub enum Command { Export(Option, File), + ExportEventsCsv( + String, + Option>, + Option>, + Option, + File, + ), CreateBucket(Bucket), DeleteBucket(String), GetBucket(String), @@ -413,6 +420,15 @@ impl DatastoreWorker { drop(writer); Ok(Response::Export(file, name)) } + Command::ExportEventsCsv(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::Export(file, Some(bucket_id))) + } Command::CreateBucket(bucket) => match ds.create_bucket(tx, bucket) { Ok(_) => { self.commit = true; @@ -695,6 +711,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_events_csv_to_file( + &self, + bucket_id: &str, + starttime_opt: Option>, + endtime_opt: Option>, + limit_opt: Option, + file: File, + ) -> Result { + match self.request(Command::ExportEventsCsv( + bucket_id.to_string(), + starttime_opt, + endtime_opt, + limit_opt, + file, + ))? { + Response::Export(file, _) => Ok(file), + _ => panic!("Invalid response"), + } + } + pub fn insert_events( &self, bucket_id: &str, diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index 58f7284a..94c49df7 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -1,5 +1,5 @@ use std::fs::File; -use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom, Write}; +use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; use std::thread; use chrono::{DateTime, Utc}; @@ -162,65 +162,42 @@ impl<'r> Responder<'r, 'static> for BucketsExportRocket { // ── CSV streaming export ────────────────────────────────────────────────────── -/// Minimal RFC-4180 CSV field escaper. -fn csv_escape(s: &str) -> String { - if s.contains([',', '"', '\n', '\r']) { - format!("\"{}\"", s.replace('"', "\"\"")) - } else { - s.to_owned() - } -} - -/// Serialize a slice of events as RFC-4180 CSV. -/// -/// Columns: `id`, `timestamp`, `duration`, then all keys from the first -/// event's `data` map (same schema the webui uses for client-side CSV). -fn events_to_csv(events: &[aw_models::Event]) -> String { - let data_keys: Vec = events - .first() - .map(|e| e.data.keys().cloned().collect()) - .unwrap_or_default(); - - let header: Vec<&str> = { - let mut h = vec!["id", "timestamp", "duration"]; - for k in &data_keys { - h.push(k.as_str()); - } - h - }; - - let mut out = header - .iter() - .map(|s| csv_escape(s)) - .collect::>() - .join(","); - out.push('\n'); - - for event in events { - let mut fields = vec![ - event.id.map(|i| i.to_string()).unwrap_or_default(), - event.timestamp.to_rfc3339(), - // duration is serialized as fractional seconds (same as JSON) - format!("{:.9}", event.duration.num_milliseconds() as f64 / 1000.0), - ]; - for key in &data_keys { - let val = match event.data.get(key) { - Some(serde_json::Value::String(s)) => s.clone(), - Some(v) => v.to_string(), - None => String::new(), +fn spawn_csv_export_stream( + datastore: aw_datastore::Datastore, + bucket_id: String, + start: Option>, + end: Option>, + limit: Option, + writer: PipeWriter, +) { + thread::spawn(move || { + let staging = match tempfile::tempfile() { + Ok(file) => file, + Err(err) => { + error!("Failed to create CSV staging file: {err}"); + return; + } + }; + // Worker writes CSV rows incrementally from SQL (no full event Vec + // and no second full-size String). Same tempfile-then-copy pattern + // as JSON export: a slow download must not stall heartbeats. + let mut staging = + match datastore.export_events_csv_to_file(&bucket_id, start, end, limit, staging) { + Ok(file) => file, + Err(err) => { + error!("CSV export stream failed: {err:?}"); + return; + } }; - fields.push(val); + if let Err(err) = staging.seek(SeekFrom::Start(0)) { + error!("CSV staging rewind failed: {err}"); + return; } - out.push_str( - &fields - .iter() - .map(|s| csv_escape(s)) - .collect::>() - .join(","), - ); - out.push('\n'); - } - out + 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 { @@ -240,9 +217,15 @@ impl BucketEventsCsvRocket { end: Option>, limit: Option, ) -> Result { - // 404-check the bucket before headers are sent so errors return - // a proper JSON body rather than truncating mid-stream. + // 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. 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 filename = format!("attachment; filename=aw-events-export-{bucket_id}.csv"); Ok(Self { datastore: datastore.clone(), @@ -269,38 +252,7 @@ impl<'r> Responder<'r, 'static> for BucketEventsCsvRocket { error!("Failed to open CSV export pipe: {err}"); Status::InternalServerError })?; - // Serialize on the datastore worker into a private tempfile, then copy - // to the client pipe. Headers are sent immediately; the download starts - // before the CSV is fully built. - thread::spawn(move || { - let events = match datastore.get_events(&bucket_id, start, end, limit) { - Ok(e) => e, - Err(err) => { - error!("CSV export: get_events failed: {err:?}"); - return; - } - }; - let csv = events_to_csv(&events); - let mut staging = match tempfile::tempfile() { - Ok(f) => f, - Err(err) => { - error!("Failed to create CSV staging file: {err}"); - return; - } - }; - if let Err(err) = staging.write_all(csv.as_bytes()) { - error!("CSV staging write 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}"); - } - }); + spawn_csv_export_stream(datastore, bucket_id, start, end, limit, writer); Response::build() .status(Status::Ok) .header(Header::new("Content-Disposition", filename)) diff --git a/aw-server/tests/api.rs b/aw-server/tests/api.rs index 761c97ca..3b000ad2 100644 --- a/aw-server/tests/api.rs +++ b/aw-server/tests/api.rs @@ -150,11 +150,15 @@ mod api_tests { })) .unwrap(); datastore.create_bucket(&bucket).unwrap(); - let mut event = aw_models::Event::default(); + 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(); @@ -179,6 +183,10 @@ mod api_tests { 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}"); From 7962992c0254fb8dd9a2ee2dc73b7e71de1b3eca Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 23 Sep 2026 13:09:48 +0000 Subject: [PATCH 03/11] fix(export): move CSV serialization off datastore worker MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The ExportEventsCsv worker command blocked the single shared datastore worker for the entire duration of the export — heartbeats and all other requests queued until serialization finished. Fix: remove ExportEventsCsv from the worker. The background thread now calls get_events() (worker holds the DB lock only for the SQL read), then serializes CSV via write_csv_from_events() off-worker. The worker is free for other requests during the write. write_csv_from_events() is the Vec-based counterpart to the connection-based write_events_csv(); both produce identical RFC-4180 output with ns-precision durations and formula neutralization. Git-Session-Id: bb28 --- aw-datastore/src/export.rs | 23 +++++++++++++++++++++ aw-datastore/src/lib.rs | 1 + aw-datastore/src/worker.rs | 36 --------------------------------- aw-server/src/endpoints/util.rs | 30 +++++++++++++++------------ 4 files changed, 41 insertions(+), 49 deletions(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index 6affd019..9d8a3920 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -206,6 +206,29 @@ fn write_csv_event( write_csv_record(writer, fields) } +/// Write an already-fetched event slice as RFC-4180 CSV. +/// +/// Unlike `write_events_csv`, this does not need a database connection — +/// call it from a background thread after fetching events via `get_events` +/// so the datastore worker is free during the (potentially long) serialization. +/// +/// Columns: `id`, `timestamp`, `duration`, then all keys from the first +/// event's data map. +pub fn write_csv_from_events( + events: &[aw_models::Event], + mut writer: impl Write, +) -> Result<(), DatastoreError> { + let data_keys: Vec = events + .first() + .map(|e| e.data.keys().cloned().collect()) + .unwrap_or_default(); + write_csv_header(&mut writer, &data_keys)?; + for event in events { + write_csv_event(&mut writer, event, &data_keys)?; + } + Ok(()) +} + /// Stream events for one bucket as RFC-4180 CSV, writing one row at a time. /// /// Columns: `id`, `timestamp`, `duration`, then all keys from the first diff --git a/aw-datastore/src/lib.rs b/aw-datastore/src/lib.rs index 1dd60961..b4bbe397 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::export::write_csv_from_events; pub use self::worker::Datastore; #[derive(Clone)] diff --git a/aw-datastore/src/worker.rs b/aw-datastore/src/worker.rs index a86cc9c0..3495f192 100644 --- a/aw-datastore/src/worker.rs +++ b/aw-datastore/src/worker.rs @@ -123,13 +123,6 @@ pub enum Response { #[derive(Debug)] pub enum Command { Export(Option, File), - ExportEventsCsv( - String, - Option>, - Option>, - Option, - File, - ), CreateBucket(Bucket), DeleteBucket(String), GetBucket(String), @@ -420,15 +413,6 @@ impl DatastoreWorker { drop(writer); Ok(Response::Export(file, name)) } - Command::ExportEventsCsv(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::Export(file, Some(bucket_id))) - } Command::CreateBucket(bucket) => match ds.create_bucket(tx, bucket) { Ok(_) => { self.commit = true; @@ -716,26 +700,6 @@ impl Datastore { /// 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_events_csv_to_file( - &self, - bucket_id: &str, - starttime_opt: Option>, - endtime_opt: Option>, - limit_opt: Option, - file: File, - ) -> Result { - match self.request(Command::ExportEventsCsv( - bucket_id.to_string(), - starttime_opt, - endtime_opt, - limit_opt, - file, - ))? { - Response::Export(file, _) => Ok(file), - _ => panic!("Invalid response"), - } - } - pub fn insert_events( &self, bucket_id: &str, diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index 94c49df7..96111df1 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -1,5 +1,5 @@ use std::fs::File; -use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; +use std::io::{copy, pipe, BufWriter, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; use std::thread; use chrono::{DateTime, Utc}; @@ -171,24 +171,28 @@ fn spawn_csv_export_stream( writer: PipeWriter, ) { thread::spawn(move || { - let staging = match tempfile::tempfile() { + // Fetch events via the worker. The worker holds the DB lock only for + // this SQL read, then is immediately free for heartbeats and other + // requests while CSV serialization runs here, off-worker. + let events = match datastore.get_events(&bucket_id, start, end, limit) { + Ok(events) => events, + Err(err) => { + error!("CSV export stream failed fetching events: {err:?}"); + return; + } + }; + let mut staging = match tempfile::tempfile() { Ok(file) => file, Err(err) => { error!("Failed to create CSV staging file: {err}"); return; } }; - // Worker writes CSV rows incrementally from SQL (no full event Vec - // and no second full-size String). Same tempfile-then-copy pattern - // as JSON export: a slow download must not stall heartbeats. - let mut staging = - match datastore.export_events_csv_to_file(&bucket_id, start, end, limit, staging) { - Ok(file) => file, - Err(err) => { - error!("CSV export stream failed: {err:?}"); - return; - } - }; + if let Err(err) = aw_datastore::write_csv_from_events(&events, BufWriter::new(&mut staging)) + { + error!("CSV export serialization failed: {err:?}"); + return; + } if let Err(err) = staging.seek(SeekFrom::Start(0)) { error!("CSV staging rewind failed: {err}"); return; From 09f23911d466c7dc099a07d0d544be003a1ff05f Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 23 Sep 2026 13:43:39 +0000 Subject: [PATCH 04/11] fix(export): surface CSV staging flush errors before copy BufWriter::drop swallows a final flush failure, so a full staging filesystem still rewound and copied a truncated CSV under 200 OK. Keep the writer in a local, flush explicitly, and flush inside write_csv_from_events / write_events_csv so the error propagates. Git-Session-Id: 01a0ce79 --- aw-datastore/src/export.rs | 32 +++++++++++++++++++++++++++++--- aw-server/src/endpoints/util.rs | 16 ++++++++++++---- 2 files changed, 41 insertions(+), 7 deletions(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index 9d8a3920..f1296974 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -226,7 +226,7 @@ pub fn write_csv_from_events( for event in events { write_csv_event(&mut writer, event, &data_keys)?; } - Ok(()) + writer.flush().map_err(csv_io_err) } /// Stream events for one bucket as RFC-4180 CSV, writing one row at a time. @@ -258,7 +258,8 @@ pub(crate) fn write_events_csv( }; if starttime_filter_ns > endtime_filter_ns { warn!("Starttime in event query was lower than endtime!"); - return write_csv_header(&mut writer, &[]); + write_csv_header(&mut writer, &[])?; + return writer.flush().map_err(csv_io_err); } let limit = match limit_opt { Some(l) => l as i64, @@ -312,7 +313,7 @@ pub(crate) fn write_events_csv( if data_keys.is_none() { write_csv_header(&mut writer, &[])?; } - Ok(()) + writer.flush().map_err(csv_io_err) } #[cfg(test)] @@ -435,6 +436,31 @@ mod tests { assert_eq!(csv_escape("=1,2"), "\"'=1,2\""); } + #[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 event = Event::new( + DateTime::from_timestamp(0, 0).unwrap(), + Duration::nanoseconds(1_500_000), + serde_json::from_value(serde_json::json!({"app": "firefox"})).unwrap(), + ); + assert!(matches!( + write_csv_from_events(&[event], 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(); diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index 96111df1..0419e4bf 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -1,5 +1,5 @@ use std::fs::File; -use std::io::{copy, pipe, BufWriter, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; +use std::io::{copy, pipe, BufWriter, Cursor, PipeReader, PipeWriter, Seek, SeekFrom, Write}; use std::thread; use chrono::{DateTime, Utc}; @@ -188,10 +188,18 @@ fn spawn_csv_export_stream( return; } }; - if let Err(err) = aw_datastore::write_csv_from_events(&events, BufWriter::new(&mut staging)) { - error!("CSV export serialization failed: {err:?}"); - return; + let mut csv_writer = BufWriter::new(&mut staging); + if let Err(err) = aw_datastore::write_csv_from_events(&events, &mut csv_writer) { + error!("CSV export serialization failed: {err:?}"); + return; + } + // BufWriter::drop swallows flush errors. Detect a full staging + // filesystem (or similar) before we rewind and copy a truncated CSV. + if let Err(err) = csv_writer.flush() { + error!("CSV export flush failed: {err}"); + return; + } } if let Err(err) = staging.seek(SeekFrom::Start(0)) { error!("CSV staging rewind failed: {err}"); From 4655638093cf9fc0816693774a19da8ca3ba844d Mon Sep 17 00:00:00 2001 From: TimeToBuildBob Date: Fri, 25 Sep 2026 16:18:26 +0000 Subject: [PATCH 05/11] fix(export): neutralize whitespace-prefixed formulas; union CSV keys across events Git-Session-Id: 4a49419b-e1d7-5734-b196-44aca96d17ff --- aw-datastore/src/export.rs | 153 ++++++++++++++++++++++++++++++++----- 1 file changed, 132 insertions(+), 21 deletions(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index f1296974..9ec77bf3 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -1,4 +1,5 @@ -use std::{collections::HashMap, io::Write}; +use std::collections::{HashMap, HashSet}; +use std::io::Write; use aw_models::{Bucket, Event}; use chrono::{DateTime, Utc}; @@ -129,9 +130,16 @@ fn csv_io_err(err: std::io::Error) -> DatastoreError { } /// 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 { - match s.chars().next() { - Some('=' | '+' | '-' | '@' | '\t' | '\r') => format!("'{s}"), + let trimmed = s.trim_start_matches(FORMULA_LEADING_WHITESPACE); + match trimmed.chars().next() { + Some('=' | '+' | '-' | '@') => format!("'{s}"), _ => s.to_owned(), } } @@ -212,16 +220,24 @@ fn write_csv_event( /// call it from a background thread after fetching events via `get_events` /// so the datastore worker is free during the (potentially long) serialization. /// -/// Columns: `id`, `timestamp`, `duration`, then all keys from the first -/// event's data map. +/// Columns: `id`, `timestamp`, `duration`, then the union of data keys across +/// all events (heterogeneous event data keeps every key's values). pub fn write_csv_from_events( events: &[aw_models::Event], mut writer: impl Write, ) -> Result<(), DatastoreError> { - let data_keys: Vec = events - .first() - .map(|e| e.data.keys().cloned().collect()) - .unwrap_or_default(); + // Union of keys across all events, so a later event's keys are not + // silently dropped just because the first event lacked them. + let mut key_order: Vec = Vec::new(); + let mut seen: HashSet<&str> = HashSet::new(); + for event in events { + for key in event.data.keys() { + if seen.insert(key.as_str()) { + key_order.push(key.clone()); + } + } + } + let data_keys = key_order; write_csv_header(&mut writer, &data_keys)?; for event in events { write_csv_event(&mut writer, event, &data_keys)?; @@ -231,8 +247,9 @@ pub fn write_csv_from_events( /// Stream events for one bucket as RFC-4180 CSV, writing one row at a time. /// -/// Columns: `id`, `timestamp`, `duration`, then all keys from the first -/// valid event's `data` map (same schema the webui uses for client-side CSV). +/// 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. /// Query filters, clipping, and corrupt-row skipping match `get_events`. pub(crate) fn write_events_csv( conn: &Connection, @@ -277,6 +294,48 @@ pub(crate) fn write_events_csv( 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; only the (small) key set is held in memory. + 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}")) + })?; + 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()); + } + } + } + } + } + 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}")) })?; @@ -292,7 +351,8 @@ pub(crate) fn write_events_csv( })?; let clip = Some((starttime_filter_ns, endtime_filter_ns)); - let mut data_keys: Option> = None; + let data_keys = 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}")) })? { @@ -303,15 +363,7 @@ pub(crate) fn write_events_csv( continue; } }; - if data_keys.is_none() { - let keys: Vec = event.data.keys().cloned().collect(); - write_csv_header(&mut writer, &keys)?; - data_keys = Some(keys); - } - write_csv_event(&mut writer, &event, data_keys.as_ref().unwrap())?; - } - if data_keys.is_none() { - write_csv_header(&mut writer, &[])?; + write_csv_event(&mut writer, &event, &data_keys)?; } writer.flush().map_err(csv_io_err) } @@ -436,6 +488,65 @@ mod tests { 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 csv_columns_use_union_of_keys_across_events() { + let events: Vec = [ + serde_json::json!({"a": 1, "b": 2}), + serde_json::json!({"a": 3, "c": 4}), + ] + .into_iter() + .map(|data| { + Event::new( + DateTime::from_timestamp(0, 0).unwrap(), + Duration::seconds(1), + serde_json::from_value(data).unwrap(), + ) + }) + .collect(); + let mut output = Vec::new(); + write_csv_from_events(&events, &mut output).unwrap(); + let csv = String::from_utf8(output).unwrap(); + let header = csv.lines().next().unwrap(); + assert!(header.contains(",a,"), "header: {header}"); + assert!(header.ends_with(",b,c"), "header: {header}"); + // Second event keeps its `c` value instead of losing it. + assert!(csv.contains(",4"), "second row lost c: {csv}"); + } + + #[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_flush_errors_propagate() { struct FlushFailingWriter; From bec0d945fd8571fd99c98c6312160fef2ec1dec5 Mon Sep 17 00:00:00 2001 From: Bob Date: Fri, 25 Sep 2026 18:40:40 +0000 Subject: [PATCH 06/11] fix(export): stream CSV export rows on the worker, bound column count MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The CSV export endpoint fetched every matching event into a `Vec` via `get_events`, serialized that slice into a second full-size CSV buffer, and then copied the staging file to the client. For the 500k-event buckets this endpoint exists to serve, the full event set plus the serialized CSV were both live in memory at once — the server-side OOM the endpoint was meant to remove, moved from the webview to the server. It also built its column set from the union of every event's data keys, so a bucket with many distinct keys produced `events × keys` cells (quadratic when keys grow with events) written to disk before any body bytes were sent. Changes: - `aw-datastore`: add `Command::ExportCsv` / `Response::ExportCsv` and `Datastore::export_csv_to_file`, which run the existing row-at-a-time `write_events_csv` on the worker straight into the staging file. - `aw-server`: `spawn_csv_export_stream` now uses that instead of `get_events` + `write_csv_from_events`, so peak memory is O(distinct keys) for the header pre-pass rather than O(events). The staging-file pattern (and its slow/dropped-download protection) is unchanged and now matches the JSON export. - `aw-datastore`: cap the union of data keys at `MAX_CSV_DATA_COLUMNS` (32). Above the cap the export emits `id,timestamp,duration,data` with each event's full data object as JSON, bounding the worst case while still dropping no keys. - Remove `write_csv_from_events` (the materializing variant) so the HTTP path cannot regress to it; its tests now exercise `write_events_csv`. Verified: `cargo test -p aw-datastore --lib export` 10 passed (incl. new `csv_falls_back_to_a_json_data_column_for_wide_schemas`), `cargo test -p aw-server --test api csv_export` passed, `cargo fmt --check` and `cargo clippy` clean. Git-Session-Id: 366e2805-96f6-5ba1-a390-5e3501978b44 --- aw-datastore/src/export.rs | 124 ++++++++++++++++---------------- aw-datastore/src/lib.rs | 2 +- aw-datastore/src/worker.rs | 37 ++++++++++ aw-server/src/endpoints/util.rs | 31 +++----- 4 files changed, 110 insertions(+), 84 deletions(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index 9ec77bf3..1a684046 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -188,6 +188,24 @@ fn write_csv_record( writer.write_all(b"\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(), @@ -202,54 +220,30 @@ 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), ]; - for key in data_keys { - fields.push(event_field_value(event, key)); - } - write_csv_record(writer, fields) -} - -/// Write an already-fetched event slice as RFC-4180 CSV. -/// -/// Unlike `write_events_csv`, this does not need a database connection — -/// call it from a background thread after fetching events via `get_events` -/// so the datastore worker is free during the (potentially long) serialization. -/// -/// Columns: `id`, `timestamp`, `duration`, then the union of data keys across -/// all events (heterogeneous event data keeps every key's values). -pub fn write_csv_from_events( - events: &[aw_models::Event], - mut writer: impl Write, -) -> Result<(), DatastoreError> { - // Union of keys across all events, so a later event's keys are not - // silently dropped just because the first event lacked them. - let mut key_order: Vec = Vec::new(); - let mut seen: HashSet<&str> = HashSet::new(); - for event in events { - for key in event.data.keys() { - if seen.insert(key.as_str()) { - key_order.push(key.clone()); - } + 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)); } } - let data_keys = key_order; - write_csv_header(&mut writer, &data_keys)?; - for event in events { - write_csv_event(&mut writer, event, &data_keys)?; - } - writer.flush().map_err(csv_io_err) + 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. +/// 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, @@ -351,7 +345,7 @@ pub(crate) fn write_events_csv( })?; let clip = Some((starttime_filter_ns, endtime_filter_ns)); - let data_keys = key_order; + 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}")) @@ -363,7 +357,7 @@ pub(crate) fn write_events_csv( continue; } }; - write_csv_event(&mut writer, &event, &data_keys)?; + write_csv_event(&mut writer, &event, &data_keys, data_json_column)?; } writer.flush().map_err(csv_io_err) } @@ -495,31 +489,6 @@ mod tests { assert_eq!(csv_escape(" plain"), " plain"); } - #[test] - fn csv_columns_use_union_of_keys_across_events() { - let events: Vec = [ - serde_json::json!({"a": 1, "b": 2}), - serde_json::json!({"a": 3, "c": 4}), - ] - .into_iter() - .map(|data| { - Event::new( - DateTime::from_timestamp(0, 0).unwrap(), - Duration::seconds(1), - serde_json::from_value(data).unwrap(), - ) - }) - .collect(); - let mut output = Vec::new(); - write_csv_from_events(&events, &mut output).unwrap(); - let csv = String::from_utf8(output).unwrap(); - let header = csv.lines().next().unwrap(); - assert!(header.contains(",a,"), "header: {header}"); - assert!(header.ends_with(",b,c"), "header: {header}"); - // Second event keeps its `c` value instead of losing it. - assert!(csv.contains(",4"), "second row lost c: {csv}"); - } - #[test] fn streamed_csv_columns_use_union_of_keys_across_events() { let (conn, mut ds) = setup(); @@ -547,6 +516,35 @@ mod tests { 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 csv_flush_errors_propagate() { struct FlushFailingWriter; @@ -561,13 +559,15 @@ mod tests { )) } } + 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!( - write_csv_from_events(&[event], FlushFailingWriter), + ds.write_events_csv(&conn, "empty", None, None, None, FlushFailingWriter), Err(DatastoreError::InternalError(msg)) if msg.contains("disk full") )); } diff --git a/aw-datastore/src/lib.rs b/aw-datastore/src/lib.rs index b4bbe397..522216f2 100644 --- a/aw-datastore/src/lib.rs +++ b/aw-datastore/src/lib.rs @@ -25,7 +25,7 @@ mod worker; pub use self::datastore::DatastoreInstance; pub use self::datastore::NEWEST_DB_VERSION; -pub use self::export::write_csv_from_events; + pub use self::worker::Datastore; #[derive(Clone)] diff --git a/aw-datastore/src/worker.rs b/aw-datastore/src/worker.rs index 3495f192..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; @@ -700,6 +717,26 @@ impl Datastore { /// 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/util.rs b/aw-server/src/endpoints/util.rs index 0419e4bf..09d65ec7 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -1,5 +1,5 @@ use std::fs::File; -use std::io::{copy, pipe, BufWriter, Cursor, PipeReader, PipeWriter, Seek, SeekFrom, Write}; +use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; use std::thread; use chrono::{DateTime, Utc}; @@ -171,36 +171,25 @@ fn spawn_csv_export_stream( writer: PipeWriter, ) { thread::spawn(move || { - // Fetch events via the worker. The worker holds the DB lock only for - // this SQL read, then is immediately free for heartbeats and other - // requests while CSV serialization runs here, off-worker. - let events = match datastore.get_events(&bucket_id, start, end, limit) { - Ok(events) => events, - Err(err) => { - error!("CSV export stream failed fetching events: {err:?}"); - return; - } - }; - let mut staging = match tempfile::tempfile() { + let staging = match tempfile::tempfile() { Ok(file) => file, Err(err) => { error!("Failed to create CSV staging file: {err}"); return; } }; + // Serialize on the datastore worker, one SQL row at a time, into the + // staging file. 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) { - let mut csv_writer = BufWriter::new(&mut staging); - if let Err(err) = aw_datastore::write_csv_from_events(&events, &mut csv_writer) { + Ok(file) => file, + Err(err) => { error!("CSV export serialization failed: {err:?}"); return; } - // BufWriter::drop swallows flush errors. Detect a full staging - // filesystem (or similar) before we rewind and copy a truncated CSV. - if let Err(err) = csv_writer.flush() { - error!("CSV export flush failed: {err}"); - return; - } - } + }; if let Err(err) = staging.seek(SeekFrom::Start(0)) { error!("CSV staging rewind failed: {err}"); return; From 942ea9cf0f98304798e1b83f3cce35b33323f3a3 Mon Sep 17 00:00:00 2001 From: Bob Date: Fri, 25 Sep 2026 20:42:12 +0000 Subject: [PATCH 07/11] perf(export): profile CSV export memory with the reproducible example Extend `export_memory` with `csv` and `csv-union` modes so the CSV export path has the same reproducible memory measurement that #677 added for the JSON export. `csv-union` gives every event a distinct data key to exercise the MAX_CSV_DATA_COLUMNS fallback. Measured with /usr/bin/time -v (release build, synthetic in-memory data, includes the in-memory input database): csv 100k events 45 MB peak RSS, 31 MB CSV csv 500k events 213 MB peak RSS, 155 MB CSV csv-union 500k events 148 MB peak RSS, 35 MB CSV stream (JSON) 500k events 214 MB peak RSS materialized 500k events 913 MB peak RSS (pre-#677 JSON baseline) Peak RSS for the CSV path is ~4x lower than the materialized baseline on the same dataset, and the 500k-distinct-key case stays bounded (one JSON data column) instead of emitting events x keys cells. Git-Session-Id: 050403af-7294-5a68-bc97-2dab10c51675 --- aw-datastore/examples/export_memory.rs | 100 +++++++++++++++++++------ 1 file changed, 77 insertions(+), 23 deletions(-) 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()); } From 7f72f6ae0fa772da5f1ba297a1881a866cf61d99 Mon Sep 17 00:00:00 2001 From: Bob Date: Fri, 25 Sep 2026 20:53:20 +0000 Subject: [PATCH 08/11] fix(export): bound the CSV key pre-pass; sanitize Content-Disposition In-band review findings on the CSV export: - The key pre-pass collected the full union of data keys before deciding to fall back to the single JSON `data` column, so a bucket whose events each carry a distinct key still held every key in memory -- the unbounded term the PR claims to remove. Stop collecting at MAX_CSV_DATA_COLUMNS, where the fallback decision is already made. Measured on the 500k-distinct-key case: 144 MiB -> 70 MiB peak RSS, 0.90 s -> 0.51 s. - `duration_csv` went through `num_nanoseconds().unwrap_or(0)`, which silently exports `0.000000000` for a duration outside the i64 nanosecond range. Compose from whole seconds and the signed subsecond remainder instead. - `Content-Disposition` interpolated the client-supplied bucket id verbatim. Ids are not restricted at creation and Rocket percent-decodes path segments, so CR/LF in an id would split the response headers. Header values now strip control characters and quote/backslash/semicolon -- applied to the CSV filename and to the pre-existing bucket-export filename (#721 had the same exposure). Tests: streamed_csv_stops_key_pass_at_the_cap_without_truncating_rows, sanitize_header_value_strips_header_metacharacters. Git-Session-Id: 050403af-7294-5a68-bc97-2dab10c51675 --- aw-datastore/src/export.rs | 58 ++++++++++++++++++++++++++++++--- aw-server/src/endpoints/util.rs | 46 ++++++++++++++++++++++++-- 2 files changed, 97 insertions(+), 7 deletions(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index 1a684046..4b44f852 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -156,11 +156,21 @@ fn csv_escape(s: &str) -> String { /// 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 { - let ns = duration.num_nanoseconds().unwrap_or(0); - let sign = if ns < 0 { "-" } else { "" }; - let ns = ns.unsigned_abs(); - format!("{sign}{}.{:09}", ns / 1_000_000_000, ns % 1_000_000_000) + let sign = if *duration < chrono::Duration::zero() { + "-" + } else { + "" + }; + format!( + "{sign}{}.{:09}", + duration.num_seconds().unsigned_abs(), + duration.subsec_nanos().unsigned_abs() + ) } fn event_field_value(event: &Event, key: &str) -> String { @@ -290,7 +300,9 @@ pub(crate) fn write_events_csv( }; // 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; only the (small) key set is held in memory. + // 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(); @@ -307,6 +319,7 @@ pub(crate) fn write_events_csv( .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}")) })? { @@ -322,10 +335,17 @@ pub(crate) fn write_events_csv( 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); @@ -545,6 +565,34 @@ mod tests { assert!(csv.contains("k32"), "last key lost: {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; diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index 09d65ec7..bfbc45b7 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -46,6 +46,22 @@ 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() +} + fn export_filename( datastore: &aw_datastore::Datastore, bucket_id: Option<&str>, @@ -61,7 +77,10 @@ fn export_filename( } }; Ok(match name { - Some(id) => format!("attachment; filename=aw-bucket-export_{id}.json"), + Some(id) => format!( + "attachment; filename=aw-bucket-export_{}.json", + sanitize_header_value(&id) + ), None => "attachment; filename=aw-buckets-export.json".into(), }) } @@ -227,7 +246,10 @@ impl BucketEventsCsvRocket { // (same tradeoff as JSON export / #721). datastore.get_bucket(bucket_id)?; datastore.get_events(bucket_id, start, end, Some(1))?; - let filename = format!("attachment; filename=aw-events-export-{bucket_id}.csv"); + let filename = format!( + "attachment; filename=aw-events-export-{}.csv", + sanitize_header_value(bucket_id) + ); Ok(Self { datastore: datastore.clone(), bucket_id: bucket_id.to_owned(), @@ -299,3 +321,23 @@ impl From for HttpErrorJson { } } } + +#[cfg(test)] +mod tests { + use super::sanitize_header_value; + + #[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"); + } +} From 8e3eaeb37c635c792088b934307ed51205d9203c Mon Sep 17 00:00:00 2001 From: Bob Date: Fri, 25 Sep 2026 20:58:20 +0000 Subject: [PATCH 09/11] fix(export): terminate CSV records with CRLF per RFC-4180 The endpoint advertises RFC-4180 CSV, which delimits each record with CRLF; the writer emitted a bare LF. Readers accept both, but the format claim should match the bytes. Pinned by streamed_csv_terminates_every_record_with_crlf. Also applies to the CSV the staging-file path copies to the response. Git-Session-Id: 050403af-7294-5a68-bc97-2dab10c51675 --- aw-datastore/src/export.rs | 22 +++++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index 4b44f852..62d6e819 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -195,7 +195,8 @@ fn write_csv_record( .write_all(csv_escape(&field).as_bytes()) .map_err(csv_io_err)?; } - writer.write_all(b"\n").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. @@ -565,6 +566,25 @@ mod tests { 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(); From 8a3051d0e0edc74ec9f0595449cc3a406b8e5504 Mon Sep 17 00:00:00 2001 From: Bob Date: Fri, 25 Sep 2026 21:04:44 +0000 Subject: [PATCH 10/11] fix(export): quote the Content-Disposition filename Bucket ids may contain spaces, so `filename=aw-events-export-my bucket.csv` was a malformed parameter that clients could drop. Filenames are now emitted as a quoted-string via a single `content_disposition()` helper, applied to the CSV export and to the bucket export (which had the same unquoted form). The sanitizer already removes `"`, `\` and control characters, so the quoted value cannot be broken out of. Tests: content_disposition_quotes_the_filename, plus the three existing Content-Disposition header assertions updated to the quoted form. Git-Session-Id: 050403af-7294-5a68-bc97-2dab10c51675 --- aw-server/src/endpoints/util.rs | 41 +++++++++++++++++++++++++-------- aw-server/tests/api.rs | 6 ++--- 2 files changed, 34 insertions(+), 13 deletions(-) diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index bfbc45b7..ef0ca1ca 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -62,6 +62,19 @@ fn sanitize_header_value(value: &str) -> String { .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>, @@ -77,11 +90,8 @@ fn export_filename( } }; Ok(match name { - Some(id) => format!( - "attachment; filename=aw-bucket-export_{}.json", - sanitize_header_value(&id) - ), - 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"), }) } @@ -246,10 +256,7 @@ impl BucketEventsCsvRocket { // (same tradeoff as JSON export / #721). datastore.get_bucket(bucket_id)?; datastore.get_events(bucket_id, start, end, Some(1))?; - let filename = format!( - "attachment; filename=aw-events-export-{}.csv", - sanitize_header_value(bucket_id) - ); + let filename = content_disposition(&format!("aw-events-export-{bucket_id}.csv")); Ok(Self { datastore: datastore.clone(), bucket_id: bucket_id.to_owned(), @@ -324,7 +331,21 @@ impl From for HttpErrorJson { #[cfg(test)] mod tests { - use super::sanitize_header_value; + 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() { diff --git a/aw-server/tests/api.rs b/aw-server/tests/api.rs index 3b000ad2..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!({})); @@ -173,7 +173,7 @@ mod api_tests { ); assert_eq!( response.headers().get_one("Content-Disposition"), - Some("attachment; filename=aw-events-export-testbucket.csv") + Some("attachment; filename=\"aw-events-export-testbucket.csv\"") ); let body = response.into_string().unwrap(); // Header row From 67deabfe2362dd1f7b688937d0a2345ce65bab2c Mon Sep 17 00:00:00 2001 From: Bob Date: Fri, 25 Sep 2026 22:30:24 +0000 Subject: [PATCH 11/11] fix(export): correct negative CSV durations and saturate the limit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two correctness fixes in the streamed CSV export path: - `duration_csv` split a negative duration into `num_seconds()` and `subsec_nanos()` while both carried their own sign handling. chrono stores -1.5s as secs=-2, nanos=500_000_000, so the row rendered as "-2.500000000": a full second plus the fraction away from the event's real duration, corrupting any timeline built from the export. Take the magnitude once, before splitting. - `Some(l) => l as i64` wrapped a `u64` limit above `i64::MAX` to a negative value, and SQLite reads a negative LIMIT as "no limit" — so `?limit=18446744073709551615` exported the whole bucket instead of capping it. Saturate via `i64::try_from` (new `sql_limit` helper). Also create the CSV staging file during preflight rather than inside the spawned writer thread, so a full staging filesystem returns a JSON 500 before the response commits `200 OK` instead of an empty body that looks like a genuinely empty bucket. Tests: `csv_duration_handles_negative_durations`, `csv_limit_saturates_instead_of_wrapping`. Git-Session-Id: a0ee77dc-cf96-5982-817e-85a571138c85 --- aw-datastore/src/export.rs | 56 +++++++++++++++++++++++++++------ aw-server/src/endpoints/util.rs | 36 ++++++++++++--------- 2 files changed, 68 insertions(+), 24 deletions(-) diff --git a/aw-datastore/src/export.rs b/aw-datastore/src/export.rs index 62d6e819..4ba6e12b 100644 --- a/aw-datastore/src/export.rs +++ b/aw-datastore/src/export.rs @@ -161,18 +161,35 @@ fn csv_escape(s: &str) -> String { /// `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 { - let sign = if *duration < chrono::Duration::zero() { - "-" + // 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}", - duration.num_seconds().unsigned_abs(), - duration.subsec_nanos().unsigned_abs() + 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(), @@ -283,10 +300,7 @@ pub(crate) fn write_events_csv( write_csv_header(&mut writer, &[])?; return writer.flush().map_err(csv_io_err); } - let limit = match limit_opt { - Some(l) => l as i64, - None => -1, - }; + 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 @@ -489,6 +503,30 @@ mod tests { 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"); diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index ef0ca1ca..4eff492b 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -197,20 +197,15 @@ fn spawn_csv_export_stream( start: Option>, end: Option>, limit: Option, + staging: File, writer: PipeWriter, ) { thread::spawn(move || { - let staging = match tempfile::tempfile() { - Ok(file) => file, - Err(err) => { - error!("Failed to create CSV staging file: {err}"); - return; - } - }; // Serialize on the datastore worker, one SQL row at a time, into the - // staging file. 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. + // 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, @@ -237,6 +232,7 @@ pub struct BucketEventsCsvRocket { end: Option>, limit: Option, filename: String, + staging: File, } impl BucketEventsCsvRocket { @@ -250,12 +246,20 @@ impl BucketEventsCsvRocket { // 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. 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). + // 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(), @@ -264,6 +268,7 @@ impl BucketEventsCsvRocket { end, limit, filename, + staging, }) } } @@ -277,12 +282,13 @@ impl<'r> Responder<'r, 'static> for BucketEventsCsvRocket { 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, writer); + spawn_csv_export_stream(datastore, bucket_id, start, end, limit, staging, writer); Response::build() .status(Status::Ok) .header(Header::new("Content-Disposition", filename))