From 74af90cda5074c3a8ee4904c63cac8029af5fcb4 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Fri, 13 Mar 2026 04:25:25 +0300 Subject: [PATCH] fjall: export streams the ops instead of collecting them --- src/mirror/fjall.rs | 18 +++++++++++++----- 1 file changed, 13 insertions(+), 5 deletions(-) diff --git a/src/mirror/fjall.rs b/src/mirror/fjall.rs index 657b082..68dce3f 100644 --- a/src/mirror/fjall.rs +++ b/src/mirror/fjall.rs @@ -276,13 +276,21 @@ async fn export( let limit = 1000; let db = fjall.clone(); - let ops = spawn_blocking(move || { + let (tx, rx) = tokio::sync::mpsc::channel::>(64); + + tokio::task::spawn_blocking(move || { let iter = db.export_ops((after + 1)..)?; - iter.take(limit).collect::>>() - }) - .await?; + for op in iter.take(limit) { + if tx.blocking_send(op).is_err() { + break; + } + } + anyhow::Ok(()) + }); - let stream = futures::stream::iter(ops).map(|op| { + // todo: its a bit annoying that errors just cut it off here... + let stream = tokio_stream::wrappers::ReceiverStream::new(rx).map(|result| { + let op = result.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?; let mut json = serde_json::to_string(&op.to_sequenced_json()).unwrap(); json.push('\n'); Ok::<_, std::io::Error>(json) -- 2.51.2