feat: fine-grained scheduling priority

This commit is contained in:
sxyazi 2023-12-29 00:45:16 +08:00
parent bfa5e0af0c
commit d00e2ed746
No known key found for this signature in database
5 changed files with 53 additions and 49 deletions

View file

@ -33,10 +33,10 @@ keymap = [
{ on = [ "H" ], exec = "back", desc = "Go back to the previous directory" },
{ on = [ "L" ], exec = "forward", desc = "Go forward to the next directory" },
{ on = [ "<A-k>" ], exec = "seek -5", desc = "Peek up 5 units in the preview" },
{ on = [ "<A-j>" ], exec = "seek 5", desc = "Peek down 5 units in the preview" },
{ on = [ "<A-PageUp>" ], exec = "seek -5", desc = "Peek up 5 units in the preview" },
{ on = [ "<A-PageDown>" ], exec = "seek 5", desc = "Peek down 5 units in the preview" },
{ on = [ "<A-k>" ], exec = "seek -5", desc = "Seek up 5 units in the preview" },
{ on = [ "<A-j>" ], exec = "seek 5", desc = "Seek down 5 units in the preview" },
{ on = [ "<A-PageUp>" ], exec = "seek -5", desc = "Seek up 5 units in the preview" },
{ on = [ "<A-PageDown>" ], exec = "seek 5", desc = "Seek down 5 units in the preview" },
{ on = [ "<Up>" ], exec = "arrow -1", desc = "Move cursor up" },
{ on = [ "<Down>" ], exec = "arrow 1", desc = "Move cursor down" },

View file

@ -13,7 +13,7 @@ pub struct File {
tx: async_channel::Sender<FileOp>,
rx: async_channel::Receiver<FileOp>,
sch: mpsc::UnboundedSender<TaskOp>,
prog: mpsc::UnboundedSender<TaskOp>,
}
#[derive(Debug)]
@ -60,9 +60,9 @@ pub struct FileOpTrash {
}
impl File {
pub fn new(sch: mpsc::UnboundedSender<TaskOp>) -> Self {
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self {
let (tx, rx) = async_channel::unbounded();
Self { tx, rx, sch }
Self { tx, rx, prog }
}
#[inline]
@ -92,7 +92,7 @@ impl File {
}
break;
}
Ok(n) => self.sch.send(TaskOp::Adv(task.id, 0, n))?,
Ok(n) => self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, meta.len()))?);
return Ok(self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, meta.len()))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, task.length))?
self.prog.send(TaskOp::Adv(task.id, 1, task.length))?
}
FileOp::Trash(task) => {
#[cfg(target_os = "macos")]
@ -178,7 +178,7 @@ impl File {
{
trash::delete(&task.target)?;
}
self.sch.send(TaskOp::Adv(task.id, 1, task.length))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(id, meta.len()))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(task.id, meta.len()))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(id, task.meta.as_ref().unwrap().len()))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(id, meta.len()))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(task.id, meta.len()))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(id, task.length))?;
self.prog.send(TaskOp::New(id, task.length))?;
self.tx.send(FileOp::Trash(task)).await?;
self.succ(id)
}
@ -343,15 +343,17 @@ impl File {
impl File {
#[inline]
fn succ(&self, id: usize) -> Result<()> { Ok(self.sch.send(TaskOp::Succ(id))?) }
fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) }
#[inline]
fn fail(&self, id: usize, reason: String) -> Result<()> {
Ok(self.sch.send(TaskOp::Fail(id, reason))?)
Ok(self.prog.send(TaskOp::Fail(id, reason))?)
}
#[inline]
fn log(&self, id: usize, line: String) -> Result<()> { Ok(self.sch.send(TaskOp::Log(id, line))?) }
fn log(&self, id: usize, line: String) -> Result<()> {
Ok(self.prog.send(TaskOp::Log(id, line))?)
}
}
impl FileOpPaste {

View file

@ -7,7 +7,7 @@ pub struct Plugin {
tx: async_channel::Sender<PluginOp>,
rx: async_channel::Receiver<PluginOp>,
sch: mpsc::UnboundedSender<TaskOp>,
prog: mpsc::UnboundedSender<TaskOp>,
}
#[derive(Debug)]
@ -22,9 +22,9 @@ pub struct PluginOpEntry {
}
impl Plugin {
pub fn new(sch: mpsc::UnboundedSender<TaskOp>) -> Self {
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self {
let (tx, rx) = async_channel::unbounded();
Self { tx, rx, sch }
Self { tx, rx, prog }
}
#[inline]
@ -44,21 +44,21 @@ impl Plugin {
}
pub async fn micro(&self, task: PluginOpEntry) -> Result<()> {
self.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, 0))?;
self.prog.send(TaskOp::Adv(task.id, 1, 0))?;
self.succ(task.id)
}
pub fn macro_(&self, task: PluginOpEntry) -> Result<()> {
let id = task.id;
self.sch.send(TaskOp::New(id, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::Succ(id))?) }
fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) }
#[inline]
fn fail(&self, id: usize, reason: String) -> Result<()> {
Ok(self.sch.send(TaskOp::Fail(id, reason))?)
Ok(self.prog.send(TaskOp::Fail(id, reason))?)
}
}

View file

@ -9,7 +9,7 @@ use yazi_shared::{fs::{calculate_size, FilesOp, Url}, Throttle};
use crate::TaskOp;
pub struct Preload {
sch: mpsc::UnboundedSender<TaskOp>,
prog: mpsc::UnboundedSender<TaskOp>,
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(sch: mpsc::UnboundedSender<TaskOp>) -> Self {
Self { sch, rule_loaded: Default::default(), size_loading: Default::default() }
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> Self {
Self { prog, rule_loaded: Default::default(), size_loading: Default::default() }
}
pub async fn rule(&self, task: PreloadOpRule) -> Result<()> {
self.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, 0))?;
self.prog.send(TaskOp::Adv(task.id, 1, 0))?;
self.succ(task.id)
}
pub async fn size(&self, task: PreloadOpSize) -> Result<()> {
self.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::Adv(task.id, 1, 0))?;
self.prog.send(TaskOp::Adv(task.id, 1, 0))?;
self.succ(task.id)
}
}
impl Preload {
#[inline]
fn succ(&self, id: usize) -> Result<()> { Ok(self.sch.send(TaskOp::Succ(id))?) }
fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) }
#[inline]
fn fail(&self, id: usize, reason: String) -> Result<()> {
Ok(self.sch.send(TaskOp::Fail(id, reason))?)
Ok(self.prog.send(TaskOp::Fail(id, reason))?)
}
}

View file

@ -7,7 +7,7 @@ use yazi_plugin::external::{self, ShellOpt};
use crate::{Scheduler, TaskOp, BLOCKER};
pub struct Process {
sch: mpsc::UnboundedSender<TaskOp>,
prog: mpsc::UnboundedSender<TaskOp>,
}
#[derive(Debug)]
@ -32,7 +32,7 @@ impl From<&mut ProcessOpOpen> for ShellOpt {
}
impl Process {
pub fn new(sch: mpsc::UnboundedSender<TaskOp>) -> Self { Self { sch } }
pub fn new(prog: mpsc::UnboundedSender<TaskOp>) -> 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.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::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.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::New(task.id, 0))?;
self.fail(task.id, format!("Failed to spawn process: {e}"))?;
}
}
return Ok(());
}
self.sch.send(TaskOp::New(task.id, 0))?;
self.prog.send(TaskOp::New(task.id, 0))?;
let mut child = external::shell(opt.with_piped())?;
let mut stdout = BufReader::new(child.stdout.take().unwrap()).lines();
@ -94,20 +94,22 @@ impl Process {
}
}
self.sch.send(TaskOp::Adv(task.id, 1, 0))?;
self.prog.send(TaskOp::Adv(task.id, 1, 0))?;
self.succ(task.id)
}
}
impl Process {
#[inline]
fn succ(&self, id: usize) -> Result<()> { Ok(self.sch.send(TaskOp::Succ(id))?) }
fn succ(&self, id: usize) -> Result<()> { Ok(self.prog.send(TaskOp::Succ(id))?) }
#[inline]
fn fail(&self, id: usize, reason: String) -> Result<()> {
Ok(self.sch.send(TaskOp::Fail(id, reason))?)
Ok(self.prog.send(TaskOp::Fail(id, reason))?)
}
#[inline]
fn log(&self, id: usize, line: String) -> Result<()> { Ok(self.sch.send(TaskOp::Log(id, line))?) }
fn log(&self, id: usize, line: String) -> Result<()> {
Ok(self.prog.send(TaskOp::Log(id, line))?)
}
}