mirror of
https://github.com/immich-app/walkrs.git
synced 2026-09-30 13:32:57 +08:00
feat: backpressure and fewer allocations (#29)
* tweaks * reduce allocations * backpressure
This commit is contained in:
+1
-3
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
+22
-17
@@ -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<String>,
|
||||
count: usize,
|
||||
buf: Vec<u8>,
|
||||
tx: UnboundedSender<Vec<u8>>,
|
||||
tx: Sender<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 new(tx: Sender<Vec<u8>>) -> 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(())
|
||||
}
|
||||
|
||||
+19
-10
@@ -4,19 +4,28 @@ 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(),
|
||||
)
|
||||
let mut extensions: Vec<String> = 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()
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
+7
-9
@@ -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<Mutex<mpsc::UnboundedReceiver<Vec<u8>>>>,
|
||||
rx: Arc<Mutex<mpsc::Receiver<Vec<u8>>>>,
|
||||
}
|
||||
|
||||
#[napi]
|
||||
@@ -43,10 +43,7 @@ impl AsyncGenerator for Walk {
|
||||
type Next = ();
|
||||
type Return = ();
|
||||
|
||||
fn next(
|
||||
&mut self,
|
||||
_value: Option<Self::Next>,
|
||||
) -> impl std::future::Future<Output = Result<Option<Self::Yield>>> + Send + 'static {
|
||||
fn next(&mut self, _value: Option<Self::Next>) -> impl 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)) }
|
||||
}
|
||||
@@ -54,7 +51,8 @@ impl AsyncGenerator for Walk {
|
||||
|
||||
#[napi]
|
||||
pub fn walk(options: WalkOptions) -> Result<Walk> {
|
||||
let (tx, rx) = mpsc::unbounded_channel::<Vec<u8>>();
|
||||
const CHANNEL_CAPACITY: usize = 16;
|
||||
let (tx, rx) = mpsc::channel::<Vec<u8>>(CHANNEL_CAPACITY);
|
||||
|
||||
if options.paths.is_empty() {
|
||||
return Ok(Walk {
|
||||
@@ -111,7 +109,7 @@ fn build_exclusion_set(exclusion_patterns: &[String]) -> Result<GlobSet> {
|
||||
}
|
||||
|
||||
fn visit(
|
||||
tx: UnboundedSender<Vec<u8>>,
|
||||
tx: Sender<Vec<u8>>,
|
||||
exclusion_set: Arc<GlobSet>,
|
||||
extension_filter: Arc<ExtensionFilter>,
|
||||
) -> Box<dyn FnMut(std::result::Result<DirEntry, ignore::Error>) -> 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;
|
||||
};
|
||||
|
||||
|
||||
Reference in New Issue
Block a user