From ae8e6c9bdb100bd6d7ea1f5194ca947c2da0c832 Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 23 Sep 2026 00:31:27 +0000 Subject: [PATCH 1/2] fix(export): send HTTP headers before serializing large exports #677 still wrote the full JSON tempfile before responding, so a ~500k-event export could sit silent until the 30s web UI timeout. Open a pipe, return 200 + Content-Disposition immediately, and serialize into the body. Missing buckets still 404 before headers. Mid-stream failures truncate the download instead of hanging the connection. ActivityWatch/aw-android#228 Git-Session-Id: a180614b-5a5a-5d29-84d8-e1c4e076b990 --- Cargo.lock | 1 - aw-server/Cargo.toml | 1 - aw-server/src/endpoints/util.rs | 96 +++++++++++++++++++++++++-------- aw-server/tests/api.rs | 18 +++++++ 4 files changed, 93 insertions(+), 23 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index f200bd93..07907aa0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -294,7 +294,6 @@ dependencies = [ "serde", "serde_json", "subtle", - "tempfile", "tokio", "toml", "uuid", diff --git a/aw-server/Cargo.toml b/aw-server/Cargo.toml index b3cf34df..13df70e7 100644 --- a/aw-server/Cargo.toml +++ b/aw-server/Cargo.toml @@ -40,7 +40,6 @@ uuid = { version = "1.3", features = ["serde", "v4"] } clap = { version = "4.1", features = ["derive", "cargo"] } log-panics = { version = "2", features = ["with-backtrace"]} subtle = "2" -tempfile = "3" rust-embed = { version = "8.0.0", features = ["interpolate-folder-path", "debug-embed"] } aw-datastore = { path = "../aw-datastore", default-features = false } aw-models = { path = "../aw-models" } diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index 8fce4582..db3ad562 100644 --- a/aw-server/src/endpoints/util.rs +++ b/aw-server/src/endpoints/util.rs @@ -1,4 +1,6 @@ -use std::io::{Cursor, Seek, SeekFrom}; +use std::fs::File; +use std::io::{pipe, Cursor, PipeReader, PipeWriter}; +use std::thread; use rocket::http::ContentType; use rocket::http::Header; @@ -38,43 +40,95 @@ impl<'r> Responder<'r, 'static> for HttpErrorJson { } pub struct BucketsExportRocket { - file: std::fs::File, + datastore: aw_datastore::Datastore, + bucket_id: Option, filename: String, } +fn export_filename( + datastore: &aw_datastore::Datastore, + bucket_id: Option<&str>, +) -> Result { + let name = match bucket_id { + Some(id) => { + datastore.get_bucket(id)?; + Some(id.to_owned()) + } + None => { + let buckets = datastore.get_buckets()?; + (buckets.len() == 1).then(|| buckets.into_keys().next().unwrap()) + } + }; + Ok(match name { + Some(id) => format!("attachment; filename=aw-bucket-export_{id}.json"), + None => "attachment; filename=aw-buckets-export.json".into(), + }) +} + +#[cfg(not(any(unix, windows)))] +compile_error!("export streaming requires unix or windows anonymous pipes"); + +fn pipe_reader_to_file(reader: PipeReader) -> File { + #[cfg(unix)] + { + File::from(std::os::fd::OwnedFd::from(reader)) + } + #[cfg(windows)] + { + File::from(std::os::windows::io::OwnedHandle::from(reader)) + } +} + +fn pipe_writer_to_file(writer: PipeWriter) -> File { + #[cfg(unix)] + { + File::from(std::os::fd::OwnedFd::from(writer)) + } + #[cfg(windows)] + { + File::from(std::os::windows::io::OwnedHandle::from(writer)) + } +} + impl BucketsExportRocket { pub fn new( datastore: &aw_datastore::Datastore, bucket_id: Option<&str>, ) -> Result { - let io_error = |err: std::io::Error| { - error!("Failed to prepare export file: {err}"); - HttpErrorJson::new( - Status::InternalServerError, - "Failed to prepare export file".into(), - ) - }; - // tempfile creates a private file and removes it when the response is - // dropped. Spooling preserves HTTP errors even if serialization or disk - // writes fail, while keeping event buffering bounded. - let file = tempfile::tempfile().map_err(io_error)?; - let (mut file, name) = datastore.export_to_file(bucket_id, file)?; - file.seek(SeekFrom::Start(0)).map_err(io_error)?; - let filename = match name { - Some(id) => format!("attachment; filename=aw-bucket-export_{id}.json"), - None => "attachment; filename=aw-buckets-export.json".into(), - }; - Ok(Self { file, filename }) + // Resolve the download name and 404 missing buckets before the + // response is built. Serialization itself runs after headers so a + // slow export does not look like a hung connection. + let filename = export_filename(datastore, bucket_id)?; + Ok(Self { + datastore: datastore.clone(), + bucket_id: bucket_id.map(str::to_owned), + filename, + }) } } impl<'r> Responder<'r, 'static> for BucketsExportRocket { fn respond_to(self, _: &Request) -> response::Result<'static> { + let (reader, writer) = pipe().map_err(|err| { + error!("Failed to open export pipe: {err}"); + Status::InternalServerError + })?; + let datastore = self.datastore; + let bucket_id = self.bucket_id; + thread::spawn(move || { + if let Err(err) = + datastore.export_to_file(bucket_id.as_deref(), pipe_writer_to_file(writer)) + { + error!("Export stream failed: {err:?}"); + } + }); Response::build() .status(Status::Ok) .header(Header::new("Content-Disposition", self.filename)) .header(ContentType::JSON) - .streamed_body(rocket::tokio::fs::File::from_std(self.file)) + .streamed_body(rocket::tokio::fs::File::from_std(pipe_reader_to_file( + reader, + ))) .ok() } } diff --git a/aw-server/tests/api.rs b/aw-server/tests/api.rs index 99690f00..58666fcb 100644 --- a/aw-server/tests/api.rs +++ b/aw-server/tests/api.rs @@ -71,6 +71,24 @@ mod api_tests { assert!(body["message"].as_str().unwrap().contains("does not exist")); } + #[test] + fn export_all_empty_uses_plural_filename_and_opens_before_body() { + let server = setup_testserver(); + let client = Client::untracked(server).unwrap(); + let response = client + .get("/api/0/export") + .header(Header::new("Host", "127.0.0.1:5600")) + .dispatch(); + assert_eq!(response.status(), Status::Ok); + assert_eq!(response.content_type(), Some(ContentType::JSON)); + assert_eq!( + response.headers().get_one("Content-Disposition"), + Some("attachment; filename=aw-buckets-export.json") + ); + let body: Value = serde_json::from_str(&response.into_string().unwrap()).unwrap(); + assert_eq!(body["buckets"], json!({})); + } + #[test] fn test_bucket() { let server = setup_testserver(); From e6e245fe9f7e5fc56cb51529cbcfe568625f62a9 Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 23 Sep 2026 00:53:24 +0000 Subject: [PATCH 2/2] fix(export): don't let a slow client stall the datastore worker Serialize to a tempfile on the worker (disk-paced, same as #677), then copy to the response pipe from the export thread. Headers still go out before the body; unread or slow downloads no longer block heartbeats. ActivityWatch/aw-android#228 Git-Session-Id: 3f22e322-f840-50e9-aa06-6db541c7f001 --- Cargo.lock | 1 + aw-server/Cargo.toml | 1 + aw-server/src/endpoints/util.rs | 46 ++++++++++++++++++++++++------- aw-server/tests/api.rs | 48 +++++++++++++++++++++++++++++++++ 4 files changed, 86 insertions(+), 10 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 07907aa0..f200bd93 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -294,6 +294,7 @@ dependencies = [ "serde", "serde_json", "subtle", + "tempfile", "tokio", "toml", "uuid", diff --git a/aw-server/Cargo.toml b/aw-server/Cargo.toml index 13df70e7..b3cf34df 100644 --- a/aw-server/Cargo.toml +++ b/aw-server/Cargo.toml @@ -40,6 +40,7 @@ uuid = { version = "1.3", features = ["serde", "v4"] } clap = { version = "4.1", features = ["derive", "cargo"] } log-panics = { version = "2", features = ["with-backtrace"]} subtle = "2" +tempfile = "3" rust-embed = { version = "8.0.0", features = ["interpolate-folder-path", "debug-embed"] } aw-datastore = { path = "../aw-datastore", default-features = false } aw-models = { path = "../aw-models" } diff --git a/aw-server/src/endpoints/util.rs b/aw-server/src/endpoints/util.rs index db3ad562..a96c583f 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::{pipe, Cursor, PipeReader, PipeWriter}; +use std::io::{copy, pipe, Cursor, PipeReader, PipeWriter, Seek, SeekFrom}; use std::thread; use rocket::http::ContentType; @@ -90,6 +90,40 @@ fn pipe_writer_to_file(writer: PipeWriter) -> File { } } +/// Serialize on the datastore worker into a private tempfile, then copy to +/// the client pipe from this thread. The worker stays disk-paced; a slow +/// or dropped download must not stall heartbeats (see `ServerState`). +fn spawn_export_stream( + datastore: aw_datastore::Datastore, + bucket_id: Option, + writer: PipeWriter, +) { + thread::spawn(move || { + let staging = match tempfile::tempfile() { + Ok(file) => file, + Err(err) => { + error!("Failed to create export staging file: {err}"); + return; + } + }; + let mut staging = match datastore.export_to_file(bucket_id.as_deref(), staging) { + Ok((file, _)) => file, + Err(err) => { + error!("Export stream failed: {err:?}"); + return; + } + }; + if let Err(err) = staging.seek(SeekFrom::Start(0)) { + error!("Failed to rewind export staging file: {err}"); + return; + } + let mut writer = pipe_writer_to_file(writer); + if let Err(err) = copy(&mut staging, &mut writer) { + error!("Export stream copy failed: {err}"); + } + }); +} + impl BucketsExportRocket { pub fn new( datastore: &aw_datastore::Datastore, @@ -113,15 +147,7 @@ impl<'r> Responder<'r, 'static> for BucketsExportRocket { error!("Failed to open export pipe: {err}"); Status::InternalServerError })?; - let datastore = self.datastore; - let bucket_id = self.bucket_id; - thread::spawn(move || { - if let Err(err) = - datastore.export_to_file(bucket_id.as_deref(), pipe_writer_to_file(writer)) - { - error!("Export stream failed: {err:?}"); - } - }); + spawn_export_stream(self.datastore, self.bucket_id, writer); Response::build() .status(Status::Ok) .header(Header::new("Content-Disposition", self.filename)) diff --git a/aw-server/tests/api.rs b/aw-server/tests/api.rs index 58666fcb..7a7b7914 100644 --- a/aw-server/tests/api.rs +++ b/aw-server/tests/api.rs @@ -89,6 +89,54 @@ mod api_tests { assert_eq!(body["buckets"], json!({})); } + #[test] + fn unread_export_body_does_not_block_datastore() { + let server = setup_testserver(); + let datastore = server + .state::() + .unwrap() + .datastore + .clone(); + let bucket: Bucket = serde_json::from_value(json!({ + "id": "big", "type": "test", "client": "test", "hostname": "test" + })) + .unwrap(); + datastore.create_bucket(&bucket).unwrap(); + let mut event = aw_models::Event::default(); + event + .data + .insert("blob".into(), json!("x".repeat(128 * 1024))); + datastore.insert_events("big", &[event]).unwrap(); + + let client = Client::untracked(server).unwrap(); + let response = client + .get("/api/0/buckets/big/export") + .header(Header::new("Host", "127.0.0.1:5600")) + .dispatch(); + assert_eq!(response.status(), Status::Ok); + + // Let Command::Export start. A client-paced pipe would fill here and + // stall the worker; staging to a tempfile must not. + std::thread::sleep(std::time::Duration::from_millis(100)); + let (tx, rx) = std::sync::mpsc::channel(); + let ds = datastore.clone(); + std::thread::spawn(move || { + let _ = tx.send(ds.get_buckets()); + }); + rx.recv_timeout(std::time::Duration::from_secs(3)) + .expect("datastore worker blocked by unread export body") + .unwrap(); + + let body: Value = serde_json::from_str(&response.into_string().unwrap()).unwrap(); + assert_eq!( + body["buckets"]["big"]["events"][0]["data"]["blob"] + .as_str() + .unwrap() + .len(), + 128 * 1024 + ); + } + #[test] fn test_bucket() { let server = setup_testserver();