From e76d8cb0f16539c9f36b658a294b8d905251b5e4 Mon Sep 17 00:00:00 2001 From: Mert <101130780+mertalev@users.noreply.github.com> Date: Wed, 18 Feb 2026 03:45:44 -0500 Subject: [PATCH] feat: backpressure and fewer allocations (#29) * tweaks * reduce allocations * backpressure --- Cargo.toml | 4 +--- rustfmt.toml | 20 -------------------- src/batch_sender.rs | 39 ++++++++++++++++++++++----------------- src/extension_filter.rs | 29 +++++++++++++++++++---------- src/lib.rs | 16 +++++++--------- 5 files changed, 49 insertions(+), 59 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 19282f3..765017a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,13 +10,11 @@ crate-type = ["cdylib", "rlib"] [dependencies] napi = { version = "3.8.3", features = ["async"] } napi-derive = "3.5.2" -tokio = { version = "1", features = ["full"] } +tokio = { version = "1", features = ["sync"] } ignore = "0.4" globset = "0.4" -anyhow = "1" serde = { version = "1", features = ["derive"] } serde_json = "1" -chrono = "0.4" [build-dependencies] napi-build = "2.3.1" diff --git a/rustfmt.toml b/rustfmt.toml index 9cb471d..cf1e170 100644 --- a/rustfmt.toml +++ b/rustfmt.toml @@ -1,27 +1,7 @@ edition = "2024" -hard_tabs = false tab_spaces = 2 max_width = 120 newline_style = "Unix" -use_small_heuristics = "Default" -reorder_imports = true -reorder_modules = true -remove_nested_parens = true - -fn_single_line = false -where_single_line = false - -imports_granularity = "Crate" - -trailing_comma = "Vertical" - -wrap_comments = false -format_code_in_doc_comments = false -normalize_comments = false - -format_strings = false -format_macro_matchers = true -format_macro_bodies = true diff --git a/src/batch_sender.rs b/src/batch_sender.rs index 759278d..1096f10 100644 --- a/src/batch_sender.rs +++ b/src/batch_sender.rs @@ -1,36 +1,41 @@ -use tokio::sync::mpsc::UnboundedSender; +use tokio::sync::mpsc::Sender; const BATCH_SIZE: usize = 4096; +const BUF_CAPACITY: usize = BATCH_SIZE * 100; pub(crate) struct BatchSender { - batch: Vec, + count: usize, buf: Vec, - tx: UnboundedSender>, + tx: Sender>, } impl BatchSender { - pub fn new(tx: UnboundedSender>) -> Self { - Self { - batch: Vec::with_capacity(BATCH_SIZE), - buf: Vec::new(), - tx, - } + pub fn new(tx: Sender>) -> Self { + let mut buf = Vec::with_capacity(BUF_CAPACITY); + buf.push(b'['); + Self { count: 0, buf, tx } } - pub fn send(&mut self, item: String) -> Result<(), ()> { - self.batch.push(item); - if self.batch.len() >= BATCH_SIZE { + pub fn send(&mut self, item: &str) -> Result<(), ()> { + if self.count > 0 { + self.buf.push(b','); + } + serde_json::to_writer(&mut self.buf, item).unwrap(); + self.count += 1; + if self.count >= 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(); + 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); + self.tx.blocking_send(buf).map_err(|_| ())?; + self.count = 0; } Ok(()) } diff --git a/src/extension_filter.rs b/src/extension_filter.rs index 8a7e59c..358d8a4 100644 --- a/src/extension_filter.rs +++ b/src/extension_filter.rs @@ -4,19 +4,28 @@ 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(), - ) + let mut extensions: Vec = extensions + .iter() + .map(|ext| ext.strip_prefix('.').unwrap_or(ext).to_lowercase()) + .collect(); + extensions.sort(); + extensions.dedup(); + Self(extensions) } 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))) + || path.extension().and_then(|e| e.to_str()).is_some_and(|ext| { + let mut buf = [0u8; 16]; + let ext = ext.as_bytes(); + if ext.len() > buf.len() { + return false; + } + let slot = &mut buf[..ext.len()]; + slot.copy_from_slice(ext); + slot.make_ascii_lowercase(); + let lowered = unsafe { std::str::from_utf8_unchecked(slot) }; + self.0.binary_search_by_key(&lowered, |s| s.as_str()).is_ok() + }) } } diff --git a/src/lib.rs b/src/lib.rs index e20ba7b..63ae494 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -9,7 +9,7 @@ use ignore::{DirEntry, WalkBuilder, WalkState}; use napi::bindgen_prelude::*; use napi_derive::napi; use tokio::sync::Mutex; -use tokio::sync::mpsc::{self, UnboundedSender}; +use tokio::sync::mpsc::{self, Sender}; use batch_sender::BatchSender; use extension_filter::ExtensionFilter; @@ -34,7 +34,7 @@ pub struct WalkOptions { #[napi(async_iterator)] pub struct Walk { - rx: Arc>>>, + rx: Arc>>>, } #[napi] @@ -43,10 +43,7 @@ impl AsyncGenerator for Walk { type Next = (); type Return = (); - fn next( - &mut self, - _value: Option, - ) -> impl std::future::Future>> + Send + 'static { + fn next(&mut self, _value: Option) -> impl Future>> + Send + 'static { let rx = Arc::clone(&self.rx); async move { Ok(rx.lock().await.recv().await.map(Into::into)) } } @@ -54,7 +51,8 @@ impl AsyncGenerator for Walk { #[napi] pub fn walk(options: WalkOptions) -> Result { - let (tx, rx) = mpsc::unbounded_channel::>(); + const CHANNEL_CAPACITY: usize = 16; + let (tx, rx) = mpsc::channel::>(CHANNEL_CAPACITY); if options.paths.is_empty() { return Ok(Walk { @@ -111,7 +109,7 @@ fn build_exclusion_set(exclusion_patterns: &[String]) -> Result { } fn visit( - tx: UnboundedSender>, + tx: Sender>, exclusion_set: Arc, extension_filter: Arc, ) -> Box) -> WalkState + Send> { @@ -144,7 +142,7 @@ fn visit( return WalkState::Continue; } - let Ok(path) = entry.into_path().into_os_string().into_string() else { + let Some(path) = path.to_str() else { return WalkState::Continue; };