From 41a2a21bd76262efc1425758e576e9e3dd94b785 Mon Sep 17 00:00:00 2001 From: Jonathan Jogenfors Date: Wed, 23 Sep 2026 22:55:17 +0200 Subject: [PATCH] fix: separate data and error arrays (#49) * fix: separate data and error arrays * serialize immediately --- .github/workflows/test.yml | 3 ++ bench/bench.ts | 5 +- example/main.ts | 2 +- lib/index.d.ts | 13 ++--- src/batch_sender.rs | 69 +++++++++++++++++------- src/batch_sender/tests.rs | 107 +++++++++++++++++++++++++++++++++++++ src/lib.rs | 19 +++---- test/walk.spec.ts | 41 +++++--------- 8 files changed, 190 insertions(+), 69 deletions(-) create mode 100644 src/batch_sender/tests.rs diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index d52f697..6b5390f 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -38,6 +38,9 @@ jobs: - name: Check rust formatting run: pnpm rust:format if: ${{ !cancelled() }} + - name: Run Rust tests + run: cargo test + if: ${{ !cancelled() }} - name: Check ts linting run: pnpm ts:lint if: ${{ !cancelled() }} diff --git a/bench/bench.ts b/bench/bench.ts index 5d6e53c..421b900 100644 --- a/bench/bench.ts +++ b/bench/bench.ts @@ -25,7 +25,10 @@ async function run(datasetPath: string, benchmarkOptions?: BenchmarkOptions): Pr let fileCount = 0; for await (const batch of walk(walkOptions)) { - fileCount += batch.length; + if (batch.errors.length > 0) { + throw new Error(`Walk encountered errors: ${batch.errors.map((e) => e.message).join(', ')}`); + } + fileCount += batch.files.length; } return fileCount; diff --git a/example/main.ts b/example/main.ts index 30247c3..6c10b3e 100644 --- a/example/main.ts +++ b/example/main.ts @@ -4,7 +4,7 @@ const path = process.argv[2] || '/'; const files: string[] = []; for await (const batch of walk({ paths: [path] })) { - files.push(...batch); + files.push(...batch.files); } console.log(files.length); diff --git a/lib/index.d.ts b/lib/index.d.ts index 20188da..a4c325f 100644 --- a/lib/index.d.ts +++ b/lib/index.d.ts @@ -1,16 +1,13 @@ export { WalkOptions } from '../dist/index.js'; -export type WalkEntry = { - type: 'entry'; - path: string; -}; - export type WalkError = { - type: 'error'; path?: string; message: string; }; -export type WalkItem = WalkEntry | WalkError; +export type WalkBatch = { + files: string[]; + errors: WalkError[]; +}; -export function walk(options: WalkOptions): AsyncGenerator; +export function walk(options: WalkOptions): AsyncGenerator; diff --git a/src/batch_sender.rs b/src/batch_sender.rs index 54a203e..a9278fe 100644 --- a/src/batch_sender.rs +++ b/src/batch_sender.rs @@ -1,42 +1,72 @@ -use crate::WalkItem; +use crate::WalkError; use tokio::sync::mpsc::Sender; const BATCH_SIZE: usize = 4096; const BUF_CAPACITY: usize = BATCH_SIZE * 100; +const FILES_PREFIX: &[u8] = br#"{"files":["#; pub(crate) struct BatchSender { - count: usize, - buf: Vec, + files_count: usize, + errors_count: usize, + files: Vec, + errors: Vec, tx: Sender>, } impl BatchSender { pub fn new(tx: Sender>) -> Self { - let mut buf = Vec::with_capacity(BUF_CAPACITY); - buf.push(b'['); - Self { count: 0, buf, tx } + let mut files = Vec::with_capacity(BUF_CAPACITY); + files.extend_from_slice(FILES_PREFIX); + Self { + files_count: 0, + errors_count: 0, + files, + errors: Vec::new(), + tx, + } } - pub fn send(&mut self, item: WalkItem) -> Result<(), ()> { - if self.count > 0 { - self.buf.push(b','); + pub fn send_entry(&mut self, path: &str) -> Result<(), ()> { + if self.files_count > 0 { + self.files.push(b','); } - serde_json::to_writer(&mut self.buf, &item).expect("WalkItem serialization should never fail"); - self.count += 1; - if self.count >= BATCH_SIZE { + + // Serialize immediately + serde_json::to_writer(&mut self.files, path).expect("Path serialization should never fail"); + self.files_count += 1; + if self.files_count + self.errors_count >= BATCH_SIZE { + self.flush()?; + } + Ok(()) + } + + pub fn send_error(&mut self, error: WalkError) -> Result<(), ()> { + if self.errors_count > 0 { + self.errors.push(b','); + } + + // Serialize immediately + serde_json::to_writer(&mut self.errors, &error).expect("WalkError serialization should never fail"); + self.errors_count += 1; + if self.files_count + self.errors_count >= BATCH_SIZE { self.flush()?; } Ok(()) } fn flush(&mut self) -> Result<(), ()> { - if self.count > 0 { - self.buf.push(b']'); - let mut new_buf = Vec::with_capacity(BUF_CAPACITY); - new_buf.push(b'['); - let buf = std::mem::replace(&mut self.buf, new_buf); + if self.files_count + self.errors_count > 0 { + // Merge file and error buffers + self.files.extend_from_slice(br#"],"errors":["#); + self.files.extend_from_slice(&self.errors); + self.files.extend_from_slice(b"]}"); + let mut files = Vec::with_capacity(BUF_CAPACITY); + files.extend_from_slice(FILES_PREFIX); + let buf = std::mem::replace(&mut self.files, files); + self.errors.clear(); + self.files_count = 0; + self.errors_count = 0; self.tx.blocking_send(buf).map_err(|_| ())?; - self.count = 0; } Ok(()) } @@ -47,3 +77,6 @@ impl Drop for BatchSender { let _ = self.flush(); } } + +#[cfg(test)] +mod tests; diff --git a/src/batch_sender/tests.rs b/src/batch_sender/tests.rs new file mode 100644 index 0000000..81d8729 --- /dev/null +++ b/src/batch_sender/tests.rs @@ -0,0 +1,107 @@ +use super::*; +use serde_json::{Value, json}; +use tokio::sync::mpsc::{channel, error::TryRecvError}; + +#[test] +fn empty_sender_does_not_send_a_batch() { + let (tx, mut rx) = channel(1); + let mut sender = BatchSender::new(tx); + sender.flush().unwrap(); + drop(sender); + assert_eq!(rx.try_recv(), Err(TryRecvError::Disconnected)); +} + +#[test] +fn escapes_interleaved_files_and_errors() { + let (tx, mut rx) = channel(1); + let mut sender = BatchSender::new(tx); + let path = "photos/\"quoted\"\\newline\n\t雪.jpg"; + let message = "failed: \"quoted\"\\newline\n\t\0雪"; + sender + .send_error(WalkError { + path: None, + message: message.into(), + }) + .unwrap(); + sender.send_entry(path).unwrap(); + sender + .send_error(WalkError { + path: Some(path.into()), + message: message.into(), + }) + .unwrap(); + sender.send_entry("").unwrap(); + drop(sender); + + let batch: Value = serde_json::from_slice(&rx.try_recv().unwrap()).unwrap(); + assert_eq!( + batch, + json!({ + "files": [path, ""], + "errors": [{"path": null, "message": message}, {"path": path, "message": message}] + }) + ); + assert_eq!(rx.try_recv(), Err(TryRecvError::Disconnected)); +} + +#[test] +fn batches_files_errors_and_mixed_items_at_the_combined_limit() { + for total in [BATCH_SIZE - 1, BATCH_SIZE, BATCH_SIZE + 1, 2 * BATCH_SIZE + 1] { + for mode in 0..3 { + let (tx, mut rx) = channel(3); + let mut sender = BatchSender::new(tx); + for i in 0..total { + if mode == 1 || (mode == 2 && i % 2 == 0) { + sender + .send_error(WalkError { + path: None, + message: i.to_string(), + }) + .unwrap(); + } else { + sender.send_entry(&i.to_string()).unwrap(); + } + assert_eq!(rx.len(), (i + 1) / BATCH_SIZE); + } + drop(sender); + + for start in (0..total).step_by(BATCH_SIZE) { + let end = (start + BATCH_SIZE).min(total); + let mut files = Vec::new(); + let mut errors = Vec::new(); + for i in start..end { + if mode == 1 || (mode == 2 && i % 2 == 0) { + errors.push(json!({"path": null, "message": i.to_string()})); + } else { + files.push(i.to_string()); + } + } + let batch: Value = serde_json::from_slice(&rx.try_recv().unwrap()).unwrap(); + assert_eq!(batch, json!({"files": files, "errors": errors})); + } + assert_eq!(rx.try_recv(), Err(TryRecvError::Disconnected)); + } + } +} + +#[test] +fn closed_receiver_returns_an_error() { + for last_item_is_error in [false, true] { + let (tx, rx) = channel(1); + let mut sender = BatchSender::new(tx); + drop(rx); + for _ in 0..BATCH_SIZE - 1 { + sender.send_entry("file.jpg").unwrap(); + } + let result = if last_item_is_error { + sender.send_error(WalkError { + path: None, + message: "error".into(), + }) + } else { + sender.send_entry("file.jpg") + }; + assert_eq!(result, Err(())); + assert_eq!(sender.flush(), Ok(())); + } +} diff --git a/src/lib.rs b/src/lib.rs index f47f710..f2f4056 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -15,12 +15,9 @@ use batch_sender::BatchSender; use extension_filter::ExtensionFilter; #[derive(Debug, Clone, serde::Serialize)] -#[serde(tag = "type", rename_all = "lowercase")] -pub(crate) enum WalkItem { - #[serde(rename = "entry")] - Entry { path: String }, - #[serde(rename = "error")] - Error { path: Option, message: String }, +pub(crate) struct WalkError { + pub path: Option, + pub message: String, } #[napi(object)] @@ -130,11 +127,11 @@ fn visit( Err(err) => { // Report the error and continue walking // The error message from ignore crate already includes the path - let error = WalkItem::Error { + let error = WalkError { path: None, message: err.to_string(), }; - if batch_sender.send(error).is_err() { + if batch_sender.send_error(error).is_err() { return WalkState::Quit; } return WalkState::Continue; @@ -167,11 +164,7 @@ fn visit( return WalkState::Continue; }; - let item = WalkItem::Entry { - path: path_str.to_string(), - }; - - if batch_sender.send(item).is_err() { + if batch_sender.send_entry(path_str).is_err() { return WalkState::Quit; } diff --git a/test/walk.spec.ts b/test/walk.spec.ts index ea9261f..bdcb978 100644 --- a/test/walk.spec.ts +++ b/test/walk.spec.ts @@ -166,10 +166,10 @@ const tests: TestCase[] = [ extensions: ['.jpg', '.jpeg', '.tiff', '.tif', '.dng', '.nef'], }, files: { - '/photos/image.jpg': true, - '/photos/image.Jpg': true, - '/photos/image.jpG': true, - '/photos/image.JPG': true, + '/photos/image1.jpg': true, + '/photos/image2.Jpg': true, + '/photos/image3.jpG': true, + '/photos/image4.JPG': true, '/photos/image.jpEg': true, '/photos/image.TIFF': true, '/photos/image.tif': true, @@ -221,9 +221,7 @@ describe('walk', () => { const actual: string[] = []; for await (const batch of walk(adjustedOptions)) { - // Filter for entries only (ignore errors) and extract paths - const paths = batch.filter((item) => item.type === 'entry').map((item) => item.path); - actual.push(...paths); + actual.push(...batch.files); } const expected = Object.entries(files) .filter((entry) => entry[1]) @@ -270,13 +268,8 @@ describe('walk', () => { const errors: Array<{ path?: string; message: string }> = []; for await (const batch of walk(options)) { - for (const item of batch) { - if (item.type === 'entry') { - entries.push(item.path); - } else if (item.type === 'error') { - errors.push({ path: item.path, message: item.message }); - } - } + entries.push(...batch.files); + errors.push(...batch.errors); } // Should have found the accessible file @@ -302,30 +295,22 @@ describe('walk', () => { extensions: ['.jpg'], }; - const entries: string[] = []; + const files: string[] = []; const errors: Array<{ path?: string; message: string }> = []; for await (const batch of walk(options)) { - for (const item of batch) { - if (item.type === 'entry') { - entries.push(item.path); - } else if (item.type === 'error') { - errors.push({ path: item.path, message: item.message }); - } - } + files.push(...batch.files); + errors.push(...batch.errors); } - // Should have found all files (directory listing doesn't require file read permissions) - expect(entries).toContain(path.join(tempDir, 'photos', 'accessible1.jpg')); - expect(entries).toContain(path.join(tempDir, 'photos', 'accessible2.jpg')); + expect(files).toContain(path.join(tempDir, 'photos', 'accessible1.jpg')); + expect(files).toContain(path.join(tempDir, 'photos', 'accessible2.jpg')); // File is still listed even with 0o000 permissions (directory walk only needs directory read permission) - expect(entries).toContain(path.join(tempDir, 'photos', 'restricted.jpg')); + expect(files).toContain(path.join(tempDir, 'photos', 'restricted.jpg')); - // No errors expected since we're only walking, not reading file contents expect(errors.length).toBe(0); - // Restore permissions for cleanup await fs.chmod(path.join(tempDir, 'photos', 'restricted.jpg'), 0o644); }); });