perf: load large folders in chunks

This commit is contained in:
sxyazi 2023-09-06 18:22:11 +08:00
parent eee68a99b4
commit 5a6179005c
No known key found for this signature in database
11 changed files with 159 additions and 215 deletions

60
Cargo.lock generated
View file

@ -225,9 +225,9 @@ checksum = "a3e2c3daef883ecc1b5d58c15adae93470a91d425f3532ba1695849656af3fc1"
[[package]]
name = "bytemuck"
version = "1.13.1"
version = "1.14.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "17febce684fd15d89027105661fec94afb475cb995fbc59d2865198446ba2eea"
checksum = "374d28ec25809ee0e23827c2ab573d729e293f281dfe393500e7ad618baa61c6"
[[package]]
name = "byteorder"
@ -264,9 +264,9 @@ checksum = "baf1de4339761588bc0619e3cbc0120ee582ebb74b53b4efbf79117bd2da40fd"
[[package]]
name = "chrono"
version = "0.4.28"
version = "0.4.29"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "95ed24df0632f708f5f6d8082675bef2596f7084dee3dd55f632290bf35bfe0f"
checksum = "d87d9d13be47a5b7c3907137f1290b0459a7f80efb26be8c52afb11963bccb02"
dependencies = [
"android-tzdata",
"iana-time-zone",
@ -305,7 +305,7 @@ dependencies = [
"heck",
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -700,7 +700,7 @@ checksum = "89ca545a94061b6365f2c7355b4b32bd20df3ff95f02da9329b34ccc3bd6ee72"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -1037,9 +1037,9 @@ dependencies = [
[[package]]
name = "memchr"
version = "2.6.2"
version = "2.6.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5486aed0026218e61b8a01d5fbd5a0a134649abb71a0e53b7bc088529dced86e"
checksum = "8f232d6ef707e1956a43342693d2a31e72989554d58299d7a88738cc95b0d35c"
[[package]]
name = "memoffset"
@ -1182,9 +1182,9 @@ dependencies = [
[[package]]
name = "object"
version = "0.32.0"
version = "0.32.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "77ac5bbd07aea88c60a577a1ce218075ffd59208b2d7ca97adf9bfc5aeb21ebe"
checksum = "9cf5f9dd3933bd50a9e1f149ec995f39ae2c496d31fd772c1fd45ebc27e902b0"
dependencies = [
"memchr",
]
@ -1281,7 +1281,7 @@ checksum = "4359fd9c9171ec6e8c62926d6faaf553a8dc3f64e1507e76da7911b4f6a04405"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -1459,9 +1459,9 @@ dependencies = [
[[package]]
name = "regex"
version = "1.9.4"
version = "1.9.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "12de2eff854e5fa4b1295edd650e227e9d8fb0c9e90b12e7f36d6a6811791a29"
checksum = "697061221ea1b4a94a624f67d0ae2bfe4e22b8a17b6a192afb11046542cc8c47"
dependencies = [
"aho-corasick",
"memchr",
@ -1471,9 +1471,9 @@ dependencies = [
[[package]]
name = "regex-automata"
version = "0.3.7"
version = "0.3.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "49530408a136e16e5b486e883fbb6ba058e8e4e8ae6621a77b048b314336e629"
checksum = "c2f401f4955220693b56f8ec66ee9c78abffd8d1c4f23dc41a23839eb88f0795"
dependencies = [
"aho-corasick",
"memchr",
@ -1542,7 +1542,7 @@ checksum = "4eca7ac642d82aa35b60049a6eccb4be6be75e599bd2e9adb5f875a737654af2"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -1706,7 +1706,7 @@ dependencies = [
"proc-macro2",
"quote",
"rustversion",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -1722,9 +1722,9 @@ dependencies = [
[[package]]
name = "syn"
version = "2.0.29"
version = "2.0.31"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c324c494eba9d92503e6f1ef2e6df781e78f6a7705a0202d9801b198807d518a"
checksum = "718fa2415bcb8d8bd775917a1bf12a7931b6dfa890753378538118181e0cb398"
dependencies = [
"proc-macro2",
"quote",
@ -1754,22 +1754,22 @@ dependencies = [
[[package]]
name = "thiserror"
version = "1.0.47"
version = "1.0.48"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "97a802ec30afc17eee47b2855fc72e0c4cd62be9b4efe6591edde0ec5bd68d8f"
checksum = "9d6d7a740b8a666a7e828dd00da9c0dc290dff53154ea77ac109281de90589b7"
dependencies = [
"thiserror-impl",
]
[[package]]
name = "thiserror-impl"
version = "1.0.47"
version = "1.0.48"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6bb623b56e39ab7dcd4b1b98bb6c8f8d907ed255b18de254088016b27a8ee19b"
checksum = "49922ecae66cc8a249b77e68d1d0623c1b2c514f0060c27cdc68bd62a1219d35"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -1863,7 +1863,7 @@ checksum = "630bdcf245f78637c13ec01ffae6187cca34625e8c63150d424b59e55af2675e"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -1943,7 +1943,7 @@ checksum = "5f4f31f56159e98206da9efd823404b79b6ef3143b4a7ab76e67b1751b25a4ab"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
]
[[package]]
@ -2109,9 +2109,9 @@ checksum = "49874b5167b65d7193b8aba1567f5c7d93d001cafc34600cee003eda787e483f"
[[package]]
name = "walkdir"
version = "2.3.3"
version = "2.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "36df944cda56c7d8d8b7496af378e6b16de9284591917d307c9b4d313c44e698"
checksum = "d71d857dc86794ca4c280d616f7da00d2dbfd8cd788846559a6813e6aa4b54ee"
dependencies = [
"same-file",
"winapi-util",
@ -2144,7 +2144,7 @@ dependencies = [
"once_cell",
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
"wasm-bindgen-shared",
]
@ -2166,7 +2166,7 @@ checksum = "54681b18a46765f095758388f2d0cf16eb8d4169b639ab575a8f5693af210c7b"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.29",
"syn 2.0.31",
"wasm-bindgen-backend",
"wasm-bindgen-shared",
]

View file

@ -1,9 +1,10 @@
use std::{process::Stdio, time::Duration};
use std::process::Stdio;
use anyhow::Result;
use shared::{StreamBuf, Url};
use tokio::{io::{AsyncBufReadExt, BufReader}, process::Command, sync::mpsc};
use tokio_stream::wrappers::UnboundedReceiverStream;
use shared::Url;
use tokio::{io::{AsyncBufReadExt, BufReader}, process::Command, sync::mpsc::{self, UnboundedReceiver}};
use crate::files::File;
pub struct FdOpt {
pub cwd: Url,
@ -12,7 +13,7 @@ pub struct FdOpt {
pub subject: String,
}
pub fn fd(opt: FdOpt) -> Result<StreamBuf<UnboundedReceiverStream<Url>>> {
pub fn fd(opt: FdOpt) -> Result<UnboundedReceiver<File>> {
let mut child = Command::new("fd")
.arg("--base-directory")
.arg(&opt.cwd)
@ -28,11 +29,12 @@ pub fn fd(opt: FdOpt) -> Result<StreamBuf<UnboundedReceiverStream<Url>>> {
let mut it = BufReader::new(child.stdout.take().unwrap()).lines();
let (tx, rx) = mpsc::unbounded_channel();
let rx = StreamBuf::new(UnboundedReceiverStream::new(rx), Duration::from_millis(300));
tokio::spawn(async move {
while let Ok(Some(line)) = it.next_line().await {
tx.send(opt.cwd.join(line)).ok();
if let Ok(file) = File::from(opt.cwd.join(line)).await {
tx.send(file).ok();
}
}
child.wait().await.ok();
});

View file

@ -1,9 +1,10 @@
use std::{process::Stdio, time::Duration};
use std::process::Stdio;
use anyhow::Result;
use shared::{StreamBuf, Url};
use tokio::{io::{AsyncBufReadExt, BufReader}, process::Command, sync::mpsc};
use tokio_stream::wrappers::UnboundedReceiverStream;
use shared::Url;
use tokio::{io::{AsyncBufReadExt, BufReader}, process::Command, sync::mpsc::{self, UnboundedReceiver}};
use crate::files::File;
pub struct RgOpt {
pub cwd: Url,
@ -11,7 +12,7 @@ pub struct RgOpt {
pub subject: String,
}
pub fn rg(opt: RgOpt) -> Result<StreamBuf<UnboundedReceiverStream<Url>>> {
pub fn rg(opt: RgOpt) -> Result<UnboundedReceiver<File>> {
let mut child = Command::new("rg")
.current_dir(&opt.cwd)
.args(["--color=never", "--files-with-matches", "--smart-case"])
@ -26,11 +27,12 @@ pub fn rg(opt: RgOpt) -> Result<StreamBuf<UnboundedReceiverStream<Url>>> {
let mut it = BufReader::new(child.stdout.take().unwrap()).lines();
let (tx, rx) = mpsc::unbounded_channel();
let rx = StreamBuf::new(UnboundedReceiverStream::new(rx), Duration::from_millis(300));
tokio::spawn(async move {
while let Ok(Some(line)) = it.next_line().await {
tx.send(opt.cwd.join(line)).ok();
if let Ok(file) = File::from(opt.cwd.join(line)).await {
tx.send(file).ok();
}
}
child.wait().await.ok();
});

View file

@ -1,9 +1,9 @@
use std::{collections::{BTreeMap, BTreeSet}, mem, ops::Deref};
use std::{collections::{BTreeMap, BTreeSet}, mem, ops::Deref, time::Duration};
use anyhow::Result;
use config::manager::SortBy;
use shared::Url;
use tokio::fs;
use tokio::{fs, select, sync::mpsc::{self, UnboundedReceiver}, time::sleep};
use super::{File, FilesSorter};
@ -40,25 +40,22 @@ impl Deref for Files {
}
impl Files {
pub async fn read(urls: Vec<Url>) -> Vec<File> {
let mut items = Vec::with_capacity(urls.len());
for url in urls {
if let Ok(file) = File::from(url).await {
items.push(file);
}
}
items
}
pub async fn read_dir(url: &Url) -> Result<Vec<File>> {
pub async fn from_dir(url: &Url) -> Result<UnboundedReceiver<File>> {
let mut it = fs::read_dir(url).await?;
let mut items = Vec::new();
let (tx, rx) = mpsc::unbounded_channel();
let url = url.clone();
tokio::spawn(async move {
while let Ok(Some(item)) = it.next_entry().await {
if let Ok(meta) = item.metadata().await {
items.push(File::from_meta(Url::new(item.path(), url), meta).await);
select! {
_ = tx.closed() => break,
Ok(meta) = item.metadata() => {
tx.send(File::from_meta(Url::new(item.path(), &url), meta).await).ok();
}
}
Ok(items)
}
});
Ok(rx)
}
}
@ -111,25 +108,7 @@ impl Files {
applied
}
pub fn update_read(&mut self, mut items: Vec<File>) -> bool {
if !self.show_hidden {
(self.hidden, items) = items.into_iter().partition(|f| f.is_hidden);
}
self.sorter.sort(&mut items);
self.items = items;
true
}
pub fn update_size(&mut self, items: BTreeMap<Url, u64>) -> bool {
self.sizes.extend(items);
if self.sorter.by == SortBy::Size {
self.sorter.sort(&mut self.items);
}
true
}
// TODO: remove this
pub fn update_search(&mut self, items: Vec<File>) -> bool {
pub fn update_read(&mut self, items: Vec<File>) -> bool {
if !items.is_empty() {
if self.show_hidden {
self.items.extend(items);
@ -148,9 +127,16 @@ impl Files {
self.hidden.clear();
return true;
}
false
}
pub fn update_size(&mut self, items: BTreeMap<Url, u64>) -> bool {
self.sizes.extend(items);
if self.sorter.by == SortBy::Size {
self.sorter.sort(&mut self.items);
}
true
}
}
impl Files {

View file

@ -330,7 +330,7 @@ impl Manager {
let cwd = self.cwd().to_owned();
let hovered = self.hovered().map(|h| h.url_owned());
let mut b = if cwd == url && !cwd.is_search() {
let mut b = if cwd == url {
self.current_mut().update(op)
} else if matches!(self.parent(), Some(p) if p.cwd == url) {
self.active_mut().parent.as_mut().unwrap().update(op)

View file

@ -1,9 +1,10 @@
use std::sync::atomic::Ordering;
use std::{sync::atomic::Ordering, time::Duration};
use adaptor::Adaptor;
use config::MANAGER;
use shared::{MimeKind, PeekError, Url, MIME_DIR};
use tokio::task::JoinHandle;
use tokio::{pin, task::JoinHandle};
use tokio_stream::{wrappers::UnboundedReceiverStream, StreamExt};
use super::{Provider, INCR};
use crate::{emit, files::{Files, FilesOp}};
@ -70,32 +71,36 @@ impl Preview {
}
self.reset(|_| true);
if files.is_some() || sequent {
emit!(Preview(PreviewLock {
url: url.clone(),
mime: MIME_DIR.to_owned(),
skip: self.skip,
data: PreviewData::Folder,
}));
}
if sequent {
return;
}
let (url, skip) = (url.clone(), self.skip);
let url = url.clone();
self.handle = Some(tokio::spawn(async move {
emit!(Files(match Files::read_dir(&url).await {
Ok(items) => FilesOp::Read(url.clone(), items),
Err(_) => FilesOp::IOErr(url.clone()),
}));
let Ok(rx) = Files::from_dir(&url).await else {
emit!(Files(FilesOp::IOErr(url)));
return;
};
emit!(Preview(PreviewLock {
url,
mime: MIME_DIR.to_owned(),
skip,
data: PreviewData::Folder,
}));
let rx = UnboundedReceiverStream::new(rx).chunks_timeout(10000, Duration::from_millis(500));
pin!(rx);
let mut first = false;
while let Some(chunk) = rx.next().await {
if first {
emit!(Files(FilesOp::clear(&url)));
first = false;
}
emit!(Files(FilesOp::Read(url.clone(), chunk)));
}
}));
}

View file

@ -1,13 +1,13 @@
use std::{borrow::Cow, collections::{BTreeMap, BTreeSet}, ffi::{OsStr, OsString}, mem};
use std::{borrow::Cow, collections::{BTreeMap, BTreeSet}, ffi::{OsStr, OsString}, mem, time::Duration};
use anyhow::{Error, Result};
use config::open::Opener;
use futures::StreamExt;
use shared::{Defer, Url};
use tokio::task::JoinHandle;
use tokio::{pin, task::JoinHandle};
use tokio_stream::{wrappers::UnboundedReceiverStream, StreamExt};
use super::{Folder, Mode, Preview, PreviewLock};
use crate::{emit, external::{self, FzfOpt, ZoxideOpt}, files::{File, Files, FilesOp, FilesSorter}, input::InputOpt, Event, BLOCKER};
use crate::{emit, external::{self, FzfOpt, ZoxideOpt}, files::{File, FilesOp, FilesSorter}, input::InputOpt, Event, BLOCKER};
pub struct Tab {
pub(super) mode: Mode,
@ -237,21 +237,24 @@ impl Tab {
handle.abort();
}
let cwd = self.current.cwd.clone();
let cwd = self.current.cwd.to_search();
let hidden = self.show_hidden;
self.search = Some(tokio::spawn(async move {
let subject = emit!(Input(InputOpt::top("Search:"))).await?;
let mut rx = if grep {
let rx = if grep {
external::rg(external::RgOpt { cwd: cwd.clone(), hidden, subject })
} else {
external::fd(external::FdOpt { cwd: cwd.clone(), hidden, glob: false, subject })
}?;
let rx = UnboundedReceiverStream::new(rx).chunks_timeout(1000, Duration::from_millis(300));
pin!(rx);
emit!(Files(FilesOp::clear(&cwd)));
while let Some(chunk) = rx.next().await {
emit!(Files(FilesOp::Read(cwd.clone(), Files::read(chunk).await)));
emit!(Files(FilesOp::Read(cwd.clone(), chunk)));
}
Ok(())
}));

View file

@ -1,14 +1,13 @@
use std::{collections::BTreeSet, sync::Arc, time::Duration};
use futures::StreamExt;
use indexmap::IndexMap;
use notify::{event::{MetadataKind, ModifyKind}, EventKind, RecommendedWatcher, RecursiveMode, Watcher as _Watcher};
use parking_lot::RwLock;
use shared::{StreamBuf, Url};
use tokio::{fs, sync::mpsc};
use tokio_stream::wrappers::UnboundedReceiverStream;
use shared::Url;
use tokio::{fs, pin, sync::mpsc::{self, UnboundedReceiver}};
use tokio_stream::{wrappers::UnboundedReceiverStream, StreamExt};
use crate::{emit, external, files::{Files, FilesOp}};
use crate::{emit, external, files::{File, Files, FilesOp}};
pub struct Watcher {
watcher: RecommendedWatcher,
@ -18,8 +17,6 @@ pub struct Watcher {
impl Watcher {
pub(super) fn start() -> Self {
let (tx, rx) = mpsc::unbounded_channel();
let rx = StreamBuf::new(UnboundedReceiverStream::new(rx), Duration::from_millis(300));
let watcher = RecommendedWatcher::new(
{
let tx = tx.clone();
@ -131,10 +128,10 @@ impl Watcher {
});
}
async fn changed(
mut rx: StreamBuf<UnboundedReceiverStream<Url>>,
watched: Arc<RwLock<IndexMap<Url, Option<Url>>>>,
) {
async fn changed(rx: UnboundedReceiver<Url>, watched: Arc<RwLock<IndexMap<Url, Option<Url>>>>) {
let rx = UnboundedReceiverStream::new(rx).chunks_timeout(100, Duration::from_millis(200));
pin!(rx);
while let Some(paths) = rx.next().await {
let (mut files, mut dirs): (Vec<_>, Vec<_>) = Default::default();
for path in paths.into_iter().collect::<BTreeSet<_>>() {
@ -163,34 +160,47 @@ impl Watcher {
}
async fn dir_changed(url: &Url, watched: Arc<RwLock<IndexMap<Url, Option<Url>>>>) {
let linked = watched
let linked: Vec<_> = watched
.read()
.iter()
.map_while(|(k, v)| v.as_ref().and_then(|v| url.strip_prefix(v)).map(|v| k.join(v)))
.collect::<Vec<_>>();
let result = Files::read_dir(url).await;
if linked.is_empty() {
emit!(Files(match result {
Ok(items) => FilesOp::Read(url.clone(), items),
Err(_) => FilesOp::IOErr(url.clone()),
}));
return;
}
.collect();
let Ok(rx) = Files::from_dir(url).await else {
emit!(Files(FilesOp::IOErr(url.clone())));
for ori in linked {
emit!(Files(match &result {
Ok(items) => {
let mut files = Vec::with_capacity(items.len());
for item in items {
let mut file = item.clone();
file.set_url(ori.join(item.url().strip_prefix(url).unwrap()));
files.push(file);
emit!(Files(FilesOp::IOErr(ori)));
}
FilesOp::Read(ori, files)
return;
};
let rx = UnboundedReceiverStream::new(rx).chunks_timeout(10000, Duration::from_millis(500));
pin!(rx);
let linked_files = |files: &[File], ori: &Url| -> Vec<File> {
let mut new = Vec::with_capacity(files.len());
for file in files {
let mut file = file.clone();
file.set_url(ori.join(file.url().strip_prefix(url).unwrap()));
new.push(file);
}
Err(_) => FilesOp::IOErr(ori),
}));
new
};
let mut first = true;
while let Some(chunk) = rx.next().await {
if first {
emit!(Files(FilesOp::clear(url)));
for ori in &linked {
emit!(Files(FilesOp::clear(ori)));
}
first = false;
}
for ori in &linked {
emit!(Files(FilesOp::Read(ori.clone(), linked_files(&chunk, ori))));
}
emit!(Files(FilesOp::Read(url.clone(), chunk)));
}
}
}

View file

@ -7,7 +7,6 @@ mod fns;
mod fs;
mod mime;
mod ro_cell;
mod stream;
mod term;
mod throttle;
mod time;
@ -20,7 +19,6 @@ pub use fns::*;
pub use fs::*;
pub use mime::*;
pub use ro_cell::*;
pub use stream::*;
pub use term::*;
pub use throttle::*;
pub use time::*;

View file

@ -1,69 +0,0 @@
use std::{mem, pin::Pin, task::{Context, Poll}, time::Duration};
use futures::{FutureExt, Stream, StreamExt};
use tokio::time::{sleep, Instant, Sleep};
pub struct StreamBuf<S>
where
S: Stream,
{
stream: S,
interval: Duration,
sleep: Sleep,
pending: Vec<S::Item>,
}
impl<S> Unpin for StreamBuf<S> where S: Stream {}
impl<S> StreamBuf<S>
where
S: Stream + Unpin,
{
pub fn new(stream: S, interval: Duration) -> StreamBuf<S> {
Self { stream, interval, sleep: sleep(Duration::ZERO), pending: Default::default() }
}
pub fn flush(&mut self) {
let sleep = unsafe { Pin::new_unchecked(&mut self.sleep) };
sleep.reset(Instant::now());
}
}
impl<S> Stream for StreamBuf<S>
where
S: Stream + Unpin,
{
type Item = Vec<S::Item>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let (mut stream, interval, mut sleep, pending) = unsafe {
let me = self.get_unchecked_mut();
(Pin::new(&mut me.stream), me.interval, Pin::new_unchecked(&mut me.sleep), &mut me.pending)
};
while let Poll::Ready(next) = stream.poll_next_unpin(cx) {
match next {
Some(next) => pending.push(next),
None => {
if pending.is_empty() {
return Poll::Ready(None);
}
break;
}
}
}
if pending.is_empty() {
return Poll::Pending;
}
match sleep.poll_unpin(cx) {
Poll::Ready(_) => {
sleep.reset(Instant::now() + interval);
Poll::Ready(Some(mem::take(pending)))
}
Poll::Pending => Poll::Pending,
}
}
}

View file

@ -88,6 +88,13 @@ impl Url {
#[inline]
pub fn is_search(&self) -> bool { self.scheme == UrlScheme::Search }
#[inline]
pub fn to_search(&self) -> Self {
let mut url = self.clone();
url.scheme = UrlScheme::Search;
url
}
#[inline]
pub fn is_archive(&self) -> bool { self.scheme == UrlScheme::Archive }