From 48a87c40cb326069e0486e4f01d93e7a1f58ba3e Mon Sep 17 00:00:00 2001 From: sxyazi Date: Fri, 29 Dec 2023 00:55:35 +0800 Subject: [PATCH] .. --- yazi-scheduler/src/scheduler.rs | 54 +++++++++++++-------------- yazi-scheduler/src/task.rs | 2 +- yazi-scheduler/src/workers/file.rs | 38 +++++++++---------- yazi-scheduler/src/workers/plugin.rs | 16 ++++---- yazi-scheduler/src/workers/preload.rs | 18 ++++----- yazi-scheduler/src/workers/process.rs | 20 +++++----- 6 files changed, 74 insertions(+), 74 deletions(-) diff --git a/yazi-scheduler/src/scheduler.rs b/yazi-scheduler/src/scheduler.rs index dffda18a..5773f1fc 100644 --- a/yazi-scheduler/src/scheduler.rs +++ b/yazi-scheduler/src/scheduler.rs @@ -6,7 +6,7 @@ use tokio::{fs, select, sync::{mpsc::{self, UnboundedReceiver}, oneshot}}; use yazi_config::{open::Opener, plugin::PluginRule, TASKS}; use yazi_shared::{emit, event::Exec, fs::{unique_path, Url}, Layer, Throttle}; -use super::{Running, TaskOp, TaskStage}; +use super::{Running, TaskProg, TaskStage}; use crate::{workers::{File, FileOpDelete, FileOpLink, FileOpPaste, FileOpTrash, Plugin, PluginOpEntry, Preload, PreloadOpRule, PreloadOpSize, Process, ProcessOpOpen}, TaskKind}; pub struct Scheduler { @@ -15,14 +15,14 @@ pub struct Scheduler { pub preload: Arc, pub process: Arc, - todo: async_channel::Sender>, - prog: mpsc::UnboundedSender, + micro: async_channel::Sender>, + prog: mpsc::UnboundedSender, pub running: Arc>, } impl Scheduler { pub fn start() -> Self { - let (todo_tx, todo_rx) = async_channel::unbounded(); + let (micro_tx, micro_rx) = async_channel::unbounded(); let (prog_tx, prog_rx) = mpsc::unbounded_channel(); let scheduler = Self { @@ -31,16 +31,16 @@ impl Scheduler { preload: Arc::new(Preload::new(prog_tx.clone())), process: Arc::new(Process::new(prog_tx.clone())), - todo: todo_tx, + micro: micro_tx, prog: prog_tx, running: Default::default(), }; for _ in 0..TASKS.micro_workers { - scheduler.schedule_micro(todo_rx.clone()); + scheduler.schedule_micro(micro_rx.clone()); } for _ in 0..TASKS.macro_workers { - scheduler.schedule_macro(todo_rx.clone()); + scheduler.schedule_macro(micro_rx.clone()); } scheduler.progress(prog_rx); scheduler @@ -79,7 +79,7 @@ impl Scheduler { continue; } if let Err(e) = file.work(&mut op).await { - prog.send(TaskOp::Fail(id, format!("Failed to work on this task: {:?}", e))).ok(); + prog.send(TaskProg::Fail(id, format!("Failed to work on this task: {:?}", e))).ok(); } } Ok((id, mut op)) = plugin.recv() => { @@ -87,7 +87,7 @@ impl Scheduler { continue; } if let Err(e) = plugin.work(&mut op).await { - prog.send(TaskOp::Fail(id, format!("Failed to work on this task: {:?}", e))).ok(); + prog.send(TaskProg::Fail(id, format!("Failed to work on this task: {:?}", e))).ok(); } } } @@ -95,20 +95,20 @@ impl Scheduler { }); } - fn progress(&self, mut rx: UnboundedReceiver) { - let todo = self.todo.clone(); + fn progress(&self, mut rx: UnboundedReceiver) { + let micro = self.micro.clone(); let running = self.running.clone(); tokio::spawn(async move { while let Some(op) = rx.recv().await { match op { - TaskOp::New(id, size) => { + TaskProg::New(id, size) => { if let Some(task) = running.write().get_mut(id) { task.total += 1; task.found += size; } } - TaskOp::Adv(id, succ, processed) => { + TaskProg::Adv(id, succ, processed) => { let mut running = running.write(); if let Some(task) = running.get_mut(id) { task.succ += succ; @@ -116,16 +116,16 @@ impl Scheduler { } if succ > 0 { if let Some(fut) = running.try_remove(id, TaskStage::Pending) { - todo.send_blocking(fut).ok(); + micro.send_blocking(fut).ok(); } } } - TaskOp::Succ(id) => { + TaskProg::Succ(id) => { if let Some(fut) = running.write().try_remove(id, TaskStage::Dispatched) { - todo.send_blocking(fut).ok(); + micro.send_blocking(fut).ok(); } } - TaskOp::Fail(id, reason) => { + TaskProg::Fail(id, reason) => { if let Some(task) = running.write().get_mut(id) { task.fail += 1; task.logs.push_str(&reason); @@ -136,7 +136,7 @@ impl Scheduler { } } } - TaskOp::Log(id, line) => { + TaskProg::Log(id, line) => { if let Some(task) = running.write().get_mut(id) { task.logs.push_str(&line); task.logs.push('\n'); @@ -156,7 +156,7 @@ impl Scheduler { let b = running.all.remove(&id).is_some(); if let Some(hook) = running.hooks.remove(&id) { - self.todo.send_blocking(hook(true)).ok(); + self.micro.send_blocking(hook(true)).ok(); } b } @@ -193,7 +193,7 @@ impl Scheduler { }) }); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let file = self.file.clone(); async move { if !force { @@ -209,7 +209,7 @@ impl Scheduler { let name = format!("Copy {:?} to {:?}", from, to); let id = self.running.write().add(TaskKind::User, name); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let file = self.file.clone(); async move { if !force { @@ -225,7 +225,7 @@ impl Scheduler { let name = format!("Link {from:?} to {to:?}"); let id = self.running.write().add(TaskKind::User, name); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let file = self.file.clone(); async move { if !force { @@ -259,7 +259,7 @@ impl Scheduler { }) }); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let file = self.file.clone(); async move { file.delete(FileOpDelete { id, target, length: 0 }).await.ok(); @@ -272,7 +272,7 @@ impl Scheduler { let name = format!("Trash {:?}", target); let id = self.running.write().add(TaskKind::User, name); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let file = self.file.clone(); async move { file.trash(FileOpTrash { id, target, length: 0 }).await.ok(); @@ -284,7 +284,7 @@ impl Scheduler { pub fn plugin_micro(&self, name: String) { let id = self.running.write().add(TaskKind::User, format!("Run micro plugin `{name}`")); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let plugin = self.plugin.clone(); async move { plugin.micro(PluginOpEntry { id, name }).await.ok(); @@ -305,7 +305,7 @@ impl Scheduler { format!("Run preloader `{}` with {} target(s)", rule.exec.cmd, targets.len()), ); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let preload = self.preload.clone(); let (rule_id, rule_multi) = (rule.id, rule.multi); @@ -323,7 +323,7 @@ impl Scheduler { let mut running = self.running.write(); for target in targets { let id = running.add(TaskKind::Preload, format!("Calculate the size of {:?}", target)); - _ = self.todo.send_blocking({ + _ = self.micro.send_blocking({ let preload = self.preload.clone(); let target = target.clone(); let throttle = throttle.clone(); diff --git a/yazi-scheduler/src/task.rs b/yazi-scheduler/src/task.rs index 6c404453..27baceea 100644 --- a/yazi-scheduler/src/task.rs +++ b/yazi-scheduler/src/task.rs @@ -59,7 +59,7 @@ impl From<&Task> for TaskSummary { } #[derive(Debug)] -pub enum TaskOp { +pub enum TaskProg { // id, size New(usize, u64), // id, processed, size diff --git a/yazi-scheduler/src/workers/file.rs b/yazi-scheduler/src/workers/file.rs index 5e3cd81e..0ceed7e8 100644 --- a/yazi-scheduler/src/workers/file.rs +++ b/yazi-scheduler/src/workers/file.rs @@ -7,13 +7,13 @@ use tracing::warn; use yazi_config::TASKS; use yazi_shared::fs::{calculate_size, copy_with_progress, path_relative_to, Url}; -use crate::TaskOp; +use crate::TaskProg; pub struct File { tx: async_channel::Sender, rx: async_channel::Receiver, - prog: mpsc::UnboundedSender, + prog: mpsc::UnboundedSender, } #[derive(Debug)] @@ -60,7 +60,7 @@ pub struct FileOpTrash { } impl File { - pub fn new(prog: mpsc::UnboundedSender) -> Self { + pub fn new(prog: mpsc::UnboundedSender) -> Self { let (tx, rx) = async_channel::unbounded(); Self { tx, rx, prog } } @@ -92,7 +92,7 @@ impl File { } break; } - Ok(n) => self.prog.send(TaskOp::Adv(task.id, 0, n))?, + Ok(n) => self.prog.send(TaskProg::Adv(task.id, 0, n))?, Err(e) if e.kind() == NotFound => { warn!("Paste task partially done: {:?}", task); break; @@ -110,7 +110,7 @@ impl File { Err(e) => Err(e)?, } } - self.prog.send(TaskOp::Adv(task.id, 1, 0))?; + self.prog.send(TaskProg::Adv(task.id, 1, 0))?; } FileOp::Link(task) => { let meta = task.meta.as_ref().unwrap(); @@ -120,7 +120,7 @@ impl File { Ok(p) => Cow::Owned(p), Err(e) if e.kind() == NotFound => { self.log(task.id, format!("Link task partially done: {:?}", task))?; - return Ok(self.prog.send(TaskOp::Adv(task.id, 1, meta.len()))?); + return Ok(self.prog.send(TaskProg::Adv(task.id, 1, meta.len()))?); } Err(e) => Err(e)?, } @@ -155,7 +155,7 @@ impl File { if task.delete { fs::remove_file(&task.from).await.ok(); } - self.prog.send(TaskOp::Adv(task.id, 1, meta.len()))?; + self.prog.send(TaskProg::Adv(task.id, 1, meta.len()))?; } FileOp::Delete(task) => { if let Err(e) = fs::remove_file(&task.target).await { @@ -164,7 +164,7 @@ impl File { Err(e)? } } - self.prog.send(TaskOp::Adv(task.id, 1, task.length))? + self.prog.send(TaskProg::Adv(task.id, 1, task.length))? } FileOp::Trash(task) => { #[cfg(target_os = "macos")] @@ -178,7 +178,7 @@ impl File { { trash::delete(&task.target)?; } - self.prog.send(TaskOp::Adv(task.id, 1, task.length))?; + self.prog.send(TaskProg::Adv(task.id, 1, task.length))?; } } Ok(()) @@ -196,7 +196,7 @@ impl File { let meta = Self::metadata(&task.from, task.follow).await?; if !meta.is_dir() { let id = task.id; - self.prog.send(TaskOp::New(id, meta.len()))?; + self.prog.send(TaskProg::New(id, meta.len()))?; if meta.is_file() { self.tx.send(FileOp::Paste(task)).await?; @@ -211,7 +211,7 @@ impl File { match $result { Ok(v) => v, Err(e) => { - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; self.fail(task.id, format!("An error occurred while pasting: {e}"))?; continue; } @@ -242,7 +242,7 @@ impl File { task.to = dest.join(src.file_name().unwrap()); task.from = src; - self.prog.send(TaskOp::New(task.id, meta.len()))?; + self.prog.send(TaskProg::New(task.id, meta.len()))?; if meta.is_file() { self.tx.send(FileOp::Paste(task.clone())).await?; @@ -260,7 +260,7 @@ impl File { task.meta = Some(fs::symlink_metadata(&task.from).await?); } - self.prog.send(TaskOp::New(id, task.meta.as_ref().unwrap().len()))?; + self.prog.send(TaskProg::New(id, task.meta.as_ref().unwrap().len()))?; self.tx.send(FileOp::Link(task)).await?; self.succ(id) } @@ -270,7 +270,7 @@ impl File { if !meta.is_dir() { let id = task.id; task.length = meta.len(); - self.prog.send(TaskOp::New(id, meta.len()))?; + self.prog.send(TaskProg::New(id, meta.len()))?; self.tx.send(FileOp::Delete(task)).await?; return self.succ(id); } @@ -295,7 +295,7 @@ impl File { task.target = Url::from(entry.path()); task.length = meta.len(); - self.prog.send(TaskOp::New(task.id, meta.len()))?; + self.prog.send(TaskProg::New(task.id, meta.len()))?; self.tx.send(FileOp::Delete(task.clone())).await?; } } @@ -306,7 +306,7 @@ impl File { let id = task.id; task.length = calculate_size(&task.target).await; - self.prog.send(TaskOp::New(id, task.length))?; + self.prog.send(TaskProg::New(id, task.length))?; self.tx.send(FileOp::Trash(task)).await?; self.succ(id) } @@ -343,16 +343,16 @@ impl File { impl File { #[inline] - fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) } + fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskProg::Succ(id))?) } #[inline] fn fail(&self, id: usize, reason: String) -> Result<()> { - Ok(self.prog.send(TaskOp::Fail(id, reason))?) + Ok(self.prog.send(TaskProg::Fail(id, reason))?) } #[inline] fn log(&self, id: usize, line: String) -> Result<()> { - Ok(self.prog.send(TaskOp::Log(id, line))?) + Ok(self.prog.send(TaskProg::Log(id, line))?) } } diff --git a/yazi-scheduler/src/workers/plugin.rs b/yazi-scheduler/src/workers/plugin.rs index 150c45f6..c28edd52 100644 --- a/yazi-scheduler/src/workers/plugin.rs +++ b/yazi-scheduler/src/workers/plugin.rs @@ -1,13 +1,13 @@ use anyhow::Result; use tokio::sync::mpsc; -use crate::TaskOp; +use crate::TaskProg; pub struct Plugin { tx: async_channel::Sender, rx: async_channel::Receiver, - prog: mpsc::UnboundedSender, + prog: mpsc::UnboundedSender, } #[derive(Debug)] @@ -22,7 +22,7 @@ pub struct PluginOpEntry { } impl Plugin { - pub fn new(prog: mpsc::UnboundedSender) -> Self { + pub fn new(prog: mpsc::UnboundedSender) -> Self { let (tx, rx) = async_channel::unbounded(); Self { tx, rx, prog } } @@ -44,21 +44,21 @@ impl Plugin { } pub async fn micro(&self, task: PluginOpEntry) -> Result<()> { - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; if let Err(e) = yazi_plugin::isolate::entry(&task.name).await { self.fail(task.id, format!("Micro plugin failed:\n{e}"))?; return Err(e.into()); } - self.prog.send(TaskOp::Adv(task.id, 1, 0))?; + self.prog.send(TaskProg::Adv(task.id, 1, 0))?; self.succ(task.id) } pub fn macro_(&self, task: PluginOpEntry) -> Result<()> { let id = task.id; - self.prog.send(TaskOp::New(id, 0))?; + self.prog.send(TaskProg::New(id, 0))?; self.tx.send_blocking(PluginOp::Entry(task))?; self.succ(id) } @@ -66,10 +66,10 @@ impl Plugin { impl Plugin { #[inline] - fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) } + fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskProg::Succ(id))?) } #[inline] fn fail(&self, id: usize, reason: String) -> Result<()> { - Ok(self.prog.send(TaskOp::Fail(id, reason))?) + Ok(self.prog.send(TaskProg::Fail(id, reason))?) } } diff --git a/yazi-scheduler/src/workers/preload.rs b/yazi-scheduler/src/workers/preload.rs index bc419a34..5a22bf5e 100644 --- a/yazi-scheduler/src/workers/preload.rs +++ b/yazi-scheduler/src/workers/preload.rs @@ -6,10 +6,10 @@ use tokio::sync::mpsc; use tracing::error; use yazi_shared::{fs::{calculate_size, FilesOp, Url}, Throttle}; -use crate::TaskOp; +use crate::TaskProg; pub struct Preload { - prog: mpsc::UnboundedSender, + prog: mpsc::UnboundedSender, pub rule_loaded: RwLock>, pub size_loading: RwLock>, @@ -32,12 +32,12 @@ pub struct PreloadOpRule { } impl Preload { - pub fn new(prog: mpsc::UnboundedSender) -> Self { + pub fn new(prog: mpsc::UnboundedSender) -> Self { Self { prog, rule_loaded: Default::default(), size_loading: Default::default() } } pub async fn rule(&self, task: PreloadOpRule) -> Result<()> { - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; let urls: Vec<_> = task.targets.iter().map(|f| f.url()).collect(); let result = yazi_plugin::isolate::preload(task.plugin, task.targets, task.rule_multi).await; @@ -57,12 +57,12 @@ impl Preload { } } - self.prog.send(TaskOp::Adv(task.id, 1, 0))?; + self.prog.send(TaskProg::Adv(task.id, 1, 0))?; self.succ(task.id) } pub async fn size(&self, task: PreloadOpSize) -> Result<()> { - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; let length = calculate_size(&task.target).await; task.throttle.done((task.target, length), |buf| { @@ -77,17 +77,17 @@ impl Preload { FilesOp::Size(parent, BTreeMap::from_iter(buf)).emit(); }); - self.prog.send(TaskOp::Adv(task.id, 1, 0))?; + self.prog.send(TaskProg::Adv(task.id, 1, 0))?; self.succ(task.id) } } impl Preload { #[inline] - fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) } + fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskProg::Succ(id))?) } #[inline] fn fail(&self, id: usize, reason: String) -> Result<()> { - Ok(self.prog.send(TaskOp::Fail(id, reason))?) + Ok(self.prog.send(TaskProg::Fail(id, reason))?) } } diff --git a/yazi-scheduler/src/workers/process.rs b/yazi-scheduler/src/workers/process.rs index 67ebd90b..acba196b 100644 --- a/yazi-scheduler/src/workers/process.rs +++ b/yazi-scheduler/src/workers/process.rs @@ -4,10 +4,10 @@ use anyhow::Result; use tokio::{io::{AsyncBufReadExt, BufReader}, select, sync::{mpsc, oneshot}}; use yazi_plugin::external::{self, ShellOpt}; -use crate::{Scheduler, TaskOp, BLOCKER}; +use crate::{Scheduler, TaskProg, BLOCKER}; pub struct Process { - prog: mpsc::UnboundedSender, + prog: mpsc::UnboundedSender, } #[derive(Debug)] @@ -32,7 +32,7 @@ impl From<&mut ProcessOpOpen> for ShellOpt { } impl Process { - pub fn new(prog: mpsc::UnboundedSender) -> Self { Self { prog } } + pub fn new(prog: mpsc::UnboundedSender) -> Self { Self { prog } } pub async fn open(&self, mut task: ProcessOpOpen) -> Result<()> { let opt = ShellOpt::from(&mut task); @@ -46,7 +46,7 @@ impl Process { self.succ(task.id)?; } Err(e) => { - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; self.fail(task.id, format!("Failed to spawn process: {e}"))?; } } @@ -57,14 +57,14 @@ impl Process { match external::shell(opt) { Ok(_) => self.succ(task.id)?, Err(e) => { - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; self.fail(task.id, format!("Failed to spawn process: {e}"))?; } } return Ok(()); } - self.prog.send(TaskOp::New(task.id, 0))?; + self.prog.send(TaskProg::New(task.id, 0))?; let mut child = external::shell(opt.with_piped())?; let mut stdout = BufReader::new(child.stdout.take().unwrap()).lines(); @@ -94,22 +94,22 @@ impl Process { } } - self.prog.send(TaskOp::Adv(task.id, 1, 0))?; + self.prog.send(TaskProg::Adv(task.id, 1, 0))?; self.succ(task.id) } } impl Process { #[inline] - fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) } + fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskProg::Succ(id))?) } #[inline] fn fail(&self, id: usize, reason: String) -> Result<()> { - Ok(self.prog.send(TaskOp::Fail(id, reason))?) + Ok(self.prog.send(TaskProg::Fail(id, reason))?) } #[inline] fn log(&self, id: usize, line: String) -> Result<()> { - Ok(self.prog.send(TaskOp::Log(id, line))?) + Ok(self.prog.send(TaskProg::Log(id, line))?) } }