fix: separate data and error arrays (#49)

* fix: separate data and error arrays

* serialize immediately
This commit is contained in:
Jonathan Jogenfors
2026-09-23 16:55:17 -04:00
committed by GitHub
parent f9c15b18c5
commit 41a2a21bd7
8 changed files with 190 additions and 69 deletions
+3
View File
@@ -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() }}
+4 -1
View File
@@ -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;
+1 -1
View File
@@ -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);
+5 -8
View File
@@ -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<WalkItem[], void, unknown>;
export function walk(options: WalkOptions): AsyncGenerator<WalkBatch, void, unknown>;
+51 -18
View File
@@ -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<u8>,
files_count: usize,
errors_count: usize,
files: Vec<u8>,
errors: Vec<u8>,
tx: Sender<Vec<u8>>,
}
impl BatchSender {
pub fn new(tx: Sender<Vec<u8>>) -> 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;
+107
View File
@@ -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(()));
}
}
+6 -13
View File
@@ -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<String>, message: String },
pub(crate) struct WalkError {
pub path: Option<String>,
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;
}
+13 -28
View File
@@ -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);
});
});