This commit is contained in:
sxyazi 2023-12-29 00:55:35 +08:00
parent d00e2ed746
commit 48a87c40cb
No known key found for this signature in database
6 changed files with 74 additions and 74 deletions

View file

@ -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<Preload>,
pub process: Arc<Process>,
todo: async_channel::Sender<BoxFuture<'static, ()>>,
prog: mpsc::UnboundedSender<TaskOp>,
micro: async_channel::Sender<BoxFuture<'static, ()>>,
prog: mpsc::UnboundedSender<TaskProg>,
pub running: Arc<RwLock<Running>>,
}
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<TaskOp>) {
let todo = self.todo.clone();
fn progress(&self, mut rx: UnboundedReceiver<TaskProg>) {
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();

View file

@ -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

View file

@ -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<FileOp>,
rx: async_channel::Receiver<FileOp>,
prog: mpsc::UnboundedSender<TaskOp>,
prog: mpsc::UnboundedSender<TaskProg>,
}
#[derive(Debug)]
@ -60,7 +60,7 @@ pub struct FileOpTrash {
}
impl File {
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self {
pub fn new(prog: mpsc::UnboundedSender<TaskProg>) -> 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))?)
}
}

View file

@ -1,13 +1,13 @@
use anyhow::Result;
use tokio::sync::mpsc;
use crate::TaskOp;
use crate::TaskProg;
pub struct Plugin {
tx: async_channel::Sender<PluginOp>,
rx: async_channel::Receiver<PluginOp>,
prog: mpsc::UnboundedSender<TaskOp>,
prog: mpsc::UnboundedSender<TaskProg>,
}
#[derive(Debug)]
@ -22,7 +22,7 @@ pub struct PluginOpEntry {
}
impl Plugin {
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self {
pub fn new(prog: mpsc::UnboundedSender<TaskProg>) -> 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))?)
}
}

View file

@ -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<TaskOp>,
prog: mpsc::UnboundedSender<TaskProg>,
pub rule_loaded: RwLock<HashMap<Url, u32>>,
pub size_loading: RwLock<BTreeSet<Url>>,
@ -32,12 +32,12 @@ pub struct PreloadOpRule {
}
impl Preload {
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self {
pub fn new(prog: mpsc::UnboundedSender<TaskProg>) -> 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))?)
}
}

View file

@ -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<TaskOp>,
prog: mpsc::UnboundedSender<TaskProg>,
}
#[derive(Debug)]
@ -32,7 +32,7 @@ impl From<&mut ProcessOpOpen> for ShellOpt {
}
impl Process {
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self { Self { prog } }
pub fn new(prog: mpsc::UnboundedSender<TaskProg>) -> 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))?)
}
}