From 5445b8ab1ff1057c3269e46a5effb1983466ed33 Mon Sep 17 00:00:00 2001 From: Thomas <9749173+uhthomas@users.noreply.github.com> Date: Sun, 15 Feb 2026 21:36:19 +0000 Subject: [PATCH] feat: use async iterator (#6) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The async iterator API is quite ergonomic, and allows us to send paths in batches. Sending paths in batches means that we can start processing the paths before we've finished searching, and also makes passing messages between rust and nodejs much faster. I tried quite a few variants, like using a Buffer and manually reconstructing strings, callbacks, etc. The fastest way to pass data between rust and nodejs is by encoding / decoding JSON. It's pretty unintuitive, but this is mainly because rust and C++ (native JSON.parse) are way faster than constructing strings manually with the v8 engine. Both callbacks and iterators are quite fast, and have their own trade-offs. The API for iterators is much more ergonomic however, so that's what I pursued. The performance improvement compared to the current implementation is huge. Up to 2-3x in some cases. ❯ hyperfine 'node example/old.ts' 'node example/stream.ts' 'node example/main.ts' Benchmark 1: node example/old.ts Time (mean ± σ): 2.007 s ± 0.478 s [User: 2.572 s, System: 2.357 s] Range (min … max): 1.619 s … 3.326 s 10 runs Benchmark 2: node example/stream.ts Time (mean ± σ): 1.234 s ± 0.258 s [User: 2.976 s, System: 2.811 s] Range (min … max): 0.913 s … 1.735 s 10 runs Benchmark 3: node example/main.ts Time (mean ± σ): 1.117 s ± 0.267 s [User: 2.425 s, System: 2.484 s] Range (min … max): 0.809 s … 1.680 s 10 runs Summary node example/main.ts ran 1.10 ± 0.35 times faster than node example/stream.ts 1.80 ± 0.61 times faster than node example/old.ts --- Cargo.toml | 7 +- README.md | 22 ++-- build.rs | 3 + example/main.ts | 10 ++ package.json | 31 ++--- pnpm-lock.yaml | 2 +- src/batch_sender.rs | 43 +++++++ src/extension_filter.rs | 22 ++++ src/lib.rs | 243 ++++++++++++++++++---------------------- 9 files changed, 221 insertions(+), 162 deletions(-) create mode 100644 build.rs create mode 100644 example/main.ts create mode 100644 src/batch_sender.rs create mode 100644 src/extension_filter.rs diff --git a/Cargo.toml b/Cargo.toml index 6b50f87..a38b944 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,17 +8,18 @@ publish = false crate-type = ["cdylib", "rlib"] [dependencies] -napi = { version = "3", features = ["async"] } -napi-derive = "3" +napi = { version = "3.8.3", features = ["async"] } +napi-derive = "3.5.2" tokio = { version = "1", features = ["full"] } ignore = "0.4" globset = "0.4" anyhow = "1" serde = { version = "1", features = ["derive"] } +serde_json = "1" chrono = "0.4" [build-dependencies] -built = { version = "0.7", features = ["chrono"] } +napi-build = "2.3.1" [profile.release] lto = true diff --git a/README.md b/README.md index c1204d1..f861aa1 100644 --- a/README.md +++ b/README.md @@ -20,17 +20,21 @@ pnpm add @immich/walkrs import { walk } from '@immich/walkrs'; // Simple usage - walk a directory -const files = await walk({ - root_paths: ['/path/to/scan'], -}); +const files: string[] = []; +for await (const batch of walk({ paths: ['/path/to/scan'] })) { + files.push(...JSON.parse(batch)); +} // Advanced usage with filtering -const photos = await walk({ - root_paths: ['/photos', '/backup/photos'], - supported_extensions: ['.jpg', '.png', '.heic', '.webp'], - exclusion_patterns: ['**/.stfolder/**'], - include_hidden: false, -}); +const photos: string[] = []; +for await (const batch of walk({ + paths: ['/photos', '/backup/photos'], + extensions: ['.jpg', '.png', '.heic', '.webp'], + exclusionPatterns: ['**/.stfolder/**'], + includeHidden: false, +})) { + photos.push(...JSON.parse(batch)); +} ``` ## Performance diff --git a/build.rs b/build.rs new file mode 100644 index 0000000..bbfc9e4 --- /dev/null +++ b/build.rs @@ -0,0 +1,3 @@ +fn main() { + napi_build::setup(); +} diff --git a/example/main.ts b/example/main.ts new file mode 100644 index 0000000..15f343c --- /dev/null +++ b/example/main.ts @@ -0,0 +1,10 @@ +import { walk } from '@immich/walkrs'; + +const path = process.argv[2] || '/'; + +const files: string[] = []; +for await (const batch of walk({ paths: [path] })) { + files.push(...JSON.parse(batch)); +} + +console.log(files.length); diff --git a/package.json b/package.json index 2bacb16..ebc740a 100644 --- a/package.json +++ b/package.json @@ -2,6 +2,7 @@ "name": "@immich/walkrs", "version": "0.0.0", "description": "Fast file tree walker for Node.js, built with ripgrep's ignore crate", + "type": "module", "main": "./dist/index.js", "types": "./dist/index.d.ts", "packageManager": "pnpm@10.28.0+sha512.05df71d1421f21399e053fde567cea34d446fa02c76571441bfc1c7956e98e363088982d940465fd34480d4d90a0668bc12362f8aa88000a64e83d0b0e47be48", @@ -18,14 +19,14 @@ "napi": { "binaryName": "walkrs", "targets": [ - "aarch64-apple-darwin", - "x86_64-apple-darwin", - "aarch64-unknown-linux-gnu", - "x86_64-unknown-linux-gnu", - "aarch64-unknown-linux-musl", - "x86_64-unknown-linux-musl", - "x86_64-pc-windows-msvc" - ], + "aarch64-apple-darwin", + "x86_64-apple-darwin", + "aarch64-unknown-linux-gnu", + "x86_64-unknown-linux-gnu", + "aarch64-unknown-linux-musl", + "x86_64-unknown-linux-musl", + "x86_64-pc-windows-msvc" + ], "package": { "name": "@immich/walkrs" } @@ -35,9 +36,9 @@ "dist" ], "scripts": { - "build": "napi build --platform --release --output-dir dist", - "build:debug": "napi build --platform --debug --output-dir dist", - "build:release": "napi build --platform --release --output-dir dist", + "build": "napi build --platform --release --esm -o dist", + "build:debug": "napi build --platform --esm -o dist", + "build:release": "napi build --platform --release --esm -o dist", "prepublishOnly": "pnpm run build:release", "test": "node --loader tsx ./test.ts", "typecheck": "tsc --noEmit", @@ -53,19 +54,19 @@ }, "devDependencies": { "@eslint/js": "^9.8.0", + "@napi-rs/cli": "^3.5.1", + "@types/node": "^24.10.11", "eslint": "^9.14.0", "eslint-config-prettier": "^10.1.8", "eslint-plugin-prettier": "^5.1.3", "eslint-plugin-unicorn": "^62.0.0", - "@napi-rs/cli": "^3.0.0", - "@types/node": "^24.10.11", "globals": "^16.0.0", "prettier": "^3.7.4", "prettier-plugin-organize-imports": "^4.0.0", "prettier-plugin-rust": "^0.1.9", - "typescript-eslint": "^8.28.0", "tsx": "^4.7.0", - "typescript": "^5.9.2" + "typescript": "^5.9.2", + "typescript-eslint": "^8.28.0" }, "volta": { "node": "24.13.0" diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index ebbd52c..2bf6146 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -12,7 +12,7 @@ importers: specifier: ^9.8.0 version: 9.39.2 '@napi-rs/cli': - specifier: ^3.0.0 + specifier: ^3.5.1 version: 3.5.1(@emnapi/runtime@1.8.1)(@types/node@24.10.13) '@types/node': specifier: ^24.10.11 diff --git a/src/batch_sender.rs b/src/batch_sender.rs new file mode 100644 index 0000000..759278d --- /dev/null +++ b/src/batch_sender.rs @@ -0,0 +1,43 @@ +use tokio::sync::mpsc::UnboundedSender; + +const BATCH_SIZE: usize = 4096; + +pub(crate) struct BatchSender { + batch: Vec, + buf: Vec, + tx: UnboundedSender>, +} + +impl BatchSender { + pub fn new(tx: UnboundedSender>) -> Self { + Self { + batch: Vec::with_capacity(BATCH_SIZE), + buf: Vec::new(), + tx, + } + } + + pub fn send(&mut self, item: String) -> Result<(), ()> { + self.batch.push(item); + if self.batch.len() >= BATCH_SIZE { + self.flush()?; + } + Ok(()) + } + + fn flush(&mut self) -> Result<(), ()> { + if !self.batch.is_empty() { + serde_json::to_writer(&mut self.buf, &self.batch).unwrap(); + self.tx.send(self.buf.clone()).map_err(|_| ())?; + self.buf.clear(); + self.batch.clear(); + } + Ok(()) + } +} + +impl Drop for BatchSender { + fn drop(&mut self) { + let _ = self.flush(); + } +} diff --git a/src/extension_filter.rs b/src/extension_filter.rs new file mode 100644 index 0000000..8a7e59c --- /dev/null +++ b/src/extension_filter.rs @@ -0,0 +1,22 @@ +use std::path::Path; + +pub(crate) struct ExtensionFilter(Vec); + +impl ExtensionFilter { + pub fn new(extensions: &[String]) -> Self { + Self( + extensions + .iter() + .map(|ext| ext.strip_prefix('.').unwrap_or(ext).to_lowercase()) + .collect(), + ) + } + + pub fn is_match(&self, path: &Path) -> bool { + self.0.is_empty() + || path + .extension() + .and_then(|e| e.to_str()) + .is_some_and(|ext| self.0.iter().any(|e| e.eq_ignore_ascii_case(ext))) + } +} diff --git a/src/lib.rs b/src/lib.rs index 528d742..55f766b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,7 +1,18 @@ +mod batch_sender; +mod extension_filter; + +use std::path::Path; +use std::sync::Arc; + +use globset::{GlobSet, GlobSetBuilder}; +use ignore::{DirEntry, WalkBuilder, WalkState}; use napi::bindgen_prelude::*; use napi_derive::napi; -use std::sync::Arc; -use std::sync::atomic::{AtomicBool, Ordering}; +use tokio::sync::Mutex; +use tokio::sync::mpsc::{self, UnboundedSender}; + +use batch_sender::BatchSender; +use extension_filter::ExtensionFilter; #[napi(object)] pub struct WalkOptions { @@ -18,157 +29,121 @@ pub struct WalkOptions { pub extensions: Option>, } -#[napi(ts_return_type = "Promise")] -pub async fn walk(options: WalkOptions) -> Result> { +#[napi(async_iterator)] +pub struct Walk { + rx: Arc>>>, +} + +impl napi::bindgen_prelude::AsyncGenerator for Walk { + type Yield = Buffer; + type Next = (); + type Return = (); + + fn next( + &mut self, + _value: Option, + ) -> impl std::future::Future>> + Send + 'static { + let rx = Arc::clone(&self.rx); + async move { Ok(rx.lock().await.recv().await.map(Into::into)) } + } +} + +#[napi] +pub fn walk(options: WalkOptions) -> Result { + let (tx, rx) = mpsc::unbounded_channel::>(); + if options.paths.is_empty() { - return Ok(Vec::new()); + return Ok(Walk { + rx: Arc::new(Mutex::new(rx)), + }); } - let root_paths = options.paths.clone(); - let include_hidden = options.include_hidden.unwrap_or(false); - let exclusion_patterns = options.exclusion_patterns.unwrap_or_default(); - let extensions = options.extensions.unwrap_or_default(); + let exclusion_set = Arc::new(build_exclusion_set(&options.exclusion_patterns.unwrap_or_default())?); + let extension_set = Arc::new(ExtensionFilter::new(&options.extensions.unwrap_or_default())); - tokio::task::spawn_blocking(move || find_files(&root_paths, include_hidden, &exclusion_patterns, &extensions)) - .await - .map_err(|e| Error::new(Status::GenericFailure, format!("Task join error: {}", e)))? -} - -struct BatchSender { - batch: Vec, - tx: std::sync::mpsc::Sender>, - limit: usize, -} - -impl BatchSender { - fn new(tx: std::sync::mpsc::Sender>, limit: usize) -> Self { - Self { - batch: Vec::with_capacity(limit), - tx, - limit, - } + let mut walk_builder = WalkBuilder::new(&options.paths[0]); + for path in &options.paths[1..] { + walk_builder.add(path); } - fn send(&mut self, item: String) -> std::result::Result<(), ()> { - self.batch.push(item); - if self.batch.len() >= self.limit { - self.flush()?; - } - Ok(()) - } - - fn flush(&mut self) -> std::result::Result<(), ()> { - if !self.batch.is_empty() { - let batch = std::mem::replace(&mut self.batch, Vec::with_capacity(self.limit)); - self.tx.send(batch).map_err(|_| ())?; - } - Ok(()) - } -} - -impl Drop for BatchSender { - fn drop(&mut self) { - let _ = self.flush(); - } -} - -fn find_files( - root_paths: &[String], - include_hidden: bool, - exclusion_patterns: &[String], - extensions: &[String], -) -> Result> { - let mut exclusion_set = globset::GlobSetBuilder::new(); - for pattern in exclusion_patterns { - if let Ok(glob) = globset::Glob::new(pattern) { - exclusion_set.add(glob); - } - } - let exclusion_set = match exclusion_set.build() { - Ok(set) => set, - Err(e) => { - eprintln!("Error building exclusion patterns: {}", e); - globset::GlobSetBuilder::new().build().unwrap() - } - }; - let exclusion_set = Arc::new(exclusion_set); - - let ext_set: std::collections::HashSet = extensions - .iter() - .map(|ext| ext.strip_prefix('.').unwrap_or(ext).to_lowercase()) - .collect(); - let ext_set = Arc::new(ext_set); - - let (tx, rx) = std::sync::mpsc::channel::>(); - let quit_flag = Arc::new(AtomicBool::new(false)); - - let mut walker_builder = ignore::WalkBuilder::new(&root_paths[0]); - for path in &root_paths[1..] { - walker_builder.add(path); - } - - let walker = walker_builder + let walker = walk_builder .git_ignore(false) - .hidden(!include_hidden) + .hidden(!options.include_hidden.unwrap_or(false)) .parents(false) .ignore(false) .git_global(false) .git_exclude(false) .build_parallel(); - walker.run(|| { - let tx = tx.clone(); - let quit_flag = quit_flag.clone(); - let exclusion_set = Arc::clone(&exclusion_set); - let ext_set = Arc::clone(&ext_set); - let mut batch_sender = BatchSender::new(tx, 256); + std::thread::spawn(move || walker.run(|| visit(tx.clone(), Arc::clone(&exclusion_set), Arc::clone(&extension_set)))); - Box::new(move |entry_result| { - if quit_flag.load(Ordering::Relaxed) { - return ignore::WalkState::Quit; - } + Ok(Walk { + rx: Arc::new(Mutex::new(rx)), + }) +} - let Ok(entry) = entry_result else { - return ignore::WalkState::Continue; +fn build_exclusion_set(exclusion_patterns: &[String]) -> Result { + let mut builder = GlobSetBuilder::new(); + for pattern in exclusion_patterns { + builder.add( + globset::GlobBuilder::new(pattern) + .case_insensitive(true) + .build() + .map_err(|e| { + Error::new( + Status::InvalidArg, + format!("Invalid exclusion pattern '{pattern}': {e}"), + ) + })?, + ); + } + builder + .build() + .map_err(|e| Error::new(Status::InvalidArg, format!("Failed to build exclusion patterns: {e}"))) +} + +fn visit( + tx: UnboundedSender>, + exclusion_set: Arc, + extension_filter: Arc, +) -> Box) -> WalkState + Send> { + let mut batch_sender = BatchSender::new(tx); + + Box::new(move |entry_result| { + let Ok(entry) = entry_result else { + return WalkState::Continue; + }; + + let Some(ft) = entry.file_type() else { + return WalkState::Continue; + }; + + let path: &Path = entry.path(); + + if exclusion_set.is_match(path) { + return if ft.is_dir() { + WalkState::Skip + } else { + WalkState::Continue }; + } - let path = entry.path(); + if !ft.is_file() { + return WalkState::Continue; + } - if !exclusion_set.is_empty() && !exclusion_set.matches(path).is_empty() { - return ignore::WalkState::Skip; - } + if !extension_filter.is_match(path) { + return WalkState::Continue; + } - let Some(file_type) = entry.file_type() else { - return ignore::WalkState::Continue; - }; + let Ok(path) = entry.into_path().into_os_string().into_string() else { + return WalkState::Continue; + }; - if !file_type.is_file() { - return ignore::WalkState::Continue; - } - - if !ext_set.is_empty() - && !path - .extension() - .and_then(|e| e.to_str()) - .is_some_and(|ext| ext_set.contains(&ext.to_lowercase())) - { - return ignore::WalkState::Continue; - } - - if batch_sender.send(path.to_string_lossy().into_owned()).is_err() { - quit_flag.store(true, Ordering::Relaxed); - return ignore::WalkState::Quit; - } - ignore::WalkState::Continue - }) - }); - - drop(tx); - - let mut all_files = Vec::new(); - while let Ok(batch) = rx.recv() { - all_files.extend(batch); - } + if batch_sender.send(path).is_err() { + return WalkState::Quit; + } - Ok(all_files) + WalkState::Continue + }) }