Files
Thing/src-tauri/src/download_engine/engine.rs
T

2046 lines
85 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use librqbit::ManagedTorrent;
use tauri::{AppHandle, Emitter, Manager};
use super::http_dl::{HttpDownloader, split_segments};
use super::rate_limit::RateLimiter;
use super::storage::{EngineState, Storage};
use super::task::{DownloadTask, DownloaderSettings, ProbeResult, Segment, TaskProtocol, TaskStatus};
use super::torrent::{TorrentDownloader, hex_encode};
use specta::Type;
/// 重复类型
#[derive(Debug, Clone, PartialEq, serde::Serialize, Type)]
#[serde(rename_all = "camelCase")]
pub enum DuplicateKind {
/// 无重复
None,
/// URL 重复(已有相同链接的任务)
Url,
/// 文件名重复(已有同名任务下载到同一目录)
Filename,
/// 磁盘文件已存在
FileExists,
}
/// 已存在的任务信息(用于前端展示)
#[derive(Debug, Clone, serde::Serialize, Type)]
#[serde(rename_all = "camelCase")]
pub struct ExistingTaskInfo {
pub id: String,
pub filename: String,
pub status: TaskStatus,
}
/// check_url 命令返回的结果
#[derive(Debug, Clone, serde::Serialize, Type)]
#[serde(rename_all = "camelCase")]
pub struct CheckUrlResult {
/// 探测是否成功
pub ok: bool,
/// 错误信息(探测失败时)
pub error: Option<String>,
/// 文件名(探测成功时)
pub filename: Option<String>,
/// 文件大小(字节)
pub total_size: Option<u64>,
/// 是否支持断点续传
pub supports_resume: bool,
/// 重复类型
pub duplicate: DuplicateKind,
/// 已存在的任务信息
pub existing: Option<ExistingTaskInfo>,
}
/// 进度事件载荷(发给前端 download-progress 事件)
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct ProgressPayload {
pub id: String,
pub completed_size: u64,
pub total_size: u64,
pub speed: u64,
pub status: TaskStatus,
/// 每个分段的已下载字节(与任务 segments 一一对应,供前端详情弹窗实时展示)
pub segments: Vec<u64>,
}
/// 完成事件载荷(发给前端 download-complete 事件)
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct CompletePayload {
pub id: String,
pub filename: String,
pub status: TaskStatus,
pub error: Option<String>,
}
/// 活跃下载句柄
struct TaskHandle {
/// 代际号:同一任务每次 start_download 递增。
/// 旧代际任务完成时不删除新任务句柄、不覆盖新任务状态(防 pause→resume 竞态)
gen: u64,
cancel: Arc<AtomicBool>,
/// 每个分段的已下载字节(与 segments 一一对应;BT 任务为空)
progress: Vec<Arc<std::sync::atomic::AtomicU64>>,
/// 下载任务的 JoinHandleNone = 已完成/已取消)
join: Mutex<Option<tauri::async_runtime::JoinHandle<()>>>,
/// BT 任务的 librqbit 句柄(协议=BitTorrent 时存在;add_task 时创建)
bt: Option<Arc<ManagedTorrent>>,
}
/// 下载引擎(进程内运行,通过 Arc 共享状态)
#[derive(Clone)]
pub struct DownloadEngine {
inner: Arc<EngineInner>,
}
struct EngineInner {
/// 所有任务状态
tasks: Mutex<HashMap<String, DownloadTask>>,
/// 活跃下载句柄
handles: Mutex<HashMap<String, TaskHandle>>,
/// 持久化存储
storage: Storage,
/// 下载设置
settings: Mutex<DownloaderSettings>,
/// 全局限速器(bytes/s0=不限)
global_limiter: Arc<RateLimiter>,
/// HTTP 下载器
http: HttpDownloader,
/// BitTorrent 下载器(封装 librqbit 全局会话)
torrent: TorrentDownloader,
/// 下载是否使用代理(true=尝试走本地 mihomo 代理(mihomo 运行时才走),false=强制直连)。
/// 不再依赖系统代理。独立于 settings 存储,避免下载过程中反复锁 settingssave_settings 时同步更新
use_proxy: AtomicBool,
/// Tauri 应用句柄(用于发事件)
app_handle: AppHandle,
/// 引擎是否已启动
started: AtomicBool,
/// 任务代际计数器(每次 start_download 递增,分配给新句柄)
next_gen: AtomicU64,
/// 持久化节流:上次保存时间
last_save: Mutex<Instant>,
/// mihomo 代理地址解析缓存(含解析时刻),避免每个任务反复探测控制端
mihomo_proxy_cache: Mutex<Option<(Instant, Option<String>)>>,
/// HTTP 任务元数据后台探测中(占位任务已建但尚未完成探测/分段):
/// 期间 schedule 跳过此类任务,防止用"未知大小/单段"提前启动;探测完成后移除并重新调度
http_pending: Mutex<HashSet<String>>,
}
impl DownloadEngine {
/// 创建引擎并从磁盘恢复状态
pub fn new(data_dir: PathBuf, app_handle: AppHandle) -> Self {
// tracker 同步目录与引擎存储同目录(Storage::new 会移动 data_dir,先 clone
let tracker_data_dir = data_dir.clone();
let storage = Storage::new(data_dir);
let state = storage.load();
// 恢复时,将 Active 任务标记为 Paused(避免自动恢复下载)
let mut tasks: HashMap<String, DownloadTask> = state
.tasks
.into_iter()
.map(|t| (t.id.clone(), t))
.collect();
for task in tasks.values_mut() {
if task.status == TaskStatus::Active {
task.status = TaskStatus::Paused;
task.speed = 0;
}
}
let settings = state.settings;
// 创建全局限速器(KB/s → bytes/s
let global_limit = if settings.global_speed_limit > 0 {
settings.global_speed_limit * 1024
} else {
0
};
let global_limiter = Arc::new(RateLimiter::new(global_limit));
// 在 settings 移入 Mutex 前读取代理开关
let use_proxy = settings.use_proxy;
// BT 会话默认目录(下载目录,每个任务用 output_folder 覆盖)
// 解析 BT 代理地址(若开启且代理模块可用)
let bt_socks = bt_proxy_addr(&app_handle, settings.use_proxy);
let torrent = TorrentDownloader::new(PathBuf::from(&settings.download_dir));
// 同步 BT 专属设置(上传限速 / 监听端口 / 代理)
torrent.set_settings(settings.bt_upload_limit_kb, settings.bt_listen_port, bt_socks);
let engine = Self {
inner: Arc::new(EngineInner {
tasks: Mutex::new(tasks),
handles: Mutex::new(HashMap::new()),
storage,
settings: Mutex::new(settings),
global_limiter,
http: HttpDownloader::new(),
torrent,
use_proxy: AtomicBool::new(use_proxy),
app_handle,
started: AtomicBool::new(false),
next_gen: AtomicU64::new(0),
last_save: Mutex::new(Instant::now()),
mihomo_proxy_cache: Mutex::new(None),
http_pending: Mutex::new(HashSet::new()),
}),
};
// 启动后台持久化定时器
engine.start_save_timer();
// 后台:每日同步动态 tracker(打开软件检测上次同步日期,非今天则更新)
{
let torrent = engine.inner.torrent.clone();
let tr_dir = tracker_data_dir;
tauri::async_runtime::spawn(async move {
torrent.sync_dynamic_trackers(&tr_dir).await;
});
}
engine.inner.started.store(true, Ordering::SeqCst);
engine
}
/// 启动后台定时持久化(每 5 秒检查一次)
fn start_save_timer(&self) {
let engine = self.clone();
tauri::async_runtime::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(5));
loop {
interval.tick().await;
engine.persist_throttled();
}
});
}
// ===================== 公开 API =====================
/// 检查 URL 重复性并探测文件信息。
/// 返回 (探测结果, 重复类型, 已存在任务信息)。
pub async fn check_url(
&self,
url: &str,
dir: Option<&str>,
headers: &HashMap<String, String>,
) -> (Result<ProbeResult, String>, DuplicateKind, Option<ExistingTaskInfo>) {
// 磁力链 / 本地 .torrent → BitTorrent 检查(按 infohash 去重)
if detect_protocol(url) == TaskProtocol::BitTorrent {
return self.check_bt_url(url, dir).await;
}
let use_proxy = self.inner.use_proxy.load(Ordering::SeqCst);
self.ensure_mihomo_proxy().await;
let probe = self.inner.http.probe(url, headers, use_proxy).await;
let settings = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).clone();
let task_dir = dir.map(|d| d.to_string()).unwrap_or_else(|| settings.download_dir.clone());
let filename = probe.as_ref().ok()
.and_then(|p| p.filename.clone())
.unwrap_or_else(|| {
url.split('?').next()
.and_then(|u| u.rsplit('/').next())
.filter(|n| !n.is_empty())
.map(|n| n.to_string())
.unwrap_or_else(|| format!("download_{}", chrono::Utc::now().timestamp()))
});
let mut duplicate = DuplicateKind::None;
let mut existing: Option<ExistingTaskInfo> = None;
{
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
for t in tasks.values() {
// URL 完全相同
if t.url == url {
duplicate = DuplicateKind::Url;
existing = Some(ExistingTaskInfo {
id: t.id.clone(),
filename: t.filename.clone(),
status: t.status.clone(),
});
break;
}
// 目标文件名 + 目录相同(可能 URL 不同但下载到同一文件)
if t.filename == filename && t.dir == task_dir {
duplicate = DuplicateKind::Filename;
existing = Some(ExistingTaskInfo {
id: t.id.clone(),
filename: t.filename.clone(),
status: t.status.clone(),
});
break;
}
}
}
// 检查磁盘文件是否已存在(仅探测成功时)
if duplicate == DuplicateKind::None {
if let Ok(ref _p) = probe {
let filepath = PathBuf::from(&task_dir).join(&filename);
if filepath.exists() {
duplicate = DuplicateKind::FileExists;
existing = Some(ExistingTaskInfo {
id: String::new(),
filename,
status: TaskStatus::Complete,
});
}
}
}
(probe, duplicate, existing)
}
/// BitTorrent 链接检查:按 infohash 去重,磁力链立即返回不阻塞等待元数据。
/// librqbit 的 add_torrent 会阻塞等待元数据解析,对冷门种子可能超时。
/// 因此磁力链跳过 inspect,直接返回占位信息,元数据由 add_bt_task 后台异步解析。
async fn check_bt_url(
&self,
url: &str,
dir: Option<&str>,
) -> (Result<ProbeResult, String>, DuplicateKind, Option<ExistingTaskInfo>) {
let settings = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).clone();
let task_dir = dir.map(|d| d.to_string()).unwrap_or_else(|| settings.download_dir.clone());
// 磁力链:提取 infohash 做去重,不阻塞等待元数据
if url.trim().to_ascii_lowercase().starts_with("magnet:") {
let placeholder = bt_placeholder_name(url);
// 从磁力链提取 infohash(用于去重)
let info_hash = extract_magnet_hash(url);
let mut duplicate = DuplicateKind::None;
let mut existing: Option<ExistingTaskInfo> = None;
{
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
for t in tasks.values() {
if let Some(ref hash) = info_hash {
if t.info_hash.as_deref() == Some(hash.as_str()) {
duplicate = DuplicateKind::Url;
existing = Some(ExistingTaskInfo {
id: t.id.clone(),
filename: t.filename.clone(),
status: t.status.clone(),
});
break;
}
}
if t.filename == placeholder && t.dir == task_dir {
duplicate = DuplicateKind::Filename;
existing = Some(ExistingTaskInfo {
id: t.id.clone(),
filename: t.filename.clone(),
status: t.status.clone(),
});
break;
}
}
}
// 立即返回占位信息,不等待元数据解析
let probe = Ok(ProbeResult {
total_size: None,
supports_resume: true,
filename: Some(placeholder),
});
return (probe, duplicate, existing);
}
// .torrent 文件:可以立即解析元数据
match self.inner.torrent.inspect(url).await {
Ok(info) => {
let mut duplicate = DuplicateKind::None;
let mut existing: Option<ExistingTaskInfo> = None;
{
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
for t in tasks.values() {
if t.info_hash.as_deref() == Some(info.info_hash.as_str()) {
duplicate = DuplicateKind::Url;
existing = Some(ExistingTaskInfo {
id: t.id.clone(),
filename: t.filename.clone(),
status: t.status.clone(),
});
break;
}
if t.filename == info.name && t.dir == task_dir {
duplicate = DuplicateKind::Filename;
existing = Some(ExistingTaskInfo {
id: t.id.clone(),
filename: t.filename.clone(),
status: t.status.clone(),
});
break;
}
}
}
let probe = Ok(ProbeResult {
total_size: Some(info.total_size),
supports_resume: true,
filename: Some(info.name),
});
(probe, duplicate, existing)
}
Err(e) => (Err(e), DuplicateKind::None, None),
}
}
/// 生成不冲突的文件名(同名时追加 (1)、(2)...)
/// 冲突来源:磁盘已有同名最终文件、已存在同名临时文件(其他任务正在下载该名)、
/// 以及任务列表中已有任务占用同名(即使最终文件尚未落盘,避免共用同一 .thingdl 临时文件)
fn generate_unique_filename(&self, dir: &str, filename: &str) -> String {
// 任务列表中已占用的文件名(同一目录下)
let taken: std::collections::HashSet<String> = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks
.values()
.filter(|t| t.dir == dir)
.map(|t| t.filename.clone())
.collect()
};
// 磁盘最终文件 / 临时文件 / 任务占用,三者任一冲突即视为被占用
let used = |name: &str| {
taken.contains(name)
|| PathBuf::from(dir).join(name).exists()
|| PathBuf::from(dir).join(format!("{}.thingdl", name)).exists()
};
if !used(filename) {
return filename.to_string();
}
let path = PathBuf::from(dir).join(filename);
let stem = path.file_stem().and_then(|s| s.to_str()).unwrap_or("download");
let ext = path.extension().and_then(|s| s.to_str());
for i in 1..1000 {
let new_name = match ext {
Some(e) => format!("{} ({}).{}", stem, i, e),
None => format!("{} ({})", stem, i),
};
if !used(&new_name) {
return new_name;
}
}
// 极端情况:追加时间戳
match ext {
Some(e) => format!("{} ({}).{}", stem, chrono::Utc::now().timestamp_millis(), e),
None => format!("{} ({})", stem, chrono::Utc::now().timestamp_millis()),
}
}
/// 添加下载任务
/// auto_rename: 同名时自动重命名(追加 (1)、(2)...),否则覆盖
/// only_files: BitTorrent 任务指定要下载的种子文件索引(None=全部文件)
pub async fn add_task(
&self,
url: String,
filename: Option<String>,
dir: Option<String>,
headers: HashMap<String, String>,
auto_rename: bool,
only_files: Option<Vec<u32>>,
) -> Result<String, String> {
// 磁力链 / 本地 .torrent → BitTorrent 流程(独立处理)
if detect_protocol(&url) == TaskProtocol::BitTorrent {
return self.add_bt_task(url, dir, only_files).await;
}
// HTTP 任务采用"立即占位 + 后台探测"
// 点击添加后立即返回 id 并让任务出现在列表,探测/分段/调度全部放后台,
// 避免添加操作阻塞在慢速探测(如 mihomo 未就绪/网络不通)上,也防止用户因"没反应"多次点击。
let settings = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).clone();
let task_dir = dir.clone().unwrap_or_else(|| settings.download_dir.clone());
// 1) 后端 URL 去重:同 URL 且非终态任务已存在 → 直接复用既有 id,不重复创建。
// 接住多次点击/重复转发造成的重复添加。
if settings.check_duplicate {
let dup_id = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks
.values()
.find(|t| {
t.protocol == TaskProtocol::Http
&& t.url == url
&& matches!(
t.status,
TaskStatus::Queued | TaskStatus::Active | TaskStatus::Paused
)
})
.map(|t| t.id.clone())
};
if let Some(id) = dup_id {
return Ok(id);
}
}
let id = self.inner.storage.next_task_id();
// 2) 占位文件名:优先用调用方指定,否则先用 URL 推导,探测完成后用 Content-Disposition 修正
let mut task_filename = filename.clone().unwrap_or_else(|| {
url.split('?')
.next()
.and_then(|u| u.rsplit('/').next())
.filter(|n| !n.is_empty())
.map(|n| n.to_string())
.unwrap_or_else(|| format!("download_{}", chrono::Utc::now().timestamp()))
});
if auto_rename {
task_filename = self.generate_unique_filename(&task_dir, &task_filename);
}
// 3) 立即插入占位任务:Queued(未探测完成前 schedule 会跳过),未知大小单段
let task = DownloadTask {
id: id.clone(),
url: url.clone(),
filename: task_filename,
dir: task_dir,
protocol: TaskProtocol::Http,
info_hash: None,
bt_files: Vec::new(),
bt_metadata_ready: false,
status: TaskStatus::Queued,
total_size: 0,
completed_size: 0,
speed: 0,
supports_resume: false,
segments: vec![Segment {
index: 0,
start: 0,
end: 0,
completed: 0,
}],
error: None,
created_at: chrono::Utc::now().timestamp_millis(),
headers: headers.clone(),
};
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.insert(id.clone(), task);
}
// 标记该任务元数据探测中,避免被并发调度提前启动
self.inner
.http_pending
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(id.clone());
self.persist_now();
self.emit_task_added(&id);
// 4) 后台探测/分段/调度:不阻塞本函数,点击添加后立即返回
let engine = self.clone();
let rid = id.clone();
let raw_url = url.clone();
let raw_headers = headers.clone();
let use_proxy = self.inner.use_proxy.load(Ordering::SeqCst);
tauri::async_runtime::spawn(async move {
engine.ensure_mihomo_proxy().await;
let probe = engine.inner.http.probe(&raw_url, &raw_headers, use_proxy).await;
// 探测结束,解除 pending 门控
engine
.inner
.http_pending
.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(&rid);
let mut should_schedule = false;
// 先在锁内应用非文件名字段,并取出目录/探测文件名(释放锁后再做唯一化命名,
// 因为 generate_unique_filename 内部会再锁 tasks,持锁调用会死锁)
let rename_from: Option<(String, Option<String>)> = {
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
match tasks.get_mut(&rid) {
Some(t) => match &probe {
Ok(p) => {
t.total_size = p.total_size.unwrap_or(0);
let resume = settings.continue_download
&& p.supports_resume
&& p.total_size.map(|s| s > 0).unwrap_or(false);
t.supports_resume = resume;
t.segments = if resume {
split_segments(p.total_size.unwrap(), settings.max_connections)
} else {
vec![Segment {
index: 0,
start: 0,
end: p.total_size.map(|s| s.saturating_sub(1)).unwrap_or(0),
completed: 0,
}]
};
if t.status == TaskStatus::Queued {
should_schedule = true;
}
// 调用方未指定文件名时,才后续用 Content-Disposition 修正
if filename.is_none() {
Some((t.dir.clone(), p.filename.clone()))
} else {
None
}
}
Err(e) => {
if t.status == TaskStatus::Queued {
t.status = TaskStatus::Error;
t.error = Some(e.clone());
}
None
}
},
None => None,
}
};
// 释放 tasks 锁后,若探测返回了真实文件名,重新做唯一化并写回
if let Some((dir_owned, Some(real_name))) = rename_from {
let final_name = if auto_rename {
engine.generate_unique_filename(&dir_owned, &real_name)
} else {
real_name
};
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(&rid) {
t.filename = final_name;
}
}
engine.persist_now();
// 探测完成:刷新大小/文件名/分段/错误态到前端
engine.emit_task_added(&rid);
if should_schedule {
engine.schedule();
}
});
Ok(id)
}
/// 通知前端任务已新增/元数据已更新(用于刷新任务列表)
fn emit_task_added(&self, id: &str) {
let _ = self.inner.app_handle.emit(
crate::constants::events::DOWNLOAD_ADDED,
serde_json::json!({ "id": id }),
);
}
/// 添加 BitTorrent 任务(磁力链 / 本地 .torrent)— 全异步添加。
/// 立即创建占位任务并**立即返回** id(不阻塞);加入会话与元数据解析全在后台进行。
/// 注意:librqbit 的 add_torrent 对磁力会在元数据就绪前阻塞返回,因此绝不能在此 await。
async fn add_bt_task(&self, url: String, dir: Option<String>, _only_files: Option<Vec<u32>>) -> Result<String, String> {
let settings = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).clone();
let task_dir = dir.unwrap_or_else(|| settings.download_dir.clone());
let id = self.inner.storage.next_task_id();
let placeholder = bt_placeholder_name(&url);
// 1. 立即创建占位任务(元数据未就绪,bt_metadata_ready=false → 调度器跳过)
let task = DownloadTask {
id: id.clone(),
url: url.clone(),
filename: placeholder,
dir: task_dir.clone(),
protocol: TaskProtocol::BitTorrent,
info_hash: None,
bt_files: Vec::new(),
bt_metadata_ready: false,
status: TaskStatus::Queued,
total_size: 0,
completed_size: 0,
speed: 0,
supports_resume: true,
segments: Vec::new(),
error: None,
created_at: chrono::Utc::now().timestamp_millis(),
headers: HashMap::new(),
};
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.insert(id.clone(), task);
}
self.persist_now();
// 通知前端有新任务加入(任务已在列表中)
let _ = self.inner.app_handle.emit(crate::constants::events::DOWNLOAD_ADDED, serde_json::json!({ "id": id }));
// 2. 后台执行:加入会话(可能等待元数据)→ 解析元数据 → 就绪通知 / 失败置 Error
let engine = self.clone();
let id_for_task = id.clone();
tauri::async_runtime::spawn(async move {
engine.bt_resolve(id_for_task, url, task_dir).await;
});
// 立即返回,命令不等待任何网络操作
Ok(id)
}
/// BT 后台解析流程(独立 task 执行,不阻塞任何命令)
/// add_torrent 对磁力会等待元数据;后续 wait_initialized 亦可能等待,均带超时兜底。
async fn bt_resolve(&self, id: String, url: String, task_dir: String) {
// 若任务在添加过程中已被取消/删除,不再执行
let alive = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.get(&id).map(|t| t.status != TaskStatus::Cancelled).unwrap_or(false)
};
if !alive {
return;
}
// 加入会话并取得句柄(磁力元数据获取带超时)
let added = tokio::time::timeout(
std::time::Duration::from_secs(crate::download_engine::torrent::INSPECT_TIMEOUT_SECS),
self.inner.torrent.add_async(&url, &task_dir),
).await;
let (_, handle) = match added {
Ok(Ok(v)) => v,
Ok(Err(e)) => {
Self::fail_bt_resolve(self, &id, e);
return;
}
Err(_) => {
Self::fail_bt_resolve(self, &id, "解析磁力元数据超时:未能从 DHT/Tracker 获取种子信息".to_string());
return;
}
};
// 存入句柄(供调度 / 暂停 / 删除使用)
{
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(h) = handles.get_mut(&id) {
h.bt = Some(handle.clone());
} else {
handles.insert(
id.clone(),
TaskHandle {
gen: self.inner.next_gen.fetch_add(1, Ordering::SeqCst) + 1,
cancel: Arc::new(AtomicBool::new(false)),
progress: Vec::new(),
join: Mutex::new(None),
bt: Some(handle.clone()),
},
);
}
}
// 等待元数据就绪
if let Err(e) = self.inner.torrent.wait_initialized(&handle).await {
Self::fail_bt_resolve(self, &id, e);
return;
}
// 元数据已在会话缓存,重新解析结构化信息(此时秒回)
match self.inner.torrent.inspect(&url).await {
Ok(info) => {
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(&id) {
t.filename = info.name.clone();
t.info_hash = Some(info.info_hash.clone());
t.bt_files = info.files.clone();
t.total_size = info.total_size;
t.bt_metadata_ready = true;
t.error = None;
}
}
self.persist_now();
// 通知前端弹文件勾选对话框
let _ = self.inner.app_handle.emit(
"download-inspect-ready",
serde_json::json!({ "id": id }),
);
}
Err(e) => Self::fail_bt_resolve(self, &id, e),
}
}
/// BT 任务重新入会:句柄缺失(元数据未解析完/句柄被移除)时,重新加入会话以取回句柄。
/// librqbit 对已加入的种子返回 AlreadyManaged(复用现有句柄,已下载片段保留),
/// 因此对暂停→继续、错误→继续等场景都幂等、安全。
async fn reenter_bt(&self, id: &str, url: &str, dir: &str) -> Result<(), String> {
let added = tokio::time::timeout(
std::time::Duration::from_secs(crate::download_engine::torrent::INSPECT_TIMEOUT_SECS),
self.inner.torrent.add_async(url, dir),
).await;
let (_, handle) = match added {
Ok(Ok(v)) => v,
Ok(Err(e)) => return Err(e),
Err(_) => return Err("解析磁力元数据超时:未能从 DHT/Tracker 获取种子信息".to_string()),
};
// 存回句柄(句柄存在则更新,缺失则新建)
{
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
match handles.get_mut(id) {
Some(h) => h.bt = Some(handle.clone()),
None => {
handles.insert(
id.to_string(),
TaskHandle {
gen: self.inner.next_gen.fetch_add(1, Ordering::SeqCst) + 1,
cancel: Arc::new(AtomicBool::new(false)),
progress: Vec::new(),
join: Mutex::new(None),
bt: Some(handle.clone()),
},
);
}
}
}
// 等待元数据就绪(本地 .torrent 秒回;磁力依赖网络,带超时)
if let Err(e) = self.inner.torrent.wait_initialized(&handle).await {
return Err(e);
}
// 刷新结构化任务元数据(filename / infohash / 文件列表 / 大小)
match self.inner.torrent.inspect(url).await {
Ok(info) => {
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(id) {
t.filename = info.name.clone();
t.info_hash = Some(info.info_hash.clone());
t.bt_files = info.files.clone();
t.total_size = info.total_size;
t.bt_metadata_ready = true;
t.error = None;
// 复位为 Queuedstart_download 只在 Queued 下继续推进(首调时已置 Active)
if matches!(t.status, TaskStatus::Active | TaskStatus::Error) {
t.status = TaskStatus::Queued;
}
}
}
Err(e) => return Err(e),
}
self.persist_now();
Ok(())
}
/// BT 元数据解析失败:任务置为 Error 并通知前端刷新
fn fail_bt_resolve(engine: &DownloadEngine, id: &str, error: String) {
{
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(id) {
t.status = TaskStatus::Error;
t.error = Some(error);
t.speed = 0;
}
}
engine.persist_now();
// 复用 complete 事件通道让前端刷新错误状态
let task = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner()).get(id).cloned();
if let Some(task) = task {
let _ = engine.inner.app_handle.emit(
"download-complete",
CompletePayload {
id: id.to_string(),
filename: task.filename.clone(),
status: task.status.clone(),
error: task.error.clone(),
},
);
}
}
/// select_bt_files 命令:用户勾选文件后设置 only_files 并开始下载
pub async fn select_bt_files(&self, id: &str, only_files: Vec<u32>) -> Result<(), String> {
let handle = {
let handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(id).and_then(|h| h.bt.clone())
};
let Some(handle) = handle else {
return Err("任务句柄不存在".to_string());
};
// 空选择兜底:前端异常/竞态导致只传了空数组时,默认全选。
// 否则 set_only_files(空) 会让 librqbit 认为"无需下载任何文件"wait_until_completed 立即返回,
// 任务被秒判为"已完成 0%"。
let only_files = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
match tasks.get(id) {
Some(t) if only_files.is_empty() && !t.bt_files.is_empty() => {
t.bt_files.iter().map(|f| f.index).collect()
}
_ => only_files,
}
};
// 应用勾选(全选也要 set,确保覆盖 add 时的默认全选)
self.inner.torrent.set_only_files(&handle, &only_files).await?;
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(id) {
t.bt_metadata_ready = true;
t.error = None;
// 总大小改为勾选文件之和
t.total_size = only_files
.iter()
.filter_map(|fi| t.bt_files.iter().find(|f| f.index == *fi))
.map(|f| f.size)
.sum();
t.status = TaskStatus::Queued;
} else {
return Err("任务不存在".to_string());
}
}
self.persist_now();
self.schedule();
Ok(())
}
/// inspect 命令入口:解析磁力链 / 本地 .torrent / 远程 .torrent URL,返回种子信息供前端勾选文件
pub async fn inspect(&self, input: &str) -> Result<crate::download_engine::torrent::TorrentInfo, String> {
self.inner.torrent.inspect(input).await
}
/// 暂停任务
pub fn pause_task(&self, id: &str) -> Result<(), String> {
// 判断是否为 BT 任务(决定暂停方式)
let is_bt = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.get(id).map(|t| t.protocol == TaskProtocol::BitTorrent).unwrap_or(false)
};
if is_bt {
// BT:中止等待 future + 暂停种子(已下载片段保留,可继续)
let bt_handle = {
let handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(handle) = handles.get(id) {
if let Ok(mut join) = handle.join.lock() {
if let Some(j) = join.take() {
j.abort();
}
}
handle.bt.clone()
} else {
None
}
};
if let Some(bt) = &bt_handle {
if let Err(e) = tauri::async_runtime::block_on(self.inner.torrent.pause(bt)) {
crate::logger::log_error("download", &format!("暂停 BT 种子失败: {}", e));
}
}
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let Some(task) = tasks.get_mut(id) else {
return Err("任务不存在".to_string());
};
if task.status == TaskStatus::Active || task.status == TaskStatus::Queued {
task.status = TaskStatus::Paused;
task.speed = 0;
}
}
self.persist_now();
self.schedule();
self.emit_task_progress(id);
return Ok(());
}
// 1. 设置取消标志 + 读取进度(锁 handles)
let progress_values: Vec<u64> = {
let handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(handle) = handles.get(id) {
handle.cancel.store(true, Ordering::SeqCst);
handle.progress.iter().map(|p| p.load(Ordering::Relaxed)).collect()
} else {
Vec::new()
}
};
// 2. 更新任务状态 + 同步进度(锁 tasks,不嵌套锁 handles
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(task) = tasks.get_mut(id) {
if task.status == TaskStatus::Active || task.status == TaskStatus::Queued {
task.status = TaskStatus::Paused;
task.speed = 0;
// 同步分段进度
for (i, val) in progress_values.iter().enumerate() {
if let Some(seg) = task.segments.get_mut(i) {
seg.completed = *val;
}
}
task.recalc_completed();
}
} else {
return Err("任务不存在".to_string());
}
}
self.persist_now();
self.schedule();
self.emit_task_progress(id);
Ok(())
}
/// 恢复任务
pub fn resume_task(&self, id: &str) -> Result<(), String> {
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(task) = tasks.get_mut(id) {
if task.status != TaskStatus::Paused && task.status != TaskStatus::Error {
return Err("任务不在可恢复状态".to_string());
}
// 如果不支持断点续传,或文件大小未知(无法定位续传起点,需从头下载),从头开始
if !task.supports_resume || task.total_size == 0 {
task.completed_size = 0;
for seg in &mut task.segments {
seg.completed = 0;
}
}
task.status = TaskStatus::Queued;
task.error = None;
} else {
return Err("任务不存在".to_string());
}
}
self.persist_now();
self.schedule();
self.emit_task_progress(id);
Ok(())
}
/// 取消任务:置为 Cancelled、清空进度并删除下载文件,但保留记录(只能再次下载)
pub fn cancel_task(&self, id: &str) -> Result<(), String> {
// 判断是否为 BT 任务
let is_bt = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.get(id).map(|t| t.protocol == TaskProtocol::BitTorrent).unwrap_or(false)
};
if is_bt {
// BT:中止等待 future + 删除种子(连同已下载文件),保留记录
{
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(handle) = handles.remove(id) {
handle.cancel.store(true, Ordering::SeqCst);
if let Ok(mut join) = handle.join.lock() {
if let Some(j) = join.take() {
j.abort();
}
}
}
}
// 从任务表中取 infohash 用于删除
let info_hash = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.get(id).and_then(|t| t.info_hash.clone())
};
if let Some(hash) = info_hash {
// delete_files=true:取消即删除已下载文件(与 HTTP 取消语义一致)
if let Err(e) = tauri::async_runtime::block_on(self.inner.torrent.delete(&hash, true)) {
crate::logger::log_error("download", &format!("取消 BT 种子删除失败: {}", e));
}
}
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let Some(task) = tasks.get_mut(id) else {
return Err("任务不存在".to_string());
};
task.status = TaskStatus::Cancelled;
task.speed = 0;
task.completed_size = 0;
task.error = None;
}
self.persist_now();
self.schedule();
self.emit_task_progress(id);
return Ok(());
}
// 1. 中止活跃下载句柄
{
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(handle) = handles.remove(id) {
handle.cancel.store(true, Ordering::SeqCst);
if let Ok(mut join) = handle.join.lock() {
if let Some(j) = join.take() {
j.abort();
}
}
}
}
// 2. 标记已取消、清空进度/错误,并取出待删路径(不在锁内做 IO)
let paths = {
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let Some(task) = tasks.get_mut(id) else {
return Err("任务不存在".to_string());
};
task.status = TaskStatus::Cancelled;
task.speed = 0;
task.completed_size = 0;
task.error = None;
for seg in &mut task.segments {
seg.completed = 0;
}
Some((task.file_path(), task.temp_file_path()))
};
// 3. 删除下载文件(临时 + 最终),取消后不保留任何文件。
// abort() 只是请求异步取消,下载 future 里的文件句柄可能尚未释放,
// Windows 上删除被占用文件会失败。用短延时重试保证"取消即删除"可靠生效。
if let Some((fp, tfp)) = &paths {
for p in [tfp, fp] {
for _ in 0..5 {
if std::fs::remove_file(p).is_ok() {
break;
}
std::thread::sleep(Duration::from_millis(20));
}
}
}
self.persist_now();
self.schedule();
self.emit_task_progress(id);
Ok(())
}
/// 重新下载已取消/出错的任务:清空进度与文件记录后重新排队
pub async fn redownload(&self, id: &str) -> Result<(), String> {
let task = {
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let Some(task) = tasks.get_mut(id) else {
return Err("任务不存在".to_string());
};
if task.status != TaskStatus::Cancelled && task.status != TaskStatus::Error {
return Err("任务不在可重新下载状态".to_string());
}
task.completed_size = 0;
task.speed = 0;
task.error = None;
for seg in &mut task.segments {
seg.completed = 0;
}
task.status = TaskStatus::Queued;
task.clone()
};
// BT:取消时已从 librqbit 会话删除种子,需重新加入。
// 用 add_async 而非 startstart 对磁力会阻塞等待元数据(同步命令线程被挂住),
// 此处带与 bt_resolve 一致的超时,超时/失败统一回滚为 Error。
if task.protocol == TaskProtocol::BitTorrent {
let added = tokio::time::timeout(
std::time::Duration::from_secs(crate::download_engine::torrent::INSPECT_TIMEOUT_SECS),
self.inner.torrent.add_async(&task.url, &task.dir),
).await;
match added {
Ok(Ok((_, handle))) => {
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.insert(
id.to_string(),
TaskHandle {
gen: self.inner.next_gen.fetch_add(1, Ordering::SeqCst) + 1,
cancel: Arc::new(AtomicBool::new(false)),
progress: Vec::new(),
join: Mutex::new(None),
bt: Some(handle),
},
);
}
Ok(Err(e)) => {
// 重新加入失败:回滚为 Error
self.mark_bt_redownload_failed(id, e);
}
Err(_) => {
self.mark_bt_redownload_failed(
id,
"重新加入种子超时:未能获取元数据,请确认有做种源或网络可直连 BT".to_string(),
);
}
}
}
self.persist_now();
self.schedule();
self.emit_task_progress(id);
Ok(())
}
/// BT 重新加入会话失败:回滚为 Error 并通知前端
fn mark_bt_redownload_failed(&self, id: &str, error: String) {
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(id) {
t.status = TaskStatus::Error;
t.error = Some(error.clone());
}
}
crate::logger::log_error("download", &format!("BT 重新加入会话失败: {}", error));
self.persist_now();
self.emit_task_progress(id);
}
/// 移除任务
pub fn remove_task(&self, id: &str, delete_files: bool) -> Result<(), String> {
// 1. 设置取消标志 + abort join handle,并取出 BT 句柄(锁 handles
let bt_handle = {
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(handle) = handles.remove(id) {
handle.cancel.store(true, Ordering::SeqCst);
if let Ok(mut join) = handle.join.lock() {
if let Some(j) = join.take() {
j.abort();
}
}
handle.bt
} else {
None
}
};
// 2. 从任务列表移除(锁 tasks)
let task = {
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.remove(id)
};
// 3. BT:从 librqbit 会话删除种子(delete_files 决定是否同时删除已下载文件)
if let Some(bt) = &bt_handle {
let hash = hex_encode(&bt.info_hash().0);
if let Err(e) = tauri::async_runtime::block_on(self.inner.torrent.delete(&hash, delete_files)) {
crate::logger::log_error("download", &format!("删除 BT 种子失败: {}", e));
}
}
// 4. HTTP:可选删除文件(临时文件 + 最终文件)
if delete_files {
if let Some(task) = &task {
let _ = std::fs::remove_file(task.file_path());
let _ = std::fs::remove_file(task.temp_file_path());
}
}
// 5. 通知前端任务已被移除(模块任务列表/专用下载窗口据此刷新去除此条目)。
// 此前 remove_task 不发射该事件,仅 HTTP 扩展删除路径会发,导致 UI 无法即时刷新。
if task.is_some() {
let _ = self.inner.app_handle.emit(
crate::constants::events::DOWNLOAD_REMOVED,
serde_json::json!({ "id": id }),
);
}
self.persist_now();
self.schedule();
Ok(())
}
/// 获取所有任务
pub fn get_tasks(&self) -> Vec<DownloadTask> {
self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner()).values().cloned().collect()
}
/// 获取设置
pub fn get_settings(&self) -> DownloaderSettings {
self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).clone()
}
/// 保存设置
pub fn save_settings(&self, settings: DownloaderSettings) {
// 更新全局限速器
let new_limit = if settings.global_speed_limit > 0 {
settings.global_speed_limit * 1024
} else {
0
};
self.inner.global_limiter.set_limit(new_limit);
// 同步代理开关(新任务/探测立即生效,正在下载的任务不受影响)
self.inner
.use_proxy
.store(settings.use_proxy, Ordering::SeqCst);
// 同步 BT 专属设置(上传限速 / 监听端口 / 代理;监听端口、代理首次创建会话时生效)。
// BT 是否走代理跟随下载设置 use_proxybt_use_proxy 不再独立控制)
let bt_socks = bt_proxy_addr(&self.inner.app_handle, settings.use_proxy);
self.inner
.torrent
.set_settings(settings.bt_upload_limit_kb, settings.bt_listen_port, bt_socks);
{
let mut s = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner());
*s = settings;
}
self.persist_now();
}
/// 引擎是否已启动
pub fn is_started(&self) -> bool {
self.inner.started.load(Ordering::SeqCst)
}
/// 读取代理模块(mihomo)配置:mixed_port 与 external_controllersettings.json)。
/// 无配置/端口非法返回 None。
fn mihomo_proxy_config(&self) -> Option<(u16, String)> {
let dir = self.inner.app_handle.path().app_data_dir().ok()?;
let p = dir.join("proxy").join("settings.json");
let s = std::fs::read_to_string(p).ok()?;
let v: serde_json::Value = serde_json::from_str(&s).ok()?;
let port = v.get("mixedPort").and_then(|x| x.as_u64())?;
if port == 0 || port > 65535 {
return None;
}
let controller = v
.get("externalController")
.and_then(|x| x.as_str())
.unwrap_or("127.0.0.1:9090")
.to_string();
Some((port as u16, controller))
}
/// 解析当前 mihomo 显式代理地址(带 TTL 缓存):
/// 代理开关开启 & mihomo 控制端在线 → `http://127.0.0.1:{mixed_port}`
/// 否则 → None(直连)。不再依赖系统代理。
async fn mihomo_proxy(&self) -> Option<String> {
// 先查缓存(短锁,不跨 await 持有)
{
let cache = self
.inner
.mihomo_proxy_cache
.lock()
.unwrap_or_else(|e| e.into_inner());
if let Some((at, addr)) = cache.as_ref() {
if at.elapsed() < Duration::from_secs(5) {
return addr.clone();
}
}
}
// 重新解析:读配置端口 + 探测控制端确认 mihomo 运行
let addr = match self.mihomo_proxy_config() {
Some((port, controller)) => {
if self.inner.http.probe_controller(&controller).await {
Some(format!("http://127.0.0.1:{}", port))
} else {
None
}
}
None => None,
};
if let Ok(mut cache) = self.inner.mihomo_proxy_cache.lock() {
*cache = Some((Instant::now(), addr.clone()));
}
addr
}
/// 下载/探测前,把解析出的 mihomo 代理地址应用到 HTTP 下载器。
/// 代理开关关闭时强制直连。
async fn ensure_mihomo_proxy(&self) {
let use_proxy = self.inner.use_proxy.load(Ordering::SeqCst);
let proxy = if use_proxy { self.mihomo_proxy().await } else { None };
self.inner.http.configure_mihomo_proxy(proxy);
}
/// 退出时清理:停止所有下载、保存状态
pub fn cleanup_on_exit(&self) {
// 取消所有活跃下载(HTTP 设 cancel 标志;BT 中止等待 future
{
let handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
for handle in handles.values() {
handle.cancel.store(true, Ordering::SeqCst);
if let Ok(mut join) = handle.join.lock() {
if let Some(j) = join.take() {
j.abort();
}
}
}
}
// 将 Active 任务标记为 Paused
{
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
for task in tasks.values_mut() {
if task.status == TaskStatus::Active {
task.status = TaskStatus::Paused;
task.speed = 0;
}
}
}
// 等待短暂时间让下载任务退出
std::thread::sleep(Duration::from_millis(200));
// 停止 BitTorrent 会话(保存 DHT 状态并退出)
tauri::async_runtime::block_on(self.inner.torrent.stop());
// 最终保存
self.persist_now();
}
// ===================== 内部调度逻辑 =====================
/// 调度:如果活跃任务数 < max_concurrent,启动排队任务
fn schedule(&self) {
let max_concurrent = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).max_concurrent as usize;
let (active_count, queued_ids) = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let pending = self.inner.http_pending.lock().unwrap_or_else(|e| e.into_inner());
let active = tasks.values().filter(|t| t.status == TaskStatus::Active).count();
let mut queued: Vec<_> = tasks
.values()
// 跳过元数据未就绪的 BT 任务(异步添加后台解析中,等待用户勾选文件)
// 及 HTTP 占位任务(后台探测未完成,未得到大小/分段前不得提前启动)
.filter(|t| {
if t.status != TaskStatus::Queued {
return false;
}
if t.protocol == TaskProtocol::BitTorrent {
return t.bt_metadata_ready;
}
!pending.contains(&t.id)
})
.collect();
queued.sort_by_key(|t| t.created_at);
(active, queued.into_iter().map(|t| t.id.clone()).collect::<Vec<_>>())
};
if active_count >= max_concurrent {
return;
}
let slots = max_concurrent.saturating_sub(active_count);
for id in queued_ids.into_iter().take(slots) {
self.start_download(id);
}
}
/// 启动单个下载任务
fn start_download(&self, id: String) {
let task = {
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
match tasks.get_mut(&id) {
Some(task) if task.status == TaskStatus::Queued => {
task.status = TaskStatus::Active;
task.speed = 0;
task.clone()
}
_ => return,
}
};
let is_bt = task.protocol == TaskProtocol::BitTorrent;
// 快照做种开关(BT 下载完成后是否继续做种;下载开始前取值,运行中改动不影响进行中任务)
let task_do_seed = if is_bt {
self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).bt_seed_after_download
} else {
false
};
// 分配代际号(同一任务每次重启递增)
let gen = self.inner.next_gen.fetch_add(1, Ordering::SeqCst) + 1;
// 取出已有 BT 句柄(add_task 时创建;pause→resume 时复用)
let existing_bt = {
let handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(&id).and_then(|h| h.bt.clone())
};
// 创建取消标志和进度计数器(HTTP 按分段;BT 无分段)
let cancel = Arc::new(AtomicBool::new(false));
let progress: Vec<Arc<std::sync::atomic::AtomicU64>> = if is_bt {
Vec::new()
} else {
task.segments
.iter()
.map(|s| Arc::new(std::sync::atomic::AtomicU64::new(s.completed)))
.collect()
};
// 存储句柄:先取消旧代际(若存在),确保旧任务尽快退出,避免新旧并发写同一临时文件
{
let mut handles = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(old) = handles.get(&id) {
old.cancel.store(true, Ordering::SeqCst);
}
handles.insert(
id.clone(),
TaskHandle {
gen,
cancel: cancel.clone(),
progress: progress.iter().map(|p| p.clone()).collect(),
join: Mutex::new(None),
bt: if is_bt { existing_bt.clone() } else { None },
},
);
}
// ===== BitTorrent 下载分支 =====
if is_bt {
let Some(bt) = existing_bt else {
// 无 BT 句柄:常见于元数据尚未解析完就点了"暂停/继续",或句柄已被移除。
// 不能直接报"句柄缺失"(否则此类任务点"继续"会立即失败),改为后台重新加入
// 会话(如已加入则复用已有句柄,已下载片段保留),元数据就绪后再进入下载。
let engine2 = self.clone();
let id2 = id.clone();
let url2 = task.url.clone();
let dir2 = task.dir.clone();
tauri::async_runtime::spawn(async move {
if let Err(e) = engine2.reenter_bt(&id2, &url2, &dir2).await {
crate::logger::log_error("download", &format!("BT 任务重入会话失败: {}", e));
let mut tasks = engine2.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = tasks.get_mut(&id2) {
if t.status == TaskStatus::Active || t.status == TaskStatus::Queued {
t.status = TaskStatus::Error;
t.error = Some(e);
}
}
engine2.persist_now();
return;
}
engine2.start_download(id2);
});
return;
};
let engine = self.clone();
let id_clone = id.clone();
let torrent = self.inner.torrent.clone();
let bt_dl = bt.clone();
// 快照做种开关:关闭则下载完成后停止上传(做种)
let seed_after = task_do_seed.clone();
// 下载 futureunpause 开始下载,并等待全部文件完成
let join = tauri::async_runtime::spawn(async move {
let my_gen = gen;
if let Err(e) = torrent.unpause(&bt_dl).await {
Self::finalize_bt(&engine, &id_clone, my_gen, TaskStatus::Error, Some(e));
return;
}
let result = bt_dl.wait_until_completed().await;
// 下载完成但未开启做种:暂停种子停止上传
if result.is_ok() && !seed_after {
if let Err(e) = torrent.pause(&bt_dl).await {
crate::logger::log_error("download", &format!("下载完成停止做种失败: {}", e));
}
}
match result {
Ok(()) => Self::finalize_bt(&engine, &id_clone, my_gen, TaskStatus::Complete, None),
Err(e) => Self::finalize_bt(
&engine,
&id_clone,
my_gen,
TaskStatus::Error,
Some(e.to_string()),
),
}
});
// 存储 JoinHandle
if let Some(h) = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner()).get_mut(&id) {
if let Ok(mut join_guard) = h.join.lock() {
*join_guard = Some(join);
}
}
// 启动进度监控(每 500ms 从 librqbit stats 读取)
let engine = self.clone();
let id_monitor = id.clone();
let cancel_monitor = cancel;
let bt_monitor = bt;
tauri::async_runtime::spawn(async move {
let my_gen = gen;
let mut last_completed = TorrentDownloader::progress(&bt_monitor).0;
let mut last_time = Instant::now();
let mut interval = tokio::time::interval(Duration::from_millis(500));
interval.tick().await; // 跳过第一次立即触发
loop {
interval.tick().await;
// 代际守卫:pause→resume 后旧监控立即退出,避免用过期进度覆盖新任务
let is_current = {
let handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(&id_monitor).map(|h| h.gen == my_gen).unwrap_or(false)
};
if !is_current {
break;
}
let is_active = {
let tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.get(&id_monitor).map(|t| t.status == TaskStatus::Active).unwrap_or(false)
};
if !is_active || cancel_monitor.load(Ordering::SeqCst) {
break;
}
// 读取进度(librqbit statsfile_progress 每文件进度供详情页展示)
let (completed, total, file_progress) = TorrentDownloader::progress_full(&bt_monitor);
// 计算速度
let now = Instant::now();
let elapsed = now.duration_since(last_time).as_secs_f64();
let speed = if elapsed > 0.0 && completed >= last_completed {
((completed - last_completed) as f64 / elapsed) as u64
} else {
0
};
last_completed = completed;
last_time = now;
// 更新任务状态 + 发送进度事件
let live_status = {
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let Some(task) = tasks.get_mut(&id_monitor) else { break };
task.completed_size = completed;
task.total_size = total.max(task.total_size);
task.speed = speed;
// 每文件进度写入 segments(与 bt_files 一一对应,供详情页实时展示)
task.segments = file_progress
.iter()
.enumerate()
.map(|(i, v)| Segment {
index: i as u32,
start: 0,
end: task
.bt_files
.get(i)
.map(|f| f.size.saturating_sub(1))
.unwrap_or(0),
completed: *v,
})
.collect();
task.status.clone()
};
let _ = engine.inner.app_handle.emit(
"download-progress",
ProgressPayload {
id: id_monitor.clone(),
completed_size: completed,
total_size: total,
speed,
status: live_status,
segments: file_progress,
},
);
}
});
return;
}
// 生成下载 future
let engine = self.clone();
let id_clone = id.clone();
let cancel_clone = cancel.clone();
let progress_clone: Vec<Arc<std::sync::atomic::AtomicU64>> =
progress.iter().map(|p| p.clone()).collect();
let limiter = self.inner.global_limiter.clone();
let url = task.url.clone();
let headers = task.headers.clone();
let segments = task.segments.clone();
let temp_file_path = task.temp_file_path();
let final_file_path = task.file_path();
let http = self.inner.http.clone();
let use_proxy = self.inner.use_proxy.load(Ordering::SeqCst);
let join = tauri::async_runtime::spawn(async move {
let my_gen = gen;
// 下载前应用当前 mihomo 代理(engine 与 http 共享同一 mihomo 客户端状态)
engine.ensure_mihomo_proxy().await;
let result = http
.download(&url, &headers, &segments, &temp_file_path, cancel_clone, &progress_clone, limiter, use_proxy)
.await;
// 代际守卫:仅最新代际的任务能更新状态 / 移除句柄 / 发完成事件。
// pause→resume 后旧代际任务才退出,此时句柄已被新代际替换,
// 若仍按旧逻辑执行会覆盖新任务状态并误删新句柄(pause/remove 失效)
let is_current = {
let handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(&id_clone).map(|h| h.gen == my_gen).unwrap_or(false)
};
if is_current {
// 下载结束,更新任务状态
let mut final_status = match &result {
Ok(()) => TaskStatus::Complete,
Err(e) if e == "已取消" => TaskStatus::Paused,
Err(_) => TaskStatus::Error,
};
let mut final_error: Option<String> = match &result {
Err(e) if e != "已取消" => Some(e.clone()),
_ => None,
};
// 下载成功后,将临时文件重命名为最终文件名
if final_status == TaskStatus::Complete {
if let Err(e) = tokio::fs::rename(&temp_file_path, &final_file_path).await {
// 重命名失败(如目标被占用/路径不可写)→ 置 Error,
// 避免"标记完成但文件缺失"的状态不一致
final_status = TaskStatus::Error;
final_error = Some(format!("移动文件到最终路径失败: {}", e));
}
}
// 同步最终进度到任务
{
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(task) = tasks.get_mut(&id_clone) {
for (i, prog) in progress_clone.iter().enumerate() {
if let Some(seg) = task.segments.get_mut(i) {
seg.completed = prog.load(Ordering::Relaxed);
}
}
task.recalc_completed();
task.speed = 0;
task.status = final_status.clone();
if let Some(e) = final_error {
task.error = Some(e);
}
if final_status == TaskStatus::Complete {
task.completed_size = task.total_size.max(task.completed_size);
}
}
}
// 从活跃句柄中移除(仅移除自己代际的句柄)
{
let mut handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(h) = handles.get(&id_clone) {
if h.gen == my_gen {
handles.remove(&id_clone);
}
}
}
// 发送完成事件
let task = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner()).get(&id_clone).cloned();
if let Some(task) = task {
let _ = engine.inner.app_handle.emit(
"download-complete",
CompletePayload {
id: id_clone.clone(),
filename: task.filename.clone(),
status: task.status.clone(),
error: task.error.clone(),
},
);
}
// 持久化 + 调度下一个
engine.persist_now();
engine.schedule();
}
});
// 存储 JoinHandle
if let Some(h) = self.inner.handles.lock().unwrap_or_else(|e| e.into_inner()).get_mut(&id) {
if let Ok(mut join_guard) = h.join.lock() {
*join_guard = Some(join);
}
}
// 启动进度监控(每 500ms 更新一次)
let engine = self.clone();
let id_monitor = id.clone();
let progress_monitor = progress;
let cancel_monitor = cancel;
tauri::async_runtime::spawn(async move {
let my_gen = gen;
// 初始化为当前已下载量,避免恢复下载时首次计算速度异常
let mut last_completed: u64 = progress_monitor
.iter()
.map(|p| p.load(Ordering::Relaxed))
.sum();
let mut last_time = Instant::now();
let mut interval = tokio::time::interval(Duration::from_millis(500));
interval.tick().await; // 跳过第一次立即触发
loop {
interval.tick().await;
// 代际守卫:pause→resume 后旧监控立即退出,避免用过期进度覆盖新任务
let is_current = {
let handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(&id_monitor).map(|h| h.gen == my_gen).unwrap_or(false)
};
if !is_current {
break;
}
// 如果任务已不在活跃状态,停止监控
let is_active = {
let tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.get(&id_monitor).map(|t| t.status == TaskStatus::Active).unwrap_or(false)
};
if !is_active || cancel_monitor.load(Ordering::SeqCst) {
break;
}
// 读取进度
let completed: u64 = progress_monitor
.iter()
.map(|p| p.load(Ordering::Relaxed))
.sum();
// 各分段实时进度(供前端详情弹窗分段条展示)
let segment_completed: Vec<u64> = progress_monitor
.iter()
.map(|p| p.load(Ordering::Relaxed))
.collect();
// 计算速度
let now = Instant::now();
let elapsed = now.duration_since(last_time).as_secs_f64();
let speed = if elapsed > 0.0 && completed >= last_completed {
((completed - last_completed) as f64 / elapsed) as u64
} else {
0
};
last_completed = completed;
last_time = now;
// 更新任务状态 + 发送进度事件
let (total_size, live_status) = {
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(task) = tasks.get_mut(&id_monitor) {
task.completed_size = completed;
task.speed = speed;
for (i, prog) in progress_monitor.iter().enumerate() {
if let Some(seg) = task.segments.get_mut(i) {
seg.completed = prog.load(Ordering::Relaxed);
}
}
// 读取任务实时状态而非硬编码 Active:
// 暂停后监控循环 break 前可能发出的最后一次事件
// 必须携带 Paused,否则前端会把任务状态覆盖回"下载中"
(task.total_size, task.status.clone())
} else {
break;
}
};
let _ = engine.inner.app_handle.emit(
"download-progress",
ProgressPayload {
id: id_monitor.clone(),
completed_size: completed,
total_size,
speed,
status: live_status,
segments: segment_completed,
},
);
}
});
}
/// BT 任务收尾:更新状态/进度、移除句柄、发完成事件、持久化并调度下一个。
/// 带代际守卫,仅最新代际执行(与 HTTP 分支的收尾逻辑一致)。
fn finalize_bt(
engine: &DownloadEngine,
id: &str,
my_gen: u64,
status: TaskStatus,
error: Option<String>,
) {
// 代际守卫:仅最新代际的任务能更新状态/移除句柄/发完成事件
let is_current = {
let handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(id).map(|h| h.gen == my_gen).unwrap_or(false)
};
if !is_current {
return;
}
// 同步最终进度
let (completed, total) = {
let bt = {
let handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
handles.get(id).and_then(|h| h.bt.clone())
};
match bt {
Some(h) => TorrentDownloader::progress(&h),
None => (0, 0),
}
};
{
let mut tasks = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
if let Some(task) = tasks.get_mut(id) {
task.completed_size = completed;
task.total_size = total.max(task.total_size);
if status == TaskStatus::Complete {
task.completed_size = total.max(completed);
}
task.speed = 0;
task.status = status.clone();
if let Some(e) = error {
task.error = Some(e);
}
}
}
// 从活跃句柄中移除(仅移除自己代际的句柄)
{
let mut handles = engine.inner.handles.lock().unwrap_or_else(|e| e.into_inner());
if let Some(h) = handles.get(id) {
if h.gen == my_gen {
handles.remove(id);
}
}
}
// 发送完成事件
let task = engine.inner.tasks.lock().unwrap_or_else(|e| e.into_inner()).get(id).cloned();
if let Some(task) = task {
let _ = engine.inner.app_handle.emit(
"download-complete",
CompletePayload {
id: id.to_string(),
filename: task.filename.clone(),
status: task.status.clone(),
error: task.error.clone(),
},
);
}
// 持久化 + 调度下一个
engine.persist_now();
engine.schedule();
}
// ===================== 状态同步通知 =====================
/// 立即发送该任务的 download-progress 快照(供前端列表即时同步暂停/继续等状态变化)。
/// 暂停/恢复这类状态切换未必能即时触发下载监控循环发事件(恢复后可能等第一批数据
/// 才开始下载),因此在此主动补发一次,保证主界面与下载窗口状态一致。
fn emit_task_progress(&self, id: &str) {
let payload = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
let Some(task) = tasks.get(id) else {
return;
};
ProgressPayload {
id: id.to_string(),
completed_size: task.completed_size,
total_size: task.total_size,
speed: task.speed,
status: task.status.clone(),
segments: task.segments.iter().map(|s| s.completed).collect(),
}
};
let _ = self.inner.app_handle.emit("download-progress", payload);
}
// ===================== 持久化 =====================
/// 节流持久化(至少间隔 3 秒)
fn persist_throttled(&self) {
let should_save = {
let last = self.inner.last_save.lock().unwrap_or_else(|e| e.into_inner());
last.elapsed() >= Duration::from_secs(3)
};
if should_save {
self.persist_now();
}
}
/// 立即持久化。
/// 保存失败(磁盘不可写/rename 失败)时把进行中任务标记为 Error,
/// 防止用户误以为任务已持久化而关闭应用导致数据丢失。
fn persist_now(&self) {
let tasks: Vec<DownloadTask> = {
let tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
tasks.values().cloned().collect()
};
let settings = self.inner.settings.lock().unwrap_or_else(|e| e.into_inner()).clone();
let state = EngineState {
tasks,
settings,
next_id: 0, // storage.save 会从 id_counter 读取
};
if let Err(e) = self.inner.storage.save(state) {
crate::logger::log_error("download", &format!("状态持久化失败,进行中任务标记为 Error: {}", e));
let mut tasks = self.inner.tasks.lock().unwrap_or_else(|e| e.into_inner());
for task in tasks.values_mut() {
if task.status == TaskStatus::Active {
task.status = TaskStatus::Error;
}
}
return; // 不更新 last_save,下次定时器会重试
}
*self.inner.last_save.lock().unwrap_or_else(|e| e.into_inner()) = Instant::now();
}
}
/// 识别 URL 所属下载协议:磁力链 / .torrent(本地文件或 http(s) 链接) → BitTorrent,其余 → HTTP
fn detect_protocol(url: &str) -> TaskProtocol {
let trimmed = url.trim();
if trimmed.to_ascii_lowercase().starts_with("magnet:") {
return TaskProtocol::BitTorrent;
}
let path_lower = trimmed
.split('?')
.next()
.unwrap_or(trimmed)
.to_ascii_lowercase();
// 本地 .torrent 文件路径
if path_lower.ends_with(".torrent") && std::path::Path::new(trimmed.split('?').next().unwrap_or(trimmed)).exists() {
return TaskProtocol::BitTorrent;
}
// HTTP(S) 链接指向 .torrent 文件(如 https://.../xxx.iso.torrent)→ BitTorrent
if (trimmed.starts_with("http://") || trimmed.starts_with("https://"))
&& path_lower.ends_with(".torrent")
{
return TaskProtocol::BitTorrent;
}
TaskProtocol::Http
}
/// 磁力/种子任务的占位文件名(元数据未解析前显示)
/// 磁力 → 取 btih=infohash 前 12 位;本地 .torrent → 取文件名
fn bt_placeholder_name(input: &str) -> String {
let t = input.trim().to_ascii_lowercase();
if t.starts_with("magnet:") {
if let Some(pos) = t.find("btih:") {
let rest = &t[pos + 5..];
let hash = rest.split(['&', ':']).next().unwrap_or(rest);
if !hash.is_empty() {
let short: String = hash.chars().take(12).collect();
return format!("磁力任务 ({})", short);
}
}
return "磁力任务".to_string();
}
// .torrent 文件名
if let Some(p) = input.split('?').next() {
if let Some(name) = p.rsplit(['/', '\\']).next() {
if !name.is_empty() {
return name.to_string();
}
}
}
"BitTorrent 任务".to_string()
}
/// 从磁力链接中提取 infohash(小写十六进制),提取失败返回 None
fn extract_magnet_hash(url: &str) -> Option<String> {
let lower = url.trim().to_ascii_lowercase();
let pos = lower.find("btih:")?;
let rest = &lower[pos + 5..];
let hash = rest.split(['&', ':']).next().unwrap_or(rest);
if hash.len() == 40 {
Some(hash.to_string())
} else if hash.len() == 32 {
// Base32 编码,转为十六进制
use base64::Engine;
let bytes = base64::engine::general_purpose::STANDARD
.decode(hash.to_uppercase())
.ok()?;
Some(super::torrent::hex_encode(&bytes))
} else {
None
}
}
/// 解析 BitTorrent SOCKS5 代理地址:开启代理时读取代理模块(mihomo)设置的 mixed_port
/// 返回 socks5://127.0.0.1:{port}。代理模块未运行/无配置/无法解析时返回 None(降级直连)。
fn bt_proxy_addr(app: &tauri::AppHandle, enabled: bool) -> Option<String> {
if !enabled {
return None;
}
let dir = app.path().app_data_dir().ok()?;
let p = dir.join("proxy").join("settings.json");
let s = std::fs::read_to_string(p).ok()?;
let v: serde_json::Value = serde_json::from_str(&s).ok()?;
let port = v.get("mixedPort").and_then(|x| x.as_u64())?;
if port == 0 || port > 65535 {
return None;
}
Some(format!("socks5://127.0.0.1:{}", port))
}