diff --git a/yazi-dds/src/pump.rs b/yazi-dds/src/pump.rs index 335ae521..c2d20410 100644 --- a/yazi-dds/src/pump.rs +++ b/yazi-dds/src/pump.rs @@ -1,9 +1,10 @@ -use std::time::Duration; +use std::{collections::HashMap, time::Duration}; use parking_lot::Mutex; use tokio::{pin, select, sync::mpsc}; use tokio_stream::{StreamExt, wrappers::UnboundedReceiverStream}; use tokio_util::sync::CancellationToken; +use yazi_fs::{File, FilesOp}; use yazi_shared::{RoCell, url::Url}; use crate::{Pubsub, body::BodyMoveItem}; @@ -12,6 +13,7 @@ static CT: RoCell = RoCell::new(); static MOVE_TX: Mutex>> = Mutex::new(None); static TRASH_TX: Mutex>> = Mutex::new(None); static DELETE_TX: Mutex>> = Mutex::new(None); +static BULK_RENAME_TX: Mutex>> = Mutex::new(None); pub struct Pump; @@ -37,15 +39,24 @@ impl Pump { } } + #[inline] + pub fn push_bulk_rename_pair(from: Url, to: File) { + if let Some(tx) = &*BULK_RENAME_TX.lock() { + tx.send((from, to)).ok(); + } + } + pub(super) fn serve() { let (move_tx, move_rx) = mpsc::unbounded_channel(); let (trash_tx, trash_rx) = mpsc::unbounded_channel(); let (delete_tx, delete_rx) = mpsc::unbounded_channel(); + let (bulk_rename_tx, bulk_rename_rx) = mpsc::unbounded_channel(); CT.with(<_>::default); MOVE_TX.lock().replace(move_tx); TRASH_TX.lock().replace(trash_tx); DELETE_TX.lock().replace(delete_tx); + BULK_RENAME_TX.lock().replace(bulk_rename_tx); tokio::spawn(async move { let move_rx = @@ -54,16 +65,23 @@ impl Pump { UnboundedReceiverStream::new(trash_rx).chunks_timeout(1000, Duration::from_millis(500)); let delete_rx = UnboundedReceiverStream::new(delete_rx).chunks_timeout(1000, Duration::from_millis(500)); + let bulk_rename_rx = UnboundedReceiverStream::new(bulk_rename_rx) + .chunks_timeout(1000, Duration::from_millis(500)); pin!(move_rx); pin!(trash_rx); pin!(delete_rx); + pin!(bulk_rename_rx); loop { select! { Some(items) = move_rx.next() => Pubsub::pub_from_move(items), Some(urls) = trash_rx.next() => Pubsub::pub_from_trash(urls), Some(urls) = delete_rx.next() => Pubsub::pub_from_delete(urls), + Some(items) = bulk_rename_rx.next() => { + Pubsub::pub_from_bulk(items.iter().map(|(from, to)| (from, &to.url)).collect()); + FilesOp::rename(HashMap::from_iter(items)); + } else => { CT.cancel(); break; @@ -77,6 +95,7 @@ impl Pump { drop(MOVE_TX.lock().take()); drop(TRASH_TX.lock().take()); drop(DELETE_TX.lock().take()); + drop(BULK_RENAME_TX.lock().take()); CT.cancelled().await; } } diff --git a/yazi-scheduler/src/scheduler.rs b/yazi-scheduler/src/scheduler.rs index 6b3e9fa6..5be7db8b 100644 --- a/yazi-scheduler/src/scheduler.rs +++ b/yazi-scheduler/src/scheduler.rs @@ -210,13 +210,32 @@ impl Scheduler { } pub fn file_rename_at(&self, root: &Path, old: &Path, new: &Path) { - let id = self.ongoing.lock().add( + let mut ongoing = self.ongoing.lock(); + let id = ongoing.add( TaskKind::User, format!("Rename at {}: {} -> {} ", root.display(), old.display(), new.display()), ); let (from, to): (Url, Url) = (root.join(old).into(), root.join(new).into()); + ongoing.hooks.insert(id, { + let from = from.clone(); + let to = to.clone(); + let ongoing = self.ongoing.clone(); + + Box::new(move |canceled: bool| { + async move { + if !canceled { + if let Ok(to) = yazi_fs::File::from(to).await { + Pump::push_bulk_rename_pair(from, to); + } + } + ongoing.lock().try_remove(id, TaskStage::Hooked); + } + .boxed() + }) + }); + let file = self.file.clone(); self.send_micro(id, LOW, async move { file.rename(FileOpRename { id, from, to }).await }); }