feat: use async iterator (#6)

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
This commit is contained in:
Thomas
2026-02-15 21:36:19 +00:00
committed by GitHub
parent 64b5f35ca9
commit 5445b8ab1f
9 changed files with 221 additions and 162 deletions
+4 -3
View File
@@ -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
+13 -9
View File
@@ -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
+3
View File
@@ -0,0 +1,3 @@
fn main() {
napi_build::setup();
}
+10
View File
@@ -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);
+16 -15
View File
@@ -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"
+1 -1
View File
@@ -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
+43
View File
@@ -0,0 +1,43 @@
use tokio::sync::mpsc::UnboundedSender;
const BATCH_SIZE: usize = 4096;
pub(crate) struct BatchSender {
batch: Vec<String>,
buf: Vec<u8>,
tx: UnboundedSender<Vec<u8>>,
}
impl BatchSender {
pub fn new(tx: UnboundedSender<Vec<u8>>) -> 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();
}
}
+22
View File
@@ -0,0 +1,22 @@
use std::path::Path;
pub(crate) struct ExtensionFilter(Vec<String>);
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)))
}
}
+109 -134
View File
@@ -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<Vec<String>>,
}
#[napi(ts_return_type = "Promise<string[]>")]
pub async fn walk(options: WalkOptions) -> Result<Vec<String>> {
#[napi(async_iterator)]
pub struct Walk {
rx: Arc<Mutex<mpsc::UnboundedReceiver<Vec<u8>>>>,
}
impl napi::bindgen_prelude::AsyncGenerator for Walk {
type Yield = Buffer;
type Next = ();
type Return = ();
fn next(
&mut self,
_value: Option<Self::Next>,
) -> impl std::future::Future<Output = Result<Option<Self::Yield>>> + 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<Walk> {
let (tx, rx) = mpsc::unbounded_channel::<Vec<u8>>();
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<String>,
tx: std::sync::mpsc::Sender<Vec<String>>,
limit: usize,
}
impl BatchSender {
fn new(tx: std::sync::mpsc::Sender<Vec<String>>, 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<Vec<String>> {
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<String> = 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::<Vec<String>>();
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<GlobSet> {
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<Vec<u8>>,
exclusion_set: Arc<GlobSet>,
extension_filter: Arc<ExtensionFilter>,
) -> Box<dyn FnMut(std::result::Result<DirEntry, ignore::Error>) -> 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
})
}