This commit is contained in:
sxyazi 2023-12-29 01:17:11 +08:00
parent 48a87c40cb
commit a0d251ec8b
No known key found for this signature in database
4 changed files with 46 additions and 13 deletions

View file

@ -7,7 +7,7 @@ use yazi_config::{open::Opener, plugin::PluginRule, TASKS};
use yazi_shared::{emit, event::Exec, fs::{unique_path, Url}, Layer, Throttle}; use yazi_shared::{emit, event::Exec, fs::{unique_path, Url}, Layer, Throttle};
use super::{Running, TaskProg, TaskStage}; use super::{Running, TaskProg, TaskStage};
use crate::{workers::{File, FileOpDelete, FileOpLink, FileOpPaste, FileOpTrash, Plugin, PluginOpEntry, Preload, PreloadOpRule, PreloadOpSize, Process, ProcessOpOpen}, TaskKind}; use crate::{workers::{File, FileOpDelete, FileOpLink, FileOpPaste, FileOpTrash, Plugin, PluginOpEntry, Preload, PreloadOpRule, PreloadOpSize, Process, ProcessOpOpen}, TaskKind, TaskOp};
pub struct Scheduler { pub struct Scheduler {
pub file: Arc<File>, pub file: Arc<File>,
@ -16,6 +16,7 @@ pub struct Scheduler {
pub process: Arc<Process>, pub process: Arc<Process>,
micro: async_channel::Sender<BoxFuture<'static, ()>>, micro: async_channel::Sender<BoxFuture<'static, ()>>,
macro_: async_channel::Sender<TaskOp>,
prog: mpsc::UnboundedSender<TaskProg>, prog: mpsc::UnboundedSender<TaskProg>,
pub running: Arc<RwLock<Running>>, pub running: Arc<RwLock<Running>>,
} }
@ -23,6 +24,7 @@ pub struct Scheduler {
impl Scheduler { impl Scheduler {
pub fn start() -> Self { pub fn start() -> Self {
let (micro_tx, micro_rx) = async_channel::unbounded(); let (micro_tx, micro_rx) = async_channel::unbounded();
let (macro_tx, macro_rx) = async_channel::unbounded();
let (prog_tx, prog_rx) = mpsc::unbounded_channel(); let (prog_tx, prog_rx) = mpsc::unbounded_channel();
let scheduler = Self { let scheduler = Self {
@ -32,6 +34,7 @@ impl Scheduler {
process: Arc::new(Process::new(prog_tx.clone())), process: Arc::new(Process::new(prog_tx.clone())),
micro: micro_tx, micro: micro_tx,
macro_: macro_tx,
prog: prog_tx, prog: prog_tx,
running: Default::default(), running: Default::default(),
}; };
@ -40,7 +43,7 @@ impl Scheduler {
scheduler.schedule_micro(micro_rx.clone()); scheduler.schedule_micro(micro_rx.clone());
} }
for _ in 0..TASKS.macro_workers { for _ in 0..TASKS.macro_workers {
scheduler.schedule_macro(micro_rx.clone()); scheduler.schedule_macro(micro_rx.clone(), macro_rx.clone());
} }
scheduler.progress(prog_rx); scheduler.progress(prog_rx);
scheduler scheduler
@ -56,7 +59,11 @@ impl Scheduler {
}); });
} }
fn schedule_macro(&self, rx: async_channel::Receiver<BoxFuture<'static, ()>>) { fn schedule_macro(
&self,
micro: async_channel::Receiver<BoxFuture<'static, ()>>,
macro_: async_channel::Receiver<TaskOp>,
) {
let file = self.file.clone(); let file = self.file.clone();
let plugin = self.plugin.clone(); let plugin = self.plugin.clone();
@ -65,15 +72,11 @@ impl Scheduler {
tokio::spawn(async move { tokio::spawn(async move {
loop { loop {
if let Ok(fut) = rx.try_recv() {
fut.await;
continue;
}
select! { select! {
Ok(fut) = rx.recv() => { Ok(fut) = micro.recv() => {
fut.await; fut.await;
} }
Ok(op) = macro_.recv() => {}
Ok((id, mut op)) = file.recv() => { Ok((id, mut op)) = file.recv() => {
if !running.read().exists(id) { if !running.read().exists(id) {
continue; continue;

View file

@ -58,6 +58,25 @@ impl From<&Task> for TaskSummary {
} }
} }
#[derive(Debug)]
pub enum TaskOp {
File(Box<crate::workers::FileOp>),
Plugin(Box<crate::workers::PluginOp>),
Preload(Box<crate::workers::PreloadOp>),
Process(Box<crate::workers::ProcessOp>),
}
impl TaskOp {
pub fn id(&self) {
match self {
TaskOp::File(op) => todo!(),
TaskOp::Plugin(op) => todo!(),
TaskOp::Preload(_) => todo!(),
TaskOp::Process(_) => todo!(),
}
}
}
#[derive(Debug)] #[derive(Debug)]
pub enum TaskProg { pub enum TaskProg {
// id, size // id, size

View file

@ -16,10 +16,9 @@ pub struct Preload {
} }
#[derive(Debug)] #[derive(Debug)]
pub struct PreloadOpSize { pub enum PreloadOp {
pub id: usize, Rule(PreloadOpRule),
pub target: Url, Size(PreloadOpSize),
pub throttle: Arc<Throttle<(Url, u64)>>,
} }
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
@ -31,6 +30,13 @@ pub struct PreloadOpRule {
pub targets: Vec<yazi_shared::fs::File>, pub targets: Vec<yazi_shared::fs::File>,
} }
#[derive(Debug)]
pub struct PreloadOpSize {
pub id: usize,
pub target: Url,
pub throttle: Arc<Throttle<(Url, u64)>>,
}
impl Preload { impl Preload {
pub fn new(prog: mpsc::UnboundedSender<TaskProg>) -> Self { pub fn new(prog: mpsc::UnboundedSender<TaskProg>) -> Self {
Self { prog, rule_loaded: Default::default(), size_loading: Default::default() } Self { prog, rule_loaded: Default::default(), size_loading: Default::default() }

View file

@ -10,6 +10,11 @@ pub struct Process {
prog: mpsc::UnboundedSender<TaskProg>, prog: mpsc::UnboundedSender<TaskProg>,
} }
#[derive(Debug)]
pub enum ProcessOp {
Open(ProcessOpOpen),
}
#[derive(Debug)] #[derive(Debug)]
pub struct ProcessOpOpen { pub struct ProcessOpOpen {
pub id: usize, pub id: usize,