mirror of
https://github.com/sxyazi/yazi.git
synced 2026-07-25 08:41:05 +00:00
perf: immediate task cancellation (#3429)
This commit is contained in:
parent
5fe58ba2b1
commit
80b3b8465d
27 changed files with 287 additions and 266 deletions
|
|
@ -84,6 +84,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/):
|
||||||
|
|
||||||
### Improved
|
### Improved
|
||||||
|
|
||||||
|
- Make copy, cut, delete, link, hardlink, download, and upload tasks immediately cancellable ([#3429])
|
||||||
- Make preload tasks discardable ([#2875])
|
- Make preload tasks discardable ([#2875])
|
||||||
- Reduce file change event frequency ([#2820])
|
- Reduce file change event frequency ([#2820])
|
||||||
- Upload and download of a single file over SFTP in chunks concurrently ([#3393])
|
- Upload and download of a single file over SFTP in chunks concurrently ([#3393])
|
||||||
|
|
@ -1556,3 +1557,4 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/):
|
||||||
[#3396]: https://github.com/sxyazi/yazi/pull/3396
|
[#3396]: https://github.com/sxyazi/yazi/pull/3396
|
||||||
[#3419]: https://github.com/sxyazi/yazi/pull/3419
|
[#3419]: https://github.com/sxyazi/yazi/pull/3419
|
||||||
[#3422]: https://github.com/sxyazi/yazi/pull/3422
|
[#3422]: https://github.com/sxyazi/yazi/pull/3422
|
||||||
|
[#3429]: https://github.com/sxyazi/yazi/pull/3429
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,9 @@ resolver = "2"
|
||||||
members = [ "yazi-*" ]
|
members = [ "yazi-*" ]
|
||||||
default-members = [ "yazi-fm", "yazi-cli" ]
|
default-members = [ "yazi-fm", "yazi-cli" ]
|
||||||
|
|
||||||
|
[profile.dev]
|
||||||
|
debug = "line-tables-only"
|
||||||
|
|
||||||
[profile.release]
|
[profile.release]
|
||||||
codegen-units = 1
|
codegen-units = 1
|
||||||
lto = true
|
lto = true
|
||||||
|
|
@ -19,6 +22,9 @@ codegen-units = 256
|
||||||
incremental = true
|
incremental = true
|
||||||
lto = false
|
lto = false
|
||||||
|
|
||||||
|
[profile.dev.package."*"]
|
||||||
|
debug = false
|
||||||
|
|
||||||
[workspace.dependencies]
|
[workspace.dependencies]
|
||||||
ansi-to-tui = "7.0.0"
|
ansi-to-tui = "7.0.0"
|
||||||
anyhow = "1.0.100"
|
anyhow = "1.0.100"
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,6 @@ use std::{mem, time::{Duration, Instant}};
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use futures::{StreamExt, stream::FuturesUnordered};
|
use futures::{StreamExt, stream::FuturesUnordered};
|
||||||
use hashbrown::HashSet;
|
use hashbrown::HashSet;
|
||||||
use tokio::sync::oneshot;
|
|
||||||
use yazi_fs::{File, FsScheme, provider::{Provider, local::Local}};
|
use yazi_fs::{File, FsScheme, provider::{Provider, local::Local}};
|
||||||
use yazi_macro::succ;
|
use yazi_macro::succ;
|
||||||
use yazi_parser::mgr::{DownloadOpt, OpenOpt};
|
use yazi_parser::mgr::{DownloadOpt, OpenOpt};
|
||||||
|
|
@ -29,9 +28,8 @@ impl Actor for Download {
|
||||||
|
|
||||||
let mut wg1 = FuturesUnordered::new();
|
let mut wg1 = FuturesUnordered::new();
|
||||||
for url in opt.urls {
|
for url in opt.urls {
|
||||||
let (tx, rx) = oneshot::channel();
|
let done = scheduler.file_download(url.to_owned());
|
||||||
scheduler.file_download(url.to_owned(), Some(tx));
|
wg1.push(async move { (done.future().await, url) });
|
||||||
wg1.push(async move { (rx.await == Ok(true), url) });
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut wg2 = vec![];
|
let mut wg2 = vec![];
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,7 @@ impl Tasks {
|
||||||
drop(loaded);
|
drop(loaded);
|
||||||
for (i, tasks) in tasks.into_iter().enumerate() {
|
for (i, tasks) in tasks.into_iter().enumerate() {
|
||||||
if !tasks.is_empty() {
|
if !tasks.is_empty() {
|
||||||
self.scheduler.fetch_paged(&YAZI.plugin.fetchers[i], tasks, None);
|
self.scheduler.fetch_paged(&YAZI.plugin.fetchers[i], tasks);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,8 +2,7 @@ use std::ffi::OsString;
|
||||||
|
|
||||||
use anyhow::anyhow;
|
use anyhow::anyhow;
|
||||||
use mlua::{ExternalError, FromLua, IntoLua, Lua, Value};
|
use mlua::{ExternalError, FromLua, IntoLua, Lua, Value};
|
||||||
use tokio::sync::oneshot;
|
use yazi_shared::{CompletionToken, event::CmdCow, url::UrlCow};
|
||||||
use yazi_shared::{event::CmdCow, url::UrlCow};
|
|
||||||
|
|
||||||
// --- Exec
|
// --- Exec
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
|
|
@ -13,7 +12,7 @@ pub struct ProcessOpenOpt {
|
||||||
pub args: Vec<UrlCow<'static>>,
|
pub args: Vec<UrlCow<'static>>,
|
||||||
pub block: bool,
|
pub block: bool,
|
||||||
pub orphan: bool,
|
pub orphan: bool,
|
||||||
pub done: Option<oneshot::Sender<()>>,
|
pub done: Option<CompletionToken>,
|
||||||
|
|
||||||
pub spread: bool, // TODO: remove
|
pub spread: bool, // TODO: remove
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,8 @@
|
||||||
use std::ffi::OsString;
|
use std::ffi::OsString;
|
||||||
|
|
||||||
use tokio::sync::oneshot;
|
|
||||||
use yazi_macro::{emit, relay};
|
use yazi_macro::{emit, relay};
|
||||||
use yazi_parser::tasks::ProcessOpenOpt;
|
use yazi_parser::tasks::ProcessOpenOpt;
|
||||||
use yazi_shared::url::{UrlBuf, UrlCow};
|
use yazi_shared::{CompletionToken, url::{UrlBuf, UrlCow}};
|
||||||
|
|
||||||
pub struct TasksProxy;
|
pub struct TasksProxy;
|
||||||
|
|
||||||
|
|
@ -20,17 +19,17 @@ impl TasksProxy {
|
||||||
block: bool,
|
block: bool,
|
||||||
orphan: bool,
|
orphan: bool,
|
||||||
) {
|
) {
|
||||||
let (tx, rx) = oneshot::channel();
|
let done = CompletionToken::new();
|
||||||
emit!(Call(relay!(tasks:process_open).with_any("opt", ProcessOpenOpt {
|
emit!(Call(relay!(tasks:process_open).with_any("opt", ProcessOpenOpt {
|
||||||
cwd,
|
cwd,
|
||||||
cmd,
|
cmd,
|
||||||
args,
|
args,
|
||||||
block,
|
block,
|
||||||
orphan,
|
orphan,
|
||||||
done: Some(tx),
|
done: Some(done.clone()),
|
||||||
spread: false
|
spread: false
|
||||||
})));
|
})));
|
||||||
rx.await.ok();
|
done.future().await;
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn update_succeed(url: impl Into<UrlBuf>) {
|
pub fn update_succeed(url: impl Into<UrlBuf>) {
|
||||||
|
|
|
||||||
|
|
@ -9,7 +9,7 @@ use yazi_shared::{path::PathCow, timestamp_us, url::{AsUrl, UrlBuf, UrlCow, UrlL
|
||||||
use yazi_vfs::{VfsCha, maybe_exists, must_be_dir, provider::{self, DirEntry}, unique_name};
|
use yazi_vfs::{VfsCha, maybe_exists, must_be_dir, provider::{self, DirEntry}, unique_name};
|
||||||
|
|
||||||
use super::{FileInCopy, FileInDelete, FileInHardlink, FileInLink, FileInTrash};
|
use super::{FileInCopy, FileInDelete, FileInHardlink, FileInLink, FileInTrash};
|
||||||
use crate::{LOW, NORMAL, TaskIn, TaskOp, TaskOps, ctx, file::{FileInCut, FileInDownload, FileInUpload, FileOutCopy, FileOutCopyDo, FileOutCut, FileOutCutDo, FileOutDelete, FileOutDeleteDo, FileOutDownload, FileOutDownloadDo, FileOutHardlink, FileOutHardlinkDo, FileOutLink, FileOutTrash, FileOutUpload, FileOutUploadDo}, hook::HookInOutCut, ok_or_not_found};
|
use crate::{LOW, NORMAL, TaskIn, TaskOp, TaskOps, ctx, file::{FileInCut, FileInDownload, FileInUpload, FileOutCopy, FileOutCopyDo, FileOutCut, FileOutCutDo, FileOutDelete, FileOutDeleteDo, FileOutDownload, FileOutDownloadDo, FileOutHardlink, FileOutHardlinkDo, FileOutLink, FileOutTrash, FileOutUpload, FileOutUploadDo}, hook::HookInOutCut, ok_or_not_found, progress_or_break};
|
||||||
|
|
||||||
pub(crate) struct File {
|
pub(crate) struct File {
|
||||||
ops: TaskOps,
|
ops: TaskOps,
|
||||||
|
|
@ -62,8 +62,8 @@ impl File {
|
||||||
let mut it =
|
let mut it =
|
||||||
ctx!(task, provider::copy_with_progress(&task.from, &task.to, task.cha.unwrap()).await)?;
|
ctx!(task, provider::copy_with_progress(&task.from, &task.to, task.cha.unwrap()).await)?;
|
||||||
|
|
||||||
while let Some(res) = it.recv().await {
|
loop {
|
||||||
match res {
|
match progress_or_break!(it, task.done) {
|
||||||
Ok(0) => break,
|
Ok(0) => break,
|
||||||
Ok(n) => self.ops.out(task.id, FileOutCopyDo::Adv(n)),
|
Ok(n) => self.ops.out(task.id, FileOutCopyDo::Adv(n)),
|
||||||
Err(e) if e.kind() == NotFound => {
|
Err(e) if e.kind() == NotFound => {
|
||||||
|
|
@ -152,8 +152,8 @@ impl File {
|
||||||
let mut it =
|
let mut it =
|
||||||
ctx!(task, provider::copy_with_progress(&task.from, &task.to, task.cha.unwrap()).await)?;
|
ctx!(task, provider::copy_with_progress(&task.from, &task.to, task.cha.unwrap()).await)?;
|
||||||
|
|
||||||
while let Some(res) = it.recv().await {
|
loop {
|
||||||
match res {
|
match progress_or_break!(it, task.done) {
|
||||||
Ok(0) => {
|
Ok(0) => {
|
||||||
provider::remove_file(&task.from).await.ok();
|
provider::remove_file(&task.from).await.ok();
|
||||||
break;
|
break;
|
||||||
|
|
@ -340,8 +340,8 @@ impl File {
|
||||||
let cache_tmp = ctx!(task, Self::tmp(&cache).await, "Cannot determine download cache")?;
|
let cache_tmp = ctx!(task, Self::tmp(&cache).await, "Cannot determine download cache")?;
|
||||||
|
|
||||||
let mut it = ctx!(task, provider::copy_with_progress(&task.url, &cache_tmp, cha).await)?;
|
let mut it = ctx!(task, provider::copy_with_progress(&task.url, &cache_tmp, cha).await)?;
|
||||||
while let Some(res) = it.recv().await {
|
loop {
|
||||||
match res {
|
match progress_or_break!(it, task.done) {
|
||||||
Ok(0) => {
|
Ok(0) => {
|
||||||
Local::regular(&cache).remove_dir_all().await.ok();
|
Local::regular(&cache).remove_dir_all().await.ok();
|
||||||
ctx!(task, provider::rename(cache_tmp, cache).await, "Cannot persist downloaded file")?;
|
ctx!(task, provider::rename(cache_tmp, cache).await, "Cannot persist downloaded file")?;
|
||||||
|
|
@ -424,8 +424,8 @@ impl File {
|
||||||
.await
|
.await
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
while let Some(res) = it.recv().await {
|
loop {
|
||||||
match res {
|
match progress_or_break!(it, task.done) {
|
||||||
Ok(0) => {
|
Ok(0) => {
|
||||||
let cha =
|
let cha =
|
||||||
ctx!(task, Self::cha(&task.url, true, None).await, "Cannot stat original file")?;
|
ctx!(task, Self::cha(&task.url, true, None).await, "Cannot stat original file")?;
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,7 @@ use std::{mem, path::PathBuf};
|
||||||
|
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
use yazi_fs::cha::Cha;
|
use yazi_fs::cha::Cha;
|
||||||
use yazi_shared::{Id, url::UrlBuf};
|
use yazi_shared::{CompletionToken, Id, url::UrlBuf};
|
||||||
|
|
||||||
// --- Copy
|
// --- Copy
|
||||||
#[derive(Clone, Debug)]
|
#[derive(Clone, Debug)]
|
||||||
|
|
@ -14,6 +14,7 @@ pub(crate) struct FileInCopy {
|
||||||
pub(crate) cha: Option<Cha>,
|
pub(crate) cha: Option<Cha>,
|
||||||
pub(crate) follow: bool,
|
pub(crate) follow: bool,
|
||||||
pub(crate) retry: u8,
|
pub(crate) retry: u8,
|
||||||
|
pub(crate) done: CompletionToken,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl FileInCopy {
|
impl FileInCopy {
|
||||||
|
|
@ -41,6 +42,7 @@ pub(crate) struct FileInCut {
|
||||||
pub(crate) cha: Option<Cha>,
|
pub(crate) cha: Option<Cha>,
|
||||||
pub(crate) follow: bool,
|
pub(crate) follow: bool,
|
||||||
pub(crate) retry: u8,
|
pub(crate) retry: u8,
|
||||||
|
pub(crate) done: CompletionToken,
|
||||||
pub(crate) drop: Option<mpsc::Sender<()>>,
|
pub(crate) drop: Option<mpsc::Sender<()>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -114,6 +116,7 @@ pub(crate) struct FileInDownload {
|
||||||
pub(crate) url: UrlBuf,
|
pub(crate) url: UrlBuf,
|
||||||
pub(crate) cha: Option<Cha>,
|
pub(crate) cha: Option<Cha>,
|
||||||
pub(crate) retry: u8,
|
pub(crate) retry: u8,
|
||||||
|
pub(crate) done: CompletionToken,
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- Upload
|
// --- Upload
|
||||||
|
|
@ -123,4 +126,5 @@ pub(crate) struct FileInUpload {
|
||||||
pub(crate) url: UrlBuf,
|
pub(crate) url: UrlBuf,
|
||||||
pub(crate) cha: Option<Cha>,
|
pub(crate) cha: Option<Cha>,
|
||||||
pub(crate) cache: Option<PathBuf>,
|
pub(crate) cache: Option<PathBuf>,
|
||||||
|
pub(crate) done: CompletionToken,
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -41,6 +41,7 @@ impl Traverse for FileInCopy {
|
||||||
cha: Some(cha),
|
cha: Some(cha),
|
||||||
follow: self.follow,
|
follow: self.follow,
|
||||||
retry: self.retry,
|
retry: self.retry,
|
||||||
|
done: self.done.clone(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -63,6 +64,7 @@ impl Traverse for FileInCut {
|
||||||
cha: Some(cha),
|
cha: Some(cha),
|
||||||
follow: self.follow,
|
follow: self.follow,
|
||||||
retry: self.retry,
|
retry: self.retry,
|
||||||
|
done: self.done.clone(),
|
||||||
drop: self.drop.clone(),
|
drop: self.drop.clone(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -113,7 +115,13 @@ impl Traverse for FileInDownload {
|
||||||
fn from(&self) -> Url<'_> { self.url.as_url() }
|
fn from(&self) -> Url<'_> { self.url.as_url() }
|
||||||
|
|
||||||
fn spawn(&self, from: UrlBuf, _to: Option<UrlBuf>, cha: Cha) -> Self {
|
fn spawn(&self, from: UrlBuf, _to: Option<UrlBuf>, cha: Cha) -> Self {
|
||||||
Self { id: self.id, url: from, cha: Some(cha), retry: self.retry }
|
Self {
|
||||||
|
id: self.id,
|
||||||
|
url: from,
|
||||||
|
cha: Some(cha),
|
||||||
|
retry: self.retry,
|
||||||
|
done: self.done.clone(),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn to(&self) -> Option<Url<'_>> { None }
|
fn to(&self) -> Option<Url<'_>> { None }
|
||||||
|
|
@ -137,7 +145,13 @@ impl Traverse for FileInUpload {
|
||||||
}
|
}
|
||||||
|
|
||||||
fn spawn(&self, from: UrlBuf, _to: Option<UrlBuf>, cha: Cha) -> Self {
|
fn spawn(&self, from: UrlBuf, _to: Option<UrlBuf>, cha: Cha) -> Self {
|
||||||
Self { id: self.id, cha: Some(cha), cache: from.cache(), url: from }
|
Self {
|
||||||
|
id: self.id,
|
||||||
|
cha: Some(cha),
|
||||||
|
cache: from.cache(),
|
||||||
|
url: from,
|
||||||
|
done: self.done.clone(),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn to(&self) -> Option<Url<'_>> { None }
|
fn to(&self) -> Option<Url<'_>> { None }
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@ use yazi_dds::Pump;
|
||||||
use yazi_proxy::TasksProxy;
|
use yazi_proxy::TasksProxy;
|
||||||
use yazi_vfs::provider;
|
use yazi_vfs::provider;
|
||||||
|
|
||||||
use crate::{Ongoing, TaskOp, TaskOps, file::{FileOutCut, FileOutDelete, FileOutDownload, FileOutTrash}, hook::{HookInOutBg, HookInOutBlock, HookInOutCut, HookInOutDelete, HookInOutDownload, HookInOutFetch, HookInOutOrphan, HookInOutTrash}, prework::PreworkOutFetch, process::{ProcessOutBg, ProcessOutBlock, ProcessOutOrphan}};
|
use crate::{Ongoing, TaskOp, TaskOps, file::{FileOutCut, FileOutDelete, FileOutDownload, FileOutTrash}, hook::{HookInOutCut, HookInOutDelete, HookInOutDownload, HookInOutTrash}};
|
||||||
|
|
||||||
pub(crate) struct Hook {
|
pub(crate) struct Hook {
|
||||||
ops: TaskOps,
|
ops: TaskOps,
|
||||||
|
|
@ -48,42 +48,6 @@ impl Hook {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn download(&self, task: HookInOutDownload) {
|
pub(crate) async fn download(&self, task: HookInOutDownload) {
|
||||||
let intact = self.ongoing.lock().intact(task.id);
|
|
||||||
task.done.send(intact).ok();
|
|
||||||
self.ops.out(task.id, FileOutDownload::Clean);
|
self.ops.out(task.id, FileOutDownload::Clean);
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- Process
|
|
||||||
pub(crate) async fn block(&self, task: HookInOutBlock) {
|
|
||||||
if let Some(tx) = task.done {
|
|
||||||
tx.send(()).ok();
|
|
||||||
}
|
|
||||||
self.ops.out(task.id, ProcessOutBlock::Clean);
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) async fn orphan(&self, task: HookInOutOrphan) {
|
|
||||||
if let Some(tx) = task.done {
|
|
||||||
tx.send(()).ok();
|
|
||||||
}
|
|
||||||
self.ops.out(task.id, ProcessOutOrphan::Clean);
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) async fn bg(&self, task: HookInOutBg) {
|
|
||||||
let intact = self.ongoing.lock().intact(task.id);
|
|
||||||
if !intact {
|
|
||||||
task.cancel.send(()).await.ok();
|
|
||||||
task.cancel.closed().await;
|
|
||||||
}
|
|
||||||
if let Some(tx) = task.done {
|
|
||||||
tx.send(()).ok();
|
|
||||||
}
|
|
||||||
self.ops.out(task.id, ProcessOutBg::Clean);
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- Prework
|
|
||||||
pub(crate) async fn fetch(&self, task: HookInOutFetch) {
|
|
||||||
let intact = self.ongoing.lock().intact(task.id);
|
|
||||||
task.done.send(intact).ok();
|
|
||||||
self.ops.out(task.id, PreworkOutFetch::Clean);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,3 @@
|
||||||
use tokio::sync::{mpsc, oneshot};
|
|
||||||
use yazi_shared::{Id, url::UrlBuf};
|
use yazi_shared::{Id, url::UrlBuf};
|
||||||
|
|
||||||
use crate::{Task, TaskProg, file::FileInCut};
|
use crate::{Task, TaskProg, file::FileInCut};
|
||||||
|
|
@ -58,8 +57,7 @@ impl HookInOutTrash {
|
||||||
// --- Download
|
// --- Download
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub(crate) struct HookInOutDownload {
|
pub(crate) struct HookInOutDownload {
|
||||||
pub(crate) id: Id,
|
pub(crate) id: Id,
|
||||||
pub(crate) done: oneshot::Sender<bool>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl HookInOutDownload {
|
impl HookInOutDownload {
|
||||||
|
|
@ -69,64 +67,3 @@ impl HookInOutDownload {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- Fetch
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub(crate) struct HookInOutFetch {
|
|
||||||
pub(crate) id: Id,
|
|
||||||
pub(crate) done: oneshot::Sender<bool>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl HookInOutFetch {
|
|
||||||
pub(crate) fn reduce(self, task: &mut Task) {
|
|
||||||
if let TaskProg::PreworkFetch(_) = &task.prog {
|
|
||||||
task.hook = Some(self.into());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- Block
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub(crate) struct HookInOutBlock {
|
|
||||||
pub(crate) id: Id,
|
|
||||||
pub(crate) done: Option<oneshot::Sender<()>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl HookInOutBlock {
|
|
||||||
pub(crate) fn reduce(self, task: &mut Task) {
|
|
||||||
if let TaskProg::ProcessBlock(_) = &task.prog {
|
|
||||||
task.hook = Some(self.into());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- Orphan
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub(crate) struct HookInOutOrphan {
|
|
||||||
pub(crate) id: Id,
|
|
||||||
pub(crate) done: Option<oneshot::Sender<()>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl HookInOutOrphan {
|
|
||||||
pub(crate) fn reduce(self, task: &mut Task) {
|
|
||||||
if let TaskProg::ProcessOrphan(_) = &task.prog {
|
|
||||||
task.hook = Some(self.into());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- Bg
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub(crate) struct HookInOutBg {
|
|
||||||
pub(crate) id: Id,
|
|
||||||
pub(crate) cancel: mpsc::Sender<()>,
|
|
||||||
pub(crate) done: Option<oneshot::Sender<()>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl HookInOutBg {
|
|
||||||
pub(crate) fn reduce(self, task: &mut Task) {
|
|
||||||
if let TaskProg::ProcessBg(_) = &task.prog {
|
|
||||||
task.hook = Some(self.into());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
use yazi_shared::Id;
|
use yazi_shared::Id;
|
||||||
|
|
||||||
use crate::{file::{FileInCopy, FileInCut, FileInDelete, FileInDownload, FileInHardlink, FileInLink, FileInTrash, FileInUpload}, hook::{HookInOutBg, HookInOutBlock, HookInOutCut, HookInOutDelete, HookInOutDownload, HookInOutFetch, HookInOutOrphan, HookInOutTrash}, impl_from_in, plugin::PluginInEntry, prework::{PreworkInFetch, PreworkInLoad, PreworkInSize}, process::{ProcessInBg, ProcessInBlock, ProcessInOrphan}};
|
use crate::{file::{FileInCopy, FileInCut, FileInDelete, FileInDownload, FileInHardlink, FileInLink, FileInTrash, FileInUpload}, hook::{HookInOutCut, HookInOutDelete, HookInOutDownload, HookInOutTrash}, impl_from_in, plugin::PluginInEntry, prework::{PreworkInFetch, PreworkInLoad, PreworkInSize}, process::{ProcessInBg, ProcessInBlock, ProcessInOrphan}};
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub(crate) enum TaskIn {
|
pub(crate) enum TaskIn {
|
||||||
|
|
@ -28,10 +28,6 @@ pub(crate) enum TaskIn {
|
||||||
HookDelete(HookInOutDelete),
|
HookDelete(HookInOutDelete),
|
||||||
HookTrash(HookInOutTrash),
|
HookTrash(HookInOutTrash),
|
||||||
HookDownload(HookInOutDownload),
|
HookDownload(HookInOutDownload),
|
||||||
HookBlock(HookInOutBlock),
|
|
||||||
HookOrphan(HookInOutOrphan),
|
|
||||||
HookBg(HookInOutBg),
|
|
||||||
HookFetch(HookInOutFetch),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl_from_in! {
|
impl_from_in! {
|
||||||
|
|
@ -44,7 +40,7 @@ impl_from_in! {
|
||||||
// Process
|
// Process
|
||||||
ProcessBlock(ProcessInBlock), ProcessOrphan(ProcessInOrphan), ProcessBg(ProcessInBg),
|
ProcessBlock(ProcessInBlock), ProcessOrphan(ProcessInOrphan), ProcessBg(ProcessInBg),
|
||||||
// Hook
|
// Hook
|
||||||
HookCut(HookInOutCut), HookDelete(HookInOutDelete), HookTrash(HookInOutTrash), HookDownload(HookInOutDownload), HookBlock(HookInOutBlock), HookOrphan(HookInOutOrphan), HookBg(HookInOutBg), HookFetch(HookInOutFetch),
|
HookCut(HookInOutCut), HookDelete(HookInOutDelete), HookTrash(HookInOutTrash), HookDownload(HookInOutDownload),
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TaskIn {
|
impl TaskIn {
|
||||||
|
|
@ -74,10 +70,35 @@ impl TaskIn {
|
||||||
Self::HookDelete(r#in) => r#in.id,
|
Self::HookDelete(r#in) => r#in.id,
|
||||||
Self::HookTrash(r#in) => r#in.id,
|
Self::HookTrash(r#in) => r#in.id,
|
||||||
Self::HookDownload(r#in) => r#in.id,
|
Self::HookDownload(r#in) => r#in.id,
|
||||||
Self::HookBlock(r#in) => r#in.id,
|
}
|
||||||
Self::HookOrphan(r#in) => r#in.id,
|
}
|
||||||
Self::HookBg(r#in) => r#in.id,
|
|
||||||
Self::HookFetch(r#in) => r#in.id,
|
pub fn is_hook(&self) -> bool {
|
||||||
|
match self {
|
||||||
|
// File
|
||||||
|
Self::FileCopy(_) => false,
|
||||||
|
Self::FileCut(_) => false,
|
||||||
|
Self::FileLink(_) => false,
|
||||||
|
Self::FileHardlink(_) => false,
|
||||||
|
Self::FileDelete(_) => false,
|
||||||
|
Self::FileTrash(_) => false,
|
||||||
|
Self::FileDownload(_) => false,
|
||||||
|
Self::FileUpload(_) => false,
|
||||||
|
// Plugin
|
||||||
|
Self::PluginEntry(_) => false,
|
||||||
|
// Prework
|
||||||
|
Self::PreworkFetch(_) => false,
|
||||||
|
Self::PreworkLoad(_) => false,
|
||||||
|
Self::PreworkSize(_) => false,
|
||||||
|
// Process
|
||||||
|
Self::ProcessBlock(_) => false,
|
||||||
|
Self::ProcessOrphan(_) => false,
|
||||||
|
Self::ProcessBg(_) => false,
|
||||||
|
// Hook
|
||||||
|
Self::HookCut(_) => true,
|
||||||
|
Self::HookDelete(_) => true,
|
||||||
|
Self::HookTrash(_) => true,
|
||||||
|
Self::HookDownload(_) => true,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,21 @@ macro_rules! ok_or_not_found {
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[macro_export]
|
||||||
|
macro_rules! progress_or_break {
|
||||||
|
($rx:ident, $done:expr) => {
|
||||||
|
tokio::select! {
|
||||||
|
r = $rx.recv() => {
|
||||||
|
match r {
|
||||||
|
Some(prog) => prog,
|
||||||
|
None => break,
|
||||||
|
}
|
||||||
|
},
|
||||||
|
false = $done.future() => break,
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
#[macro_export]
|
#[macro_export]
|
||||||
macro_rules! impl_from_in {
|
macro_rules! impl_from_in {
|
||||||
($($variant:ident($type:ty)),* $(,)?) => {
|
($($variant:ident($type:ty)),* $(,)?) => {
|
||||||
|
|
|
||||||
|
|
@ -1,15 +1,15 @@
|
||||||
use hashbrown::HashMap;
|
use hashbrown::{HashMap, hash_map::RawEntryMut};
|
||||||
use ordered_float::OrderedFloat;
|
use ordered_float::OrderedFloat;
|
||||||
use yazi_config::YAZI;
|
use yazi_config::YAZI;
|
||||||
use yazi_parser::app::TaskSummary;
|
use yazi_parser::app::TaskSummary;
|
||||||
use yazi_shared::{Id, Ids};
|
use yazi_shared::{CompletionToken, Id, Ids};
|
||||||
|
|
||||||
use super::Task;
|
use super::Task;
|
||||||
use crate::TaskProg;
|
use crate::{TaskIn, TaskProg};
|
||||||
|
|
||||||
#[derive(Default)]
|
#[derive(Default)]
|
||||||
pub struct Ongoing {
|
pub struct Ongoing {
|
||||||
pub(super) inner: HashMap<Id, Task>,
|
inner: HashMap<Id, Task>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Ongoing {
|
impl Ongoing {
|
||||||
|
|
@ -23,11 +23,39 @@ impl Ongoing {
|
||||||
self.inner.entry(id).insert(Task::new::<T>(id, name)).into_mut()
|
self.inner.entry(id).insert(Task::new::<T>(id, name)).into_mut()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(super) fn cancel(&mut self, id: Id) -> Option<TaskIn> {
|
||||||
|
match self.inner.raw_entry_mut().from_key(&id) {
|
||||||
|
RawEntryMut::Occupied(mut oe) => {
|
||||||
|
let task = oe.get_mut();
|
||||||
|
task.done.complete(false);
|
||||||
|
|
||||||
|
if let Some(hook) = task.hook.take() {
|
||||||
|
return Some(hook);
|
||||||
|
}
|
||||||
|
|
||||||
|
oe.remove();
|
||||||
|
}
|
||||||
|
RawEntryMut::Vacant(_) => {}
|
||||||
|
}
|
||||||
|
None
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(super) fn fulfill(&mut self, id: Id) -> Option<Task> {
|
||||||
|
let task = self.inner.remove(&id)?;
|
||||||
|
task.done.complete(true);
|
||||||
|
Some(task)
|
||||||
|
}
|
||||||
|
|
||||||
#[inline]
|
#[inline]
|
||||||
pub fn get_mut(&mut self, id: Id) -> Option<&mut Task> { self.inner.get_mut(&id) }
|
pub fn get_mut(&mut self, id: Id) -> Option<&mut Task> { self.inner.get_mut(&id) }
|
||||||
|
|
||||||
pub fn get_id(&self, idx: usize) -> Option<Id> { self.values().nth(idx).map(|t| t.id) }
|
pub fn get_id(&self, idx: usize) -> Option<Id> { self.values().nth(idx).map(|t| t.id) }
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
pub fn get_token(&self, id: Id) -> Option<CompletionToken> {
|
||||||
|
self.inner.get(&id).map(|t| t.done.clone())
|
||||||
|
}
|
||||||
|
|
||||||
pub fn len(&self) -> usize {
|
pub fn len(&self) -> usize {
|
||||||
if YAZI.tasks.suppress_preload {
|
if YAZI.tasks.suppress_preload {
|
||||||
self.inner.values().filter(|&t| t.prog.is_user()).count()
|
self.inner.values().filter(|&t| t.prog.is_user()).count()
|
||||||
|
|
@ -40,7 +68,9 @@ impl Ongoing {
|
||||||
pub fn exists(&self, id: Id) -> bool { self.inner.contains_key(&id) }
|
pub fn exists(&self, id: Id) -> bool { self.inner.contains_key(&id) }
|
||||||
|
|
||||||
#[inline]
|
#[inline]
|
||||||
pub fn intact(&self, id: Id) -> bool { self.inner.get(&id).is_some_and(|t| !t.canceled) }
|
pub fn intact(&self, id: Id) -> bool {
|
||||||
|
self.inner.get(&id).is_some_and(|t| t.done.completed() != Some(false))
|
||||||
|
}
|
||||||
|
|
||||||
pub fn values(&self) -> Box<dyn Iterator<Item = &Task> + '_> {
|
pub fn values(&self) -> Box<dyn Iterator<Item = &Task> + '_> {
|
||||||
if YAZI.tasks.suppress_preload {
|
if YAZI.tasks.suppress_preload {
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
use crate::{Task, file::{FileOutCopy, FileOutCopyDo, FileOutCut, FileOutCutDo, FileOutDelete, FileOutDeleteDo, FileOutDownload, FileOutDownloadDo, FileOutHardlink, FileOutHardlinkDo, FileOutLink, FileOutTrash, FileOutUpload, FileOutUploadDo}, hook::{HookInOutBg, HookInOutBlock, HookInOutCut, HookInOutDelete, HookInOutDownload, HookInOutFetch, HookInOutOrphan, HookInOutTrash}, impl_from_out, plugin::PluginOutEntry, prework::{PreworkOutFetch, PreworkOutLoad, PreworkOutSize}, process::{ProcessOutBg, ProcessOutBlock, ProcessOutOrphan}};
|
use crate::{Task, file::{FileOutCopy, FileOutCopyDo, FileOutCut, FileOutCutDo, FileOutDelete, FileOutDeleteDo, FileOutDownload, FileOutDownloadDo, FileOutHardlink, FileOutHardlinkDo, FileOutLink, FileOutTrash, FileOutUpload, FileOutUploadDo}, hook::{HookInOutCut, HookInOutDelete, HookInOutDownload, HookInOutTrash}, impl_from_out, plugin::PluginOutEntry, prework::{PreworkOutFetch, PreworkOutLoad, PreworkOutSize}, process::{ProcessOutBg, ProcessOutBlock, ProcessOutOrphan}};
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub(super) enum TaskOut {
|
pub(super) enum TaskOut {
|
||||||
|
|
@ -32,10 +32,6 @@ pub(super) enum TaskOut {
|
||||||
HookDelete(HookInOutDelete),
|
HookDelete(HookInOutDelete),
|
||||||
HookTrash(HookInOutTrash),
|
HookTrash(HookInOutTrash),
|
||||||
HookDownload(HookInOutDownload),
|
HookDownload(HookInOutDownload),
|
||||||
HookBlock(HookInOutBlock),
|
|
||||||
HookOrphan(HookInOutOrphan),
|
|
||||||
HookBg(HookInOutBg),
|
|
||||||
HookFetch(HookInOutFetch),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl_from_out! {
|
impl_from_out! {
|
||||||
|
|
@ -48,7 +44,7 @@ impl_from_out! {
|
||||||
// Process
|
// Process
|
||||||
ProcessBlock(ProcessOutBlock), ProcessOrphan(ProcessOutOrphan), ProcessBg(ProcessOutBg),
|
ProcessBlock(ProcessOutBlock), ProcessOrphan(ProcessOutOrphan), ProcessBg(ProcessOutBg),
|
||||||
// Hook
|
// Hook
|
||||||
HookCut(HookInOutCut), HookDelete(HookInOutDelete), HookTrash(HookInOutTrash), HookDownload(HookInOutDownload), HookBlock(HookInOutBlock), HookOrphan(HookInOutOrphan), HookBg(HookInOutBg), HookFetch(HookInOutFetch),
|
HookCut(HookInOutCut), HookDelete(HookInOutDelete), HookTrash(HookInOutTrash), HookDownload(HookInOutDownload),
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TaskOut {
|
impl TaskOut {
|
||||||
|
|
@ -84,10 +80,6 @@ impl TaskOut {
|
||||||
Self::HookDelete(out) => out.reduce(task),
|
Self::HookDelete(out) => out.reduce(task),
|
||||||
Self::HookTrash(out) => out.reduce(task),
|
Self::HookTrash(out) => out.reduce(task),
|
||||||
Self::HookDownload(out) => out.reduce(task),
|
Self::HookDownload(out) => out.reduce(task),
|
||||||
Self::HookBlock(out) => out.reduce(task),
|
|
||||||
Self::HookOrphan(out) => out.reduce(task),
|
|
||||||
Self::HookBg(out) => out.reduce(task),
|
|
||||||
Self::HookFetch(out) => out.reduce(task),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -5,7 +5,6 @@ use crate::{Task, TaskProg};
|
||||||
pub(crate) enum PreworkOutFetch {
|
pub(crate) enum PreworkOutFetch {
|
||||||
Succ,
|
Succ,
|
||||||
Fail(String),
|
Fail(String),
|
||||||
Clean,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<mlua::Error> for PreworkOutFetch {
|
impl From<mlua::Error> for PreworkOutFetch {
|
||||||
|
|
@ -23,9 +22,6 @@ impl PreworkOutFetch {
|
||||||
prog.state = Some(false);
|
prog.state = Some(false);
|
||||||
task.log(reason);
|
task.log(reason);
|
||||||
}
|
}
|
||||||
Self::Clean => {
|
|
||||||
prog.cleaned = true;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,8 +4,7 @@ use yazi_parser::app::TaskSummary;
|
||||||
// --- Fetch
|
// --- Fetch
|
||||||
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
||||||
pub struct PreworkProgFetch {
|
pub struct PreworkProgFetch {
|
||||||
pub state: Option<bool>,
|
pub state: Option<bool>,
|
||||||
pub cleaned: bool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<PreworkProgFetch> for TaskSummary {
|
impl From<PreworkProgFetch> for TaskSummary {
|
||||||
|
|
@ -26,7 +25,7 @@ impl PreworkProgFetch {
|
||||||
|
|
||||||
pub fn failed(self) -> bool { self.state == Some(false) }
|
pub fn failed(self) -> bool { self.state == Some(false) }
|
||||||
|
|
||||||
pub fn cleaned(self) -> bool { self.cleaned }
|
pub fn cleaned(self) -> bool { false }
|
||||||
|
|
||||||
pub fn percent(self) -> Option<f32> { None }
|
pub fn percent(self) -> Option<f32> { None }
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,6 @@
|
||||||
use std::ffi::OsString;
|
use std::ffi::OsString;
|
||||||
|
|
||||||
use tokio::sync::mpsc;
|
use yazi_shared::{CompletionToken, Id, url::UrlCow};
|
||||||
use yazi_shared::{Id, url::UrlCow};
|
|
||||||
|
|
||||||
use super::ShellOpt;
|
use super::ShellOpt;
|
||||||
|
|
||||||
|
|
@ -38,11 +37,11 @@ impl From<ProcessInOrphan> for ShellOpt {
|
||||||
// --- Bg
|
// --- Bg
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub(crate) struct ProcessInBg {
|
pub(crate) struct ProcessInBg {
|
||||||
pub(crate) id: Id,
|
pub(crate) id: Id,
|
||||||
pub(crate) cwd: UrlCow<'static>,
|
pub(crate) cwd: UrlCow<'static>,
|
||||||
pub(crate) cmd: OsString,
|
pub(crate) cmd: OsString,
|
||||||
pub(crate) args: Vec<UrlCow<'static>>,
|
pub(crate) args: Vec<UrlCow<'static>>,
|
||||||
pub(crate) cancel: mpsc::Receiver<()>,
|
pub(crate) done: CompletionToken,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<ProcessInBg> for ShellOpt {
|
impl From<ProcessInBg> for ShellOpt {
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,6 @@ use crate::{Task, TaskProg};
|
||||||
pub(crate) enum ProcessOutBlock {
|
pub(crate) enum ProcessOutBlock {
|
||||||
Succ,
|
Succ,
|
||||||
Fail(String),
|
Fail(String),
|
||||||
Clean,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<std::io::Error> for ProcessOutBlock {
|
impl From<std::io::Error> for ProcessOutBlock {
|
||||||
|
|
@ -22,9 +21,6 @@ impl ProcessOutBlock {
|
||||||
prog.state = Some(false);
|
prog.state = Some(false);
|
||||||
task.log(reason);
|
task.log(reason);
|
||||||
}
|
}
|
||||||
Self::Clean => {
|
|
||||||
prog.cleaned = true;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -34,7 +30,6 @@ impl ProcessOutBlock {
|
||||||
pub(crate) enum ProcessOutOrphan {
|
pub(crate) enum ProcessOutOrphan {
|
||||||
Succ,
|
Succ,
|
||||||
Fail(String),
|
Fail(String),
|
||||||
Clean,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<anyhow::Error> for ProcessOutOrphan {
|
impl From<anyhow::Error> for ProcessOutOrphan {
|
||||||
|
|
@ -52,9 +47,6 @@ impl ProcessOutOrphan {
|
||||||
prog.state = Some(false);
|
prog.state = Some(false);
|
||||||
task.log(reason);
|
task.log(reason);
|
||||||
}
|
}
|
||||||
Self::Clean => {
|
|
||||||
prog.cleaned = true;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -65,7 +57,6 @@ pub(crate) enum ProcessOutBg {
|
||||||
Log(String),
|
Log(String),
|
||||||
Succ,
|
Succ,
|
||||||
Fail(String),
|
Fail(String),
|
||||||
Clean,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<anyhow::Error> for ProcessOutBg {
|
impl From<anyhow::Error> for ProcessOutBg {
|
||||||
|
|
@ -86,9 +77,6 @@ impl ProcessOutBg {
|
||||||
prog.state = Some(false);
|
prog.state = Some(false);
|
||||||
task.log(reason);
|
task.log(reason);
|
||||||
}
|
}
|
||||||
Self::Clean => {
|
|
||||||
prog.cleaned = true;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -55,14 +55,13 @@ impl Process {
|
||||||
})
|
})
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
|
let done = task.done;
|
||||||
let mut stdout = BufReader::new(child.stdout.take().unwrap()).lines();
|
let mut stdout = BufReader::new(child.stdout.take().unwrap()).lines();
|
||||||
let mut stderr = BufReader::new(child.stderr.take().unwrap()).lines();
|
let mut stderr = BufReader::new(child.stderr.take().unwrap()).lines();
|
||||||
let mut cancel = task.cancel;
|
|
||||||
loop {
|
loop {
|
||||||
select! {
|
select! {
|
||||||
_ = cancel.recv() => {
|
false = done.future() => {
|
||||||
child.start_kill().ok();
|
child.start_kill().ok();
|
||||||
cancel.close();
|
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
Ok(Some(line)) = stdout.next_line() => {
|
Ok(Some(line)) = stdout.next_line() => {
|
||||||
|
|
|
||||||
|
|
@ -4,8 +4,7 @@ use yazi_parser::app::TaskSummary;
|
||||||
// --- Block
|
// --- Block
|
||||||
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
||||||
pub struct ProcessProgBlock {
|
pub struct ProcessProgBlock {
|
||||||
pub state: Option<bool>,
|
pub state: Option<bool>,
|
||||||
pub cleaned: bool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<ProcessProgBlock> for TaskSummary {
|
impl From<ProcessProgBlock> for TaskSummary {
|
||||||
|
|
@ -26,7 +25,7 @@ impl ProcessProgBlock {
|
||||||
|
|
||||||
pub fn failed(self) -> bool { self.state == Some(false) }
|
pub fn failed(self) -> bool { self.state == Some(false) }
|
||||||
|
|
||||||
pub fn cleaned(self) -> bool { self.cleaned }
|
pub fn cleaned(self) -> bool { false }
|
||||||
|
|
||||||
pub fn percent(self) -> Option<f32> { None }
|
pub fn percent(self) -> Option<f32> { None }
|
||||||
}
|
}
|
||||||
|
|
@ -34,8 +33,7 @@ impl ProcessProgBlock {
|
||||||
// --- Orphan
|
// --- Orphan
|
||||||
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
||||||
pub struct ProcessProgOrphan {
|
pub struct ProcessProgOrphan {
|
||||||
pub state: Option<bool>,
|
pub state: Option<bool>,
|
||||||
pub cleaned: bool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<ProcessProgOrphan> for TaskSummary {
|
impl From<ProcessProgOrphan> for TaskSummary {
|
||||||
|
|
@ -56,7 +54,7 @@ impl ProcessProgOrphan {
|
||||||
|
|
||||||
pub fn failed(self) -> bool { self.state == Some(false) }
|
pub fn failed(self) -> bool { self.state == Some(false) }
|
||||||
|
|
||||||
pub fn cleaned(self) -> bool { self.cleaned }
|
pub fn cleaned(self) -> bool { false }
|
||||||
|
|
||||||
pub fn percent(self) -> Option<f32> { None }
|
pub fn percent(self) -> Option<f32> { None }
|
||||||
}
|
}
|
||||||
|
|
@ -64,8 +62,7 @@ impl ProcessProgOrphan {
|
||||||
// --- Bg
|
// --- Bg
|
||||||
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
|
||||||
pub struct ProcessProgBg {
|
pub struct ProcessProgBg {
|
||||||
pub state: Option<bool>,
|
pub state: Option<bool>,
|
||||||
pub cleaned: bool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<ProcessProgBg> for TaskSummary {
|
impl From<ProcessProgBg> for TaskSummary {
|
||||||
|
|
@ -86,7 +83,7 @@ impl ProcessProgBg {
|
||||||
|
|
||||||
pub fn failed(self) -> bool { self.state == Some(false) }
|
pub fn failed(self) -> bool { self.state == Some(false) }
|
||||||
|
|
||||||
pub fn cleaned(self) -> bool { self.cleaned }
|
pub fn cleaned(self) -> bool { false }
|
||||||
|
|
||||||
pub fn percent(self) -> Option<f32> { None }
|
pub fn percent(self) -> Option<f32> { None }
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -38,10 +38,6 @@ impl Runner {
|
||||||
TaskIn::HookDelete(r#in) => Ok(self.hook.delete(r#in).await),
|
TaskIn::HookDelete(r#in) => Ok(self.hook.delete(r#in).await),
|
||||||
TaskIn::HookTrash(r#in) => Ok(self.hook.trash(r#in).await),
|
TaskIn::HookTrash(r#in) => Ok(self.hook.trash(r#in).await),
|
||||||
TaskIn::HookDownload(r#in) => Ok(self.hook.download(r#in).await),
|
TaskIn::HookDownload(r#in) => Ok(self.hook.download(r#in).await),
|
||||||
TaskIn::HookBlock(r#in) => Ok(self.hook.block(r#in).await),
|
|
||||||
TaskIn::HookOrphan(r#in) => Ok(self.hook.orphan(r#in).await),
|
|
||||||
TaskIn::HookBg(r#in) => Ok(self.hook.bg(r#in).await),
|
|
||||||
TaskIn::HookFetch(r#in) => Ok(self.hook.fetch(r#in).await),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -71,10 +67,6 @@ impl Runner {
|
||||||
TaskIn::HookDelete(_in) => unreachable!(),
|
TaskIn::HookDelete(_in) => unreachable!(),
|
||||||
TaskIn::HookTrash(_in) => unreachable!(),
|
TaskIn::HookTrash(_in) => unreachable!(),
|
||||||
TaskIn::HookDownload(_in) => unreachable!(),
|
TaskIn::HookDownload(_in) => unreachable!(),
|
||||||
TaskIn::HookBlock(_in) => unreachable!(),
|
|
||||||
TaskIn::HookOrphan(_in) => unreachable!(),
|
|
||||||
TaskIn::HookBg(_in) => unreachable!(),
|
|
||||||
TaskIn::HookFetch(_in) => unreachable!(),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,14 +1,13 @@
|
||||||
use std::{sync::Arc, time::Duration};
|
use std::{sync::Arc, time::Duration};
|
||||||
|
|
||||||
use hashbrown::hash_map::RawEntryMut;
|
|
||||||
use parking_lot::Mutex;
|
use parking_lot::Mutex;
|
||||||
use tokio::{select, sync::{mpsc::{self, UnboundedReceiver}, oneshot}, task::JoinHandle};
|
use tokio::{select, sync::mpsc::{self, UnboundedReceiver}, task::JoinHandle};
|
||||||
use yazi_config::{YAZI, plugin::{Fetcher, Preloader}};
|
use yazi_config::{YAZI, plugin::{Fetcher, Preloader}};
|
||||||
use yazi_parser::{app::PluginOpt, tasks::ProcessOpenOpt};
|
use yazi_parser::{app::PluginOpt, tasks::ProcessOpenOpt};
|
||||||
use yazi_shared::{Id, Throttle, url::{UrlBuf, UrlLike}};
|
use yazi_shared::{CompletionToken, Id, Throttle, url::{UrlBuf, UrlLike}};
|
||||||
|
|
||||||
use super::{Ongoing, TaskOp};
|
use super::{Ongoing, TaskOp};
|
||||||
use crate::{HIGH, LOW, NORMAL, Runner, TaskIn, TaskOps, file::{File, FileInCopy, FileInCut, FileInDelete, FileInDownload, FileInHardlink, FileInLink, FileInTrash, FileInUpload, FileOutCopy, FileOutCut, FileOutDownload, FileOutHardlink, FileOutUpload, FileProgCopy, FileProgCut, FileProgDelete, FileProgDownload, FileProgHardlink, FileProgLink, FileProgTrash, FileProgUpload}, hook::{Hook, HookInOutBg, HookInOutBlock, HookInOutDelete, HookInOutDownload, HookInOutFetch, HookInOutOrphan, HookInOutTrash}, plugin::{Plugin, PluginInEntry, PluginProgEntry}, prework::{Prework, PreworkInFetch, PreworkInLoad, PreworkInSize, PreworkProgFetch, PreworkProgLoad, PreworkProgSize}, process::{Process, ProcessInBg, ProcessInBlock, ProcessInOrphan, ProcessProgBg, ProcessProgBlock, ProcessProgOrphan}};
|
use crate::{HIGH, LOW, NORMAL, Runner, TaskIn, TaskOps, file::{File, FileInCopy, FileInCut, FileInDelete, FileInDownload, FileInHardlink, FileInLink, FileInTrash, FileInUpload, FileOutCopy, FileOutCut, FileOutDownload, FileOutHardlink, FileOutUpload, FileProgCopy, FileProgCut, FileProgDelete, FileProgDownload, FileProgHardlink, FileProgLink, FileProgTrash, FileProgUpload}, hook::{Hook, HookInOutDelete, HookInOutDownload, HookInOutTrash}, plugin::{Plugin, PluginInEntry, PluginProgEntry}, prework::{Prework, PreworkInFetch, PreworkInLoad, PreworkInSize, PreworkProgFetch, PreworkProgLoad, PreworkProgSize}, process::{Process, ProcessInBg, ProcessInBlock, ProcessInOrphan, ProcessProgBg, ProcessProgBlock, ProcessProgOrphan}};
|
||||||
|
|
||||||
pub struct Scheduler {
|
pub struct Scheduler {
|
||||||
ops: TaskOps,
|
ops: TaskOps,
|
||||||
|
|
@ -54,22 +53,12 @@ impl Scheduler {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn cancel(&self, id: Id) -> bool {
|
pub fn cancel(&self, id: Id) -> bool {
|
||||||
let mut ongoing = self.ongoing.lock();
|
if let Some(hook) = self.ongoing.lock().cancel(id) {
|
||||||
|
self.micro.try_send(hook, HIGH).ok();
|
||||||
match ongoing.inner.raw_entry_mut().from_key(&id) {
|
return false;
|
||||||
RawEntryMut::Occupied(mut oe) => {
|
|
||||||
let task = oe.get_mut();
|
|
||||||
if let Some(hook) = task.hook.take() {
|
|
||||||
task.canceled = true;
|
|
||||||
self.micro.try_send(hook, HIGH).ok();
|
|
||||||
false
|
|
||||||
} else {
|
|
||||||
oe.remove();
|
|
||||||
true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
RawEntryMut::Vacant(_) => false,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
true
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn shutdown(&self) {
|
pub fn shutdown(&self) {
|
||||||
|
|
@ -90,7 +79,17 @@ impl Scheduler {
|
||||||
|
|
||||||
let follow = !from.scheme().covariant(to.scheme());
|
let follow = !from.scheme().covariant(to.scheme());
|
||||||
self.queue(
|
self.queue(
|
||||||
FileInCut { id: task.id, from, to, force, cha: None, follow, retry: 0, drop: None },
|
FileInCut {
|
||||||
|
id: task.id,
|
||||||
|
from,
|
||||||
|
to,
|
||||||
|
force,
|
||||||
|
cha: None,
|
||||||
|
follow,
|
||||||
|
retry: 0,
|
||||||
|
drop: None,
|
||||||
|
done: task.done.clone(),
|
||||||
|
},
|
||||||
LOW,
|
LOW,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
@ -106,7 +105,19 @@ impl Scheduler {
|
||||||
}
|
}
|
||||||
|
|
||||||
let follow = follow || !from.scheme().covariant(to.scheme());
|
let follow = follow || !from.scheme().covariant(to.scheme());
|
||||||
self.queue(FileInCopy { id: task.id, from, to, force, cha: None, follow, retry: 0 }, LOW);
|
self.queue(
|
||||||
|
FileInCopy {
|
||||||
|
id: task.id,
|
||||||
|
from,
|
||||||
|
to,
|
||||||
|
force,
|
||||||
|
cha: None,
|
||||||
|
follow,
|
||||||
|
retry: 0,
|
||||||
|
done: task.done.clone(),
|
||||||
|
},
|
||||||
|
LOW,
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn file_link(&self, from: UrlBuf, to: UrlBuf, relative: bool, force: bool) {
|
pub fn file_link(&self, from: UrlBuf, to: UrlBuf, relative: bool, force: bool) {
|
||||||
|
|
@ -164,21 +175,21 @@ impl Scheduler {
|
||||||
self.queue(FileInTrash { id: task.id, target }, LOW);
|
self.queue(FileInTrash { id: task.id, target }, LOW);
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn file_download(&self, url: UrlBuf, done: Option<oneshot::Sender<bool>>) {
|
pub fn file_download(&self, url: UrlBuf) -> CompletionToken {
|
||||||
let mut ongoing = self.ongoing.lock();
|
let mut ongoing = self.ongoing.lock();
|
||||||
let task = ongoing.add::<FileProgDownload>(format!("Download {}", url.display()));
|
let task = ongoing.add::<FileProgDownload>(format!("Download {}", url.display()));
|
||||||
|
|
||||||
if !url.kind().is_remote() {
|
task.set_hook(HookInOutDownload { id: task.id });
|
||||||
return self
|
if url.kind().is_remote() {
|
||||||
.ops
|
self.queue(
|
||||||
.out(task.id, FileOutDownload::Fail("Cannot download non-remote file".to_owned()));
|
FileInDownload { id: task.id, url, cha: None, retry: 0, done: task.done.clone() },
|
||||||
};
|
LOW,
|
||||||
|
);
|
||||||
if let Some(done) = done {
|
} else {
|
||||||
task.set_hook(HookInOutDownload { id: task.id, done });
|
self.ops.out(task.id, FileOutDownload::Fail("Cannot download non-remote file".to_owned()));
|
||||||
}
|
}
|
||||||
|
|
||||||
self.queue(FileInDownload { id: task.id, url, cha: None, retry: 0 }, LOW);
|
task.done.clone()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn file_upload(&self, url: UrlBuf) {
|
pub fn file_upload(&self, url: UrlBuf) {
|
||||||
|
|
@ -191,7 +202,10 @@ impl Scheduler {
|
||||||
.out(task.id, FileOutUpload::Fail("Cannot upload non-remote file".to_owned()));
|
.out(task.id, FileOutUpload::Fail("Cannot upload non-remote file".to_owned()));
|
||||||
};
|
};
|
||||||
|
|
||||||
self.queue(FileInUpload { id: task.id, url, cha: None, cache: None }, LOW);
|
self.queue(
|
||||||
|
FileInUpload { id: task.id, url, cha: None, cache: None, done: task.done.clone() },
|
||||||
|
LOW,
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn plugin_entry(&self, opt: PluginOpt) {
|
pub fn plugin_entry(&self, opt: PluginOpt) {
|
||||||
|
|
@ -205,8 +219,7 @@ impl Scheduler {
|
||||||
&self,
|
&self,
|
||||||
fetcher: &'static Fetcher,
|
fetcher: &'static Fetcher,
|
||||||
targets: Vec<yazi_fs::File>,
|
targets: Vec<yazi_fs::File>,
|
||||||
done: Option<oneshot::Sender<bool>>,
|
) -> CompletionToken {
|
||||||
) {
|
|
||||||
let mut ongoing = self.ongoing.lock();
|
let mut ongoing = self.ongoing.lock();
|
||||||
let task = ongoing.add::<PreworkProgFetch>(format!(
|
let task = ongoing.add::<PreworkProgFetch>(format!(
|
||||||
"Run fetcher `{}` with {} target(s)",
|
"Run fetcher `{}` with {} target(s)",
|
||||||
|
|
@ -214,24 +227,19 @@ impl Scheduler {
|
||||||
targets.len()
|
targets.len()
|
||||||
));
|
));
|
||||||
|
|
||||||
if let Some(done) = done {
|
|
||||||
task.set_hook(HookInOutFetch { id: task.id, done });
|
|
||||||
}
|
|
||||||
|
|
||||||
self.queue(PreworkInFetch { id: task.id, plugin: fetcher, targets }, NORMAL);
|
self.queue(PreworkInFetch { id: task.id, plugin: fetcher, targets }, NORMAL);
|
||||||
|
task.done.clone()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn fetch_mimetype(&self, targets: Vec<yazi_fs::File>) -> bool {
|
pub async fn fetch_mimetype(&self, targets: Vec<yazi_fs::File>) -> bool {
|
||||||
let mut wg = vec![];
|
let mut wg = vec![];
|
||||||
for (fetcher, targets) in YAZI.plugin.mime_fetchers(targets) {
|
for (fetcher, targets) in YAZI.plugin.mime_fetchers(targets) {
|
||||||
let (tx, rx) = oneshot::channel();
|
wg.push(self.fetch_paged(fetcher, targets));
|
||||||
self.fetch_paged(fetcher, targets, Some(tx));
|
|
||||||
wg.push(rx);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
for rx in wg {
|
for done in wg {
|
||||||
if rx.await != Ok(true) {
|
if !done.future().await {
|
||||||
return false; // Canceled or error
|
return false; // Canceled
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
true
|
true
|
||||||
|
|
@ -278,24 +286,24 @@ impl Scheduler {
|
||||||
ongoing.add::<ProcessProgBg>(name)
|
ongoing.add::<ProcessProgBg>(name)
|
||||||
};
|
};
|
||||||
|
|
||||||
|
if let Some(done) = opt.done {
|
||||||
|
task.done = done;
|
||||||
|
}
|
||||||
|
|
||||||
if opt.block {
|
if opt.block {
|
||||||
task.set_hook(HookInOutBlock { id: task.id, done: opt.done });
|
|
||||||
self
|
self
|
||||||
.queue(ProcessInBlock { id: task.id, cwd: opt.cwd, cmd: opt.cmd, args: opt.args }, NORMAL);
|
.queue(ProcessInBlock { id: task.id, cwd: opt.cwd, cmd: opt.cmd, args: opt.args }, NORMAL);
|
||||||
} else if opt.orphan {
|
} else if opt.orphan {
|
||||||
task.set_hook(HookInOutOrphan { id: task.id, done: opt.done });
|
|
||||||
self
|
self
|
||||||
.queue(ProcessInOrphan { id: task.id, cwd: opt.cwd, cmd: opt.cmd, args: opt.args }, NORMAL);
|
.queue(ProcessInOrphan { id: task.id, cwd: opt.cwd, cmd: opt.cmd, args: opt.args }, NORMAL);
|
||||||
} else {
|
} else {
|
||||||
let (cancel_tx, cancel_rx) = mpsc::channel(1);
|
|
||||||
task.set_hook(HookInOutBg { id: task.id, cancel: cancel_tx, done: opt.done });
|
|
||||||
self.queue(
|
self.queue(
|
||||||
ProcessInBg {
|
ProcessInBg {
|
||||||
id: task.id,
|
id: task.id,
|
||||||
cwd: opt.cwd,
|
cwd: opt.cwd,
|
||||||
cmd: opt.cmd,
|
cmd: opt.cmd,
|
||||||
args: opt.args,
|
args: opt.args,
|
||||||
cancel: cancel_rx,
|
done: task.done.clone(),
|
||||||
},
|
},
|
||||||
NORMAL,
|
NORMAL,
|
||||||
);
|
);
|
||||||
|
|
@ -311,11 +319,19 @@ impl Scheduler {
|
||||||
loop {
|
loop {
|
||||||
if let Ok((r#in, _)) = rx.recv().await {
|
if let Ok((r#in, _)) = rx.recv().await {
|
||||||
let id = r#in.id();
|
let id = r#in.id();
|
||||||
if !ongoing.lock().exists(id) {
|
let Some(token) = ongoing.lock().get_token(id) else {
|
||||||
continue;
|
continue;
|
||||||
}
|
};
|
||||||
|
|
||||||
|
let result = if r#in.is_hook() {
|
||||||
|
runner.micro(r#in).await
|
||||||
|
} else {
|
||||||
|
select! {
|
||||||
|
r = runner.micro(r#in) => r,
|
||||||
|
false = token.future() => Ok(())
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let result = runner.micro(r#in).await;
|
|
||||||
if let Err(out) = result {
|
if let Err(out) = result {
|
||||||
ops.out(id, out);
|
ops.out(id, out);
|
||||||
}
|
}
|
||||||
|
|
@ -341,11 +357,24 @@ impl Scheduler {
|
||||||
};
|
};
|
||||||
|
|
||||||
let id = r#in.id();
|
let id = r#in.id();
|
||||||
if !ongoing.lock().exists(id) {
|
let Some(token) = ongoing.lock().get_token(id) else {
|
||||||
continue;
|
continue;
|
||||||
}
|
};
|
||||||
|
|
||||||
|
let result = if r#in.is_hook() {
|
||||||
|
if micro { runner.micro(r#in).await } else { runner.r#macro(r#in).await }
|
||||||
|
} else if micro {
|
||||||
|
select! {
|
||||||
|
r = runner.micro(r#in) => r,
|
||||||
|
false = token.future() => Ok(()),
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
select! {
|
||||||
|
r = runner.r#macro(r#in) => r,
|
||||||
|
false = token.future() => Ok(()),
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let result = if micro { runner.micro(r#in).await } else { runner.r#macro(r#in).await };
|
|
||||||
if let Err(out) = result {
|
if let Err(out) = result {
|
||||||
ops.out(id, out);
|
ops.out(id, out);
|
||||||
}
|
}
|
||||||
|
|
@ -368,7 +397,7 @@ impl Scheduler {
|
||||||
} else if let Some(hook) = task.hook.take() {
|
} else if let Some(hook) = task.hook.take() {
|
||||||
micro.try_send(hook, LOW).ok();
|
micro.try_send(hook, LOW).ok();
|
||||||
} else {
|
} else {
|
||||||
ongoing.inner.remove(&op.id);
|
ongoing.fulfill(op.id);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
use yazi_shared::Id;
|
use yazi_shared::{CompletionToken, Id};
|
||||||
|
|
||||||
use crate::{TaskIn, TaskProg};
|
use crate::{TaskIn, TaskProg};
|
||||||
|
|
||||||
|
|
@ -9,7 +9,7 @@ pub struct Task {
|
||||||
pub name: String,
|
pub name: String,
|
||||||
pub(crate) prog: TaskProg,
|
pub(crate) prog: TaskProg,
|
||||||
pub(crate) hook: Option<TaskIn>,
|
pub(crate) hook: Option<TaskIn>,
|
||||||
pub canceled: bool,
|
pub done: CompletionToken,
|
||||||
|
|
||||||
pub logs: String,
|
pub logs: String,
|
||||||
pub logger: Option<mpsc::UnboundedSender<String>>,
|
pub logger: Option<mpsc::UnboundedSender<String>>,
|
||||||
|
|
@ -25,7 +25,7 @@ impl Task {
|
||||||
name,
|
name,
|
||||||
prog: T::default().into(),
|
prog: T::default().into(),
|
||||||
hook: None,
|
hook: None,
|
||||||
canceled: false,
|
done: CompletionToken::new(),
|
||||||
|
|
||||||
logs: Default::default(),
|
logs: Default::default(),
|
||||||
logger: Default::default(),
|
logger: Default::default(),
|
||||||
|
|
|
||||||
38
yazi-shared/src/completion_token.rs
Normal file
38
yazi-shared/src/completion_token.rs
Normal file
|
|
@ -0,0 +1,38 @@
|
||||||
|
use std::sync::{Arc, atomic::{AtomicU8, Ordering}};
|
||||||
|
|
||||||
|
use tokio::sync::Notify;
|
||||||
|
|
||||||
|
#[derive(Clone, Debug)]
|
||||||
|
pub struct CompletionToken {
|
||||||
|
inner: Arc<(AtomicU8, Notify)>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl CompletionToken {
|
||||||
|
pub fn new() -> Self { Self { inner: Arc::new((AtomicU8::new(0), Notify::new())) } }
|
||||||
|
|
||||||
|
pub fn complete(&self, success: bool) {
|
||||||
|
let new = if success { 1 } else { 2 };
|
||||||
|
self.inner.0.compare_exchange(0, new, Ordering::Relaxed, Ordering::Relaxed).ok();
|
||||||
|
self.inner.1.notify_waiters();
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn completed(&self) -> Option<bool> {
|
||||||
|
let state = self.inner.0.load(Ordering::Relaxed);
|
||||||
|
if state == 0 { None } else { Some(state == 1) }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn future(&self) -> bool {
|
||||||
|
loop {
|
||||||
|
if let Some(state) = self.completed() {
|
||||||
|
return state;
|
||||||
|
}
|
||||||
|
|
||||||
|
let notified = self.inner.1.notified();
|
||||||
|
if let Some(state) = self.completed() {
|
||||||
|
return state;
|
||||||
|
}
|
||||||
|
|
||||||
|
notified.await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
yazi_macro::mod_pub!(data errors event loc path pool scheme shell strand translit url wtf8);
|
yazi_macro::mod_pub!(data errors event loc path pool scheme shell strand translit url wtf8);
|
||||||
|
|
||||||
yazi_macro::mod_flat!(alias bytes chars condition debounce either env id layer localset natsort os predictor ro_cell source sync_cell terminal tests throttle time utf8);
|
yazi_macro::mod_flat!(alias bytes chars completion_token condition debounce either env id layer localset natsort os predictor ro_cell source sync_cell terminal tests throttle time utf8);
|
||||||
|
|
||||||
pub fn init() {
|
pub fn init() {
|
||||||
LOCAL_SET.with(tokio::task::LocalSet::new);
|
LOCAL_SET.with(tokio::task::LocalSet::new);
|
||||||
|
|
|
||||||
|
|
@ -55,7 +55,7 @@ pub(super) fn copy_with_progress_impl(
|
||||||
};
|
};
|
||||||
|
|
||||||
let chunks = (cha.len + PER_CHUNK - 1) / PER_CHUNK;
|
let chunks = (cha.len + PER_CHUNK - 1) / PER_CHUNK;
|
||||||
let mut result = futures::stream::iter(0..chunks)
|
let it = futures::stream::iter(0..chunks)
|
||||||
.map(|i| {
|
.map(|i| {
|
||||||
let acc_ = acc_.clone();
|
let acc_ = acc_.clone();
|
||||||
let (from, to) = (from.clone(), to.clone());
|
let (from, to) = (from.clone(), to.clone());
|
||||||
|
|
@ -104,8 +104,13 @@ pub(super) fn copy_with_progress_impl(
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
.buffer_unordered(4)
|
.buffer_unordered(4)
|
||||||
.try_fold(None, |first, file| async { Ok(first.or(file)) })
|
.try_fold(None, |first, file| async { Ok(first.or(file)) });
|
||||||
.await;
|
|
||||||
|
let mut result = select! {
|
||||||
|
r = it => r,
|
||||||
|
_ = prog_tx_.closed() => return,
|
||||||
|
};
|
||||||
|
done_tx.send(()).ok();
|
||||||
|
|
||||||
let n = acc_.swap(0, Ordering::SeqCst);
|
let n = acc_.swap(0, Ordering::SeqCst);
|
||||||
if n > 0 {
|
if n > 0 {
|
||||||
|
|
@ -122,8 +127,6 @@ pub(super) fn copy_with_progress_impl(
|
||||||
} else {
|
} else {
|
||||||
prog_tx_.send(Ok(0)).await.ok();
|
prog_tx_.send(Ok(0)).await.ok();
|
||||||
}
|
}
|
||||||
|
|
||||||
done_tx.send(()).ok();
|
|
||||||
});
|
});
|
||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue