feat: 提高下载任务的恢复效率

This commit is contained in:
lanyeeee
2026-08-15 04:16:27 +08:00
parent 4c563de2cf
commit 608366a0bd
7 changed files with 178 additions and 50 deletions
+40
View File
@@ -301,6 +301,7 @@ dependencies = [
"parking_lot 0.12.4", "parking_lot 0.12.4",
"prost", "prost",
"rand 0.9.1", "rand 0.9.1",
"rayon",
"reqwest", "reqwest",
"reqwest-middleware", "reqwest-middleware",
"reqwest-retry", "reqwest-retry",
@@ -678,6 +679,25 @@ dependencies = [
"crossbeam-utils", "crossbeam-utils",
] ]
[[package]]
name = "crossbeam-deque"
version = "0.8.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5181e0de7b61eb03a81e347d6dd8797bae9da5146707b51077e2d71a54ec0ceb"
dependencies = [
"crossbeam-epoch",
"crossbeam-utils",
]
[[package]]
name = "crossbeam-epoch"
version = "0.9.20"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2d6914041f254d6e9176c01941b21115dcfb7089e55135a35411081bd106ef3f"
dependencies = [
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-utils" name = "crossbeam-utils"
version = "0.8.21" version = "0.8.21"
@@ -3360,6 +3380,26 @@ version = "0.6.2"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "20675572f6f24e9e76ef639bc5552774ed45f1c30e2951e1e99c59888861c539" checksum = "20675572f6f24e9e76ef639bc5552774ed45f1c30e2951e1e99c59888861c539"
[[package]]
name = "rayon"
version = "1.12.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fb39b166781f92d482534ef4b4b1b2568f42613b53e5b6c160e24cfbfa30926d"
dependencies = [
"either",
"rayon-core",
]
[[package]]
name = "rayon-core"
version = "1.13.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "22e18b0f0062d30d4230b2e85ff77fdfe4326feb054b9783a3460d8435c8ab91"
dependencies = [
"crossbeam-deque",
"crossbeam-utils",
]
[[package]] [[package]]
name = "redox_syscall" name = "redox_syscall"
version = "0.2.16" version = "0.2.16"
+1
View File
@@ -60,6 +60,7 @@ memchr = { version = "2.7.5" }
md-5 = { version = "0.10.6" } md-5 = { version = "0.10.6" }
rand = { version = "0.9.1" } rand = { version = "0.9.1" }
base64 = { version = "0.22.1" } base64 = { version = "0.22.1" }
rayon = { version = "1.12.0" }
[profile.release] [profile.release]
strip = true strip = true
+5 -3
View File
@@ -11,6 +11,7 @@ use tracing::instrument;
use crate::{ use crate::{
config::Config, config::Config,
downloader::download_task::RestoredDownloadTask,
errors::{CommandError, CommandResult}, errors::{CommandError, CommandResult},
extensions::AppHandleExt, extensions::AppHandleExt,
logger, logger,
@@ -310,13 +311,14 @@ pub fn restart_download_task(app: AppHandle, params: RestartDownloadTaskParams)
#[tauri::command(async)] #[tauri::command(async)]
#[specta::specta] #[specta::specta]
#[instrument(level = "error", skip_all)] #[instrument(level = "error", skip_all)]
pub fn restore_download_tasks(app: AppHandle) -> CommandResult<()> { pub async fn restore_download_tasks(app: AppHandle) -> CommandResult<Vec<RestoredDownloadTask>> {
let download_manager = app.get_download_manager(); let download_manager = app.get_download_manager();
download_manager let restored_tasks = download_manager
.restore_download_tasks() .restore_download_tasks()
.await
.map_err(|err| CommandError::from("恢复下载任务失败", err))?; .map_err(|err| CommandError::from("恢复下载任务失败", err))?;
tracing::debug!("恢复下载任务成功"); tracing::debug!("恢复下载任务成功");
Ok(()) Ok(restored_tasks)
} }
#[tauri::command(async)] #[tauri::command(async)]
+78 -22
View File
@@ -10,12 +10,17 @@ use std::{
use eyre::{WrapErr, eyre}; use eyre::{WrapErr, eyre};
use parking_lot::RwLock; use parking_lot::RwLock;
use rayon::prelude::*;
use tauri::{AppHandle, Manager}; use tauri::{AppHandle, Manager};
use tauri_specta::Event; use tauri_specta::Event;
use tokio::sync::Semaphore; use tokio::{
use tracing::instrument; sync::{Semaphore, oneshot},
task::JoinSet,
};
use tracing::{Instrument, instrument};
use crate::{ use crate::{
downloader::download_progress::DownloadProgress,
events::DownloadEvent, events::DownloadEvent,
extensions::{AppHandleExt, EyreReportToMessage}, extensions::{AppHandleExt, EyreReportToMessage},
types::{ types::{
@@ -24,7 +29,7 @@ use crate::{
}, },
}; };
use super::{download_progress::DownloadProgress, download_task::DownloadTask}; use super::download_task::{DownloadTask, RestoredDownloadTask};
pub struct DownloadManager { pub struct DownloadManager {
pub app: AppHandle, pub app: AppHandle,
@@ -58,41 +63,92 @@ impl DownloadManager {
} }
#[instrument(level = "error", skip_all)] #[instrument(level = "error", skip_all)]
pub fn restore_download_tasks(&self) -> eyre::Result<()> { pub async fn restore_download_tasks(&self) -> eyre::Result<Vec<RestoredDownloadTask>> {
struct TaskFile {
path: PathBuf,
content: String,
}
let task_dir = self.get_task_dir()?; let task_dir = self.get_task_dir()?;
std::fs::create_dir_all(&task_dir) std::fs::create_dir_all(&task_dir)
.wrap_err(format!("创建下载任务目录`{}`失败", task_dir.display()))?; .wrap_err(format!("创建下载任务目录`{}`失败", task_dir.display()))?;
let mut tasks = self.download_tasks.write(); let mut join_set = JoinSet::new();
for entry in std::fs::read_dir(&task_dir)?.filter_map(Result::ok) { for entry in std::fs::read_dir(&task_dir)?.filter_map(Result::ok) {
let path = entry.path(); let read_task_file = async move {
let extension = path.extension().and_then(|s| s.to_str()); let path = entry.path();
if extension != Some("json") {
// 如果这个文件不是json则删除
let _ = std::fs::remove_file(&path);
continue;
}
let progress_json = std::fs::read_to_string(&path)?; let extension = path.extension().and_then(|s| s.to_str());
if extension != Some("json") {
// 如果这个文件不是json则删除
let _ = tokio::fs::remove_file(path).await;
return None;
}
let progress: DownloadProgress = let content = match tokio::fs::read_to_string(&path)
if let Ok(progress) = serde_json::from_str(&progress_json) { .await
progress .map_err(eyre::Report::from)
} else { {
// 如果这个json解析失败则删除 Ok(content) => content,
let _ = std::fs::remove_file(&path); Err(err) => {
continue; let err_title = format!("读取下载任务文件`{}`失败", path.display());
let message = err.to_message();
tracing::error!(err_title, message);
return None;
}
}; };
let new_task = DownloadTask::from_progress(self.app.clone(), progress); Some(TaskFile { path, content })
};
join_set.spawn(read_task_file.in_current_span());
}
let mut task_files = Vec::new();
while let Some(join_result) = join_set.join_next().await {
let Ok(Some(task_file)) = join_result else {
continue;
};
task_files.push(task_file);
}
let (progresses_sender, progresses_receiver) = oneshot::channel();
rayon::spawn(move || {
let progresses = task_files
.into_par_iter()
.filter_map(|task_file| {
if let Ok(progress) = serde_json::from_str(&task_file.content) {
Some(progress)
} else {
// 如果这个json解析失败则删除
let _ = std::fs::remove_file(&task_file.path);
None
}
})
.collect();
let _ = progresses_sender.send(progresses);
});
let progresses: Vec<DownloadProgress> = progresses_receiver.await?;
let mut tasks = self.download_tasks.write();
let mut restored_tasks = Vec::new();
for progress in progresses {
let new_task = DownloadTask::from_progress(self.app.clone(), progress.clone());
let state = *new_task.state_sender.borrow();
let old_task = tasks.insert(new_task.task_id.clone(), new_task); let old_task = tasks.insert(new_task.task_id.clone(), new_task);
if let Some(old_task) = old_task { if let Some(old_task) = old_task {
// 如果同一个ID的下载任务已经存在,则取消旧的任务 // 如果同一个ID的下载任务已经存在,则取消旧的任务
old_task.cancel(); old_task.cancel();
} }
restored_tasks.push(RestoredDownloadTask { state, progress });
} }
Ok(()) Ok(restored_tasks)
} }
pub fn create_download_tasks(&self, params: &CreateDownloadTaskParams) { pub fn create_download_tasks(&self, params: &CreateDownloadTaskParams) {
+17 -9
View File
@@ -1,7 +1,9 @@
use std::{sync::Arc, time::Duration}; use std::{sync::Arc, time::Duration};
use eyre::WrapErr; use eyre::WrapErr;
use parking_lot::RwLock; use parking_lot::RwLock;
use serde::Serialize;
use specta::Type;
use tauri::{AppHandle, Manager}; use tauri::{AppHandle, Manager};
use tauri_specta::Event; use tauri_specta::Event;
use tokio::{ use tokio::{
@@ -33,6 +35,12 @@ pub struct DownloadTask {
pub progress: RwLock<DownloadProgress>, pub progress: RwLock<DownloadProgress>,
} }
#[derive(Debug, Clone, Serialize, Type)]
pub struct RestoredDownloadTask {
pub state: DownloadTaskState,
pub progress: DownloadProgress,
}
impl DownloadTask { impl DownloadTask {
#[allow(clippy::too_many_lines)] #[allow(clippy::too_many_lines)]
#[instrument(level = "error", skip_all)] #[instrument(level = "error", skip_all)]
@@ -168,7 +176,7 @@ impl DownloadTask {
progress: RwLock::new(progress), progress: RwLock::new(progress),
}); });
tauri::async_runtime::spawn(task.clone().process()); tauri::async_runtime::spawn(task.clone().process(true));
tasks.push(task); tasks.push(task);
} }
@@ -198,7 +206,7 @@ impl DownloadTask {
progress: RwLock::new(progress), progress: RwLock::new(progress),
}); });
tauri::async_runtime::spawn(task.clone().process()); tauri::async_runtime::spawn(task.clone().process(false));
task task
} }
@@ -287,14 +295,14 @@ impl DownloadTask {
up_uid = self.trace_fields.up_uid, up_uid = self.trace_fields.up_uid,
) )
)] )]
async fn process(self: Arc<Self>) { async fn process(self: Arc<Self>, emit_create_events: bool) {
let state = *self.state_sender.borrow(); if emit_create_events {
let progress = self.progress.read().clone(); let state = *self.state_sender.borrow();
let _ = DownloadEvent::TaskCreate { state, progress }.emit(&self.app); let progress = self.progress.read().clone();
let _ = DownloadEvent::TaskCreate { state, progress }.emit(&self.app);
}
let mut state_receiver = self.state_sender.subscribe(); let mut state_receiver = self.state_sender.subscribe();
state_receiver.mark_changed();
let mut restart_receiver = self.restart_sender.subscribe(); let mut restart_receiver = self.restart_sender.subscribe();
let mut cancel_receiver = self.cancel_sender.subscribe(); let mut cancel_receiver = self.cancel_sender.subscribe();
let mut delete_receiver = self.delete_sender.subscribe(); let mut delete_receiver = self.delete_sender.subscribe();
+2 -1
View File
@@ -125,7 +125,7 @@ async restartDownloadTasks(taskIds: string[]) : Promise<void> {
async restartDownloadTask(params: RestartDownloadTaskParams) : Promise<void> { async restartDownloadTask(params: RestartDownloadTaskParams) : Promise<void> {
await TAURI_INVOKE("restart_download_task", { params }); await TAURI_INVOKE("restart_download_task", { params });
}, },
async restoreDownloadTasks() : Promise<Result<null, CommandError>> { async restoreDownloadTasks() : Promise<Result<RestoredDownloadTask[], CommandError>> {
try { try {
return { status: "ok", data: await TAURI_INVOKE("restore_download_tasks") }; return { status: "ok", data: await TAURI_INVOKE("restore_download_tasks") };
} catch (e) { } catch (e) {
@@ -429,6 +429,7 @@ export type RatingInBangumi = { count: number; score: number }
export type RatingInBangumiFollow = { score: number; count: number } export type RatingInBangumiFollow = { score: number; count: number }
export type RecommendSeason = { cover: string; ep_count: string; id: number; season_url: string; subtitle: string; title: string; view: number } export type RecommendSeason = { cover: string; ep_count: string; id: number; season_url: string; subtitle: string; title: string; view: number }
export type RestartDownloadTaskParams = { task_id: string; video_task_selected: boolean; audio_task_selected: boolean; merge_selected: boolean; embed_chapter_selected: boolean; embed_skip_selected: boolean; subtitle_task_selected: boolean; xml_danmaku_selected: boolean; ass_danmaku_selected: boolean; json_danmaku_selected: boolean; cover_task_selected: boolean; nfo_task_selected: boolean; json_task_selected: boolean; video_quality: VideoQuality; codec_type: CodecType; audio_quality: AudioQuality } export type RestartDownloadTaskParams = { task_id: string; video_task_selected: boolean; audio_task_selected: boolean; merge_selected: boolean; embed_chapter_selected: boolean; embed_skip_selected: boolean; subtitle_task_selected: boolean; xml_danmaku_selected: boolean; ass_danmaku_selected: boolean; json_danmaku_selected: boolean; cover_task_selected: boolean; nfo_task_selected: boolean; json_task_selected: boolean; video_quality: VideoQuality; codec_type: CodecType; audio_quality: AudioQuality }
export type RestoredDownloadTask = { state: DownloadTaskState; progress: DownloadProgress }
export type Rights = { bp: number; elec: number; download: number; movie: number; pay: number; hd5: number; no_reprint: number; autoplay: number; ugc_pay: number; is_cooperation: number; ugc_pay_preview: number; no_background: number; clean_mode: number; is_stein_gate: number; is_360: number; no_share: number; arc_pay: number; free_watch: number } export type Rights = { bp: number; elec: number; download: number; movie: number; pay: number; hd5: number; no_reprint: number; autoplay: number; ugc_pay: number; is_cooperation: number; ugc_pay_preview: number; no_background: number; clean_mode: number; is_stein_gate: number; is_360: number; no_share: number; arc_pay: number; free_watch: number }
export type RightsInBangumi = { allow_bp: number; allow_bp_rank: number; allow_download: number; allow_review: number; area_limit: number; ban_area_show: number; can_watch: number; copyright: string; forbid_pre: number; freya_white: number; is_cover_show: number; is_preview: number; only_vip_download: number; resource: string; watch_platform: number } export type RightsInBangumi = { allow_bp: number; allow_bp_rank: number; allow_download: number; allow_review: number; area_limit: number; ban_area_show: number; can_watch: number; copyright: string; forbid_pre: number; freya_white: number; is_cover_show: number; is_preview: number; only_vip_download: number; resource: string; watch_platform: number }
export type RightsInBangumiEp = { allow_dm: number; allow_download: number; area_limit: number } export type RightsInBangumiEp = { allow_dm: number; allow_download: number; area_limit: number }
+35 -15
View File
@@ -33,7 +33,7 @@ onMounted(async () => {
...progress, ...progress,
state, state,
percentage: 0, percentage: 0,
stateIndicator: '', stateIndicator: getStateIndicator(state),
taskIndicator: '', taskIndicator: '',
} }
store.updateProgresses((progresses) => { store.updateProgresses((progresses) => {
@@ -48,21 +48,8 @@ onMounted(async () => {
return return
} }
let stateIndicator = ''
if (state === 'Pending') {
stateIndicator = '排队中'
} else if (state === 'Downloading') {
stateIndicator = '下载中'
} else if (state === 'Paused') {
stateIndicator = '已暂停'
} else if (state === 'Completed') {
stateIndicator = '下载完成'
} else if (state === 'Failed') {
stateIndicator = '下载失败'
}
progressData.state = state progressData.state = state
progressData.stateIndicator = stateIndicator progressData.stateIndicator = getStateIndicator(state)
}) })
} else if (event === 'TaskSleeping') { } else if (event === 'TaskSleeping') {
const { task_id, remaining_sec } = data const { task_id, remaining_sec } = data
@@ -151,8 +138,41 @@ onMounted(async () => {
const result = await commands.restoreDownloadTasks() const result = await commands.restoreDownloadTasks()
if (result.status === 'error') { if (result.status === 'error') {
console.error(result.error) console.error(result.error)
return
} }
store.updateProgresses((progresses) => {
for (const { state, progress } of result.data) {
const progressData: ProgressData = {
...progress,
state,
percentage: 0,
stateIndicator: getStateIndicator(state),
taskIndicator: '',
}
progresses.set(progress.task_id, progressData)
}
})
}) })
function getStateIndicator(state: DownloadTaskState) {
let stateIndicator = ''
if (state === 'Pending') {
stateIndicator = '排队中'
} else if (state === 'Downloading') {
stateIndicator = '下载中'
} else if (state === 'Paused') {
stateIndicator = '已暂停'
} else if (state === 'Completed') {
stateIndicator = '下载完成'
} else if (state === 'Failed') {
stateIndicator = '下载失败'
}
return stateIndicator
}
</script> </script>
<template> <template>