From 09ee8dc2b757dc8bea533de293d5518e2138577b Mon Sep 17 00:00:00 2001 From: sxyazi Date: Fri, 29 Dec 2023 02:16:30 +0800 Subject: [PATCH] .. --- yazi-scheduler/src/plugin/plugin.rs | 22 ++++++---------------- yazi-scheduler/src/preload/preload.rs | 23 ++++++++++++++++++----- yazi-scheduler/src/scheduler.rs | 11 +++++------ 3 files changed, 29 insertions(+), 27 deletions(-) diff --git a/yazi-scheduler/src/plugin/plugin.rs b/yazi-scheduler/src/plugin/plugin.rs index fc5e168b..e4d02d39 100644 --- a/yazi-scheduler/src/plugin/plugin.rs +++ b/yazi-scheduler/src/plugin/plugin.rs @@ -2,26 +2,16 @@ use anyhow::Result; use tokio::sync::mpsc; use super::{PluginOp, PluginOpEntry}; -use crate::TaskProg; +use crate::{TaskOp, TaskProg}; pub struct Plugin { - tx: async_channel::Sender, - rx: async_channel::Receiver, - - prog: mpsc::UnboundedSender, + macro_: async_channel::Sender, + prog: mpsc::UnboundedSender, } impl Plugin { - pub fn new(prog: mpsc::UnboundedSender) -> Self { - let (tx, rx) = async_channel::unbounded(); - Self { tx, rx, prog } - } - - #[inline] - pub async fn recv(&self) -> Result<(usize, PluginOp)> { - Ok(match self.rx.recv().await? { - PluginOp::Entry(t) => (t.id, PluginOp::Entry(t)), - }) + pub fn new(macro_: async_channel::Sender, prog: mpsc::UnboundedSender) -> Self { + Self { macro_, prog } } pub async fn work(&self, op: &mut PluginOp) -> Result<()> { @@ -49,7 +39,7 @@ impl Plugin { let id = task.id; self.prog.send(TaskProg::New(id, 0))?; - self.tx.send_blocking(PluginOp::Entry(task))?; + self.macro_.send_blocking(PluginOp::Entry(task).into())?; self.succ(id) } } diff --git a/yazi-scheduler/src/preload/preload.rs b/yazi-scheduler/src/preload/preload.rs index dfbd5090..83faedc4 100644 --- a/yazi-scheduler/src/preload/preload.rs +++ b/yazi-scheduler/src/preload/preload.rs @@ -6,19 +6,32 @@ use tokio::sync::mpsc; use tracing::error; use yazi_shared::fs::{calculate_size, FilesOp, Url}; -use super::{PreloadOpRule, PreloadOpSize}; -use crate::TaskProg; +use super::{PreloadOp, PreloadOpRule, PreloadOpSize}; +use crate::{TaskOp, TaskProg}; pub struct Preload { - prog: mpsc::UnboundedSender, + macro_: async_channel::Sender, + prog: mpsc::UnboundedSender, pub rule_loaded: RwLock>, pub size_loading: RwLock>, } impl Preload { - pub fn new(prog: mpsc::UnboundedSender) -> Self { - Self { prog, rule_loaded: Default::default(), size_loading: Default::default() } + pub fn new(macro_: async_channel::Sender, prog: mpsc::UnboundedSender) -> Self { + Self { macro_, prog, rule_loaded: Default::default(), size_loading: Default::default() } + } + + pub async fn work(&self, op: &mut PreloadOp) -> Result<()> { + match op { + PreloadOp::Rule(task) => { + todo!() + } + PreloadOp::Size(task) => { + todo!() + } + } + Ok(()) } pub async fn rule(&self, task: PreloadOpRule) -> Result<()> { diff --git a/yazi-scheduler/src/scheduler.rs b/yazi-scheduler/src/scheduler.rs index 54473dc0..5557c33f 100644 --- a/yazi-scheduler/src/scheduler.rs +++ b/yazi-scheduler/src/scheduler.rs @@ -16,7 +16,6 @@ pub struct Scheduler { pub process: Arc, micro: async_channel::Sender>, - macro_: async_channel::Sender, prog: mpsc::UnboundedSender, pub running: Arc>, } @@ -29,12 +28,11 @@ impl Scheduler { let scheduler = Self { file: Arc::new(File::new(macro_tx.clone(), prog_tx.clone())), - plugin: Arc::new(Plugin::new(prog_tx.clone())), - preload: Arc::new(Preload::new(prog_tx.clone())), + plugin: Arc::new(Plugin::new(macro_tx.clone(), prog_tx.clone())), + preload: Arc::new(Preload::new(macro_tx.clone(), prog_tx.clone())), process: Arc::new(Process::new(prog_tx.clone())), micro: micro_tx, - macro_: macro_tx, prog: prog_tx, running: Default::default(), }; @@ -66,6 +64,7 @@ impl Scheduler { ) { let file = self.file.clone(); let plugin = self.plugin.clone(); + let preload = self.preload.clone(); let prog = self.prog.clone(); let running = self.running.clone(); @@ -84,8 +83,8 @@ impl Scheduler { let result = match op { TaskOp::File(mut op) => file.work(&mut op).await, - TaskOp::Plugin(mut op) => todo!(), - TaskOp::Preload(mut op) => todo!(), + TaskOp::Plugin(mut op) => plugin.work(&mut op).await, + TaskOp::Preload(mut op) => preload.work(&mut op).await, }; if let Err(e) = result {