feat(backend): implement download queue and resolve legacy parity gaps

- Add SSRF protection and migrate sensitive credentials to tempfiles
- Refactor aria2c process management to use ephemeral JSON-RPC
- Implement sysinfo recursive process termination for graceful stopping
- Add DownloadManager and MPSC queue for strict concurrency limits
- Configure yt-dlp to delegate download chunking to aria2c
This commit is contained in:
NimBold
2026-06-13 12:19:05 +03:30
parent f8c89e2f88
commit c79dcfa74d
4 changed files with 383 additions and 106 deletions
+1
View File
@@ -3909,6 +3909,7 @@ dependencies = [
"tauri-plugin-notification",
"tauri-plugin-opener",
"tauri-plugin-single-instance",
"tempfile",
"tokio",
]
+2 -1
View File
@@ -25,7 +25,7 @@ serde = { version = "1", features = ["derive"] }
serde_json = "1"
tokio = { version = "1", features = ["process", "io-util", "rt", "macros"] }
regex = "1.10"
reqwest = "0.12"
reqwest = { version = "0.12", features = ["json"] }
tauri-plugin-notification = "2.3.3"
sysinfo = "0.39.3"
keyring = "3"
@@ -33,3 +33,4 @@ hmac = "0.12"
sha2 = "0.10"
tauri-plugin-deep-link = "2"
tauri-plugin-single-instance = { version = "2", features = ["deep-link"] }
tempfile = "3"
+310 -105
View File
@@ -29,6 +29,26 @@ async fn fetch_metadata(url: String, user_agent: Option<String>, username: Optio
let client = builder.build().map_err(|e| e.to_string())?;
if let Ok(parsed) = reqwest::Url::parse(&url) {
if let Some(host) = parsed.host_str() {
let port = parsed.port_or_known_default().unwrap_or(80);
if let Ok(addrs) = std::net::ToSocketAddrs::to_socket_addrs(&(host, port)) {
for addr in addrs {
let ip = addr.ip();
if ip.is_loopback() || ip.is_multicast() || ip.is_unspecified() {
return Err("SSRF blocked: Private/local IP not allowed".to_string());
}
match ip {
std::net::IpAddr::V4(ipv4) if ipv4.is_private() || ipv4.is_link_local() => {
return Err("SSRF blocked: Private/local IP not allowed".to_string());
}
_ => {}
}
}
}
}
}
let mut head_req = client.head(&url);
if let Some(ref user) = username {
if !user.is_empty() {
@@ -114,16 +134,25 @@ async fn fetch_media_metadata(app_handle: tauri::AppHandle, url: String, cookie_
cmd.arg("--cookies-from-browser").arg(&browser);
}
}
let mut config_file = tempfile::Builder::new().prefix("ytdlp-").suffix(".conf").tempfile().map_err(|e| e.to_string())?;
let mut config_content = String::new();
if let Some(user) = username {
if !user.is_empty() {
cmd.arg("--username").arg(&user);
config_content.push_str(&format!("--username\n{}\n", user));
}
}
if let Some(pass) = password {
if !pass.is_empty() {
cmd.arg("--password").arg(&pass);
config_content.push_str(&format!("--password\n{}\n", pass));
}
}
use std::io::Write;
config_file.write_all(config_content.as_bytes()).map_err(|e| e.to_string())?;
let config_path = config_file.into_temp_path();
if !config_content.is_empty() {
cmd.arg("--config-locations").arg(&config_path);
}
cmd.arg(&url);
@@ -347,10 +376,13 @@ use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, RwLock};
struct AppState {
tasks: Mutex<HashMap<String, u32>>,
extension_pairing_token: extension_server::SharedExtensionToken,
extension_frontend_ready: extension_server::SharedFrontendReady,
mod queue;
pub struct AppState {
pub tasks: Arc<Mutex<HashMap<String, u32>>>,
pub cmd_tx: tokio::sync::mpsc::UnboundedSender<queue::DownloadCommand>,
pub extension_pairing_token: extension_server::SharedExtensionToken,
pub extension_frontend_ready: extension_server::SharedFrontendReady,
}
#[derive(Clone, Serialize)]
@@ -372,9 +404,32 @@ fn collect_download_uris(url: &str, mirrors: Option<&str>) -> Vec<String> {
uris
}
async fn rpc_call(port: u16, secret: &str, method: &str, params: serde_json::Value) -> Result<serde_json::Value, String> {
let url = format!("http://127.0.0.1:{}/jsonrpc", port);
let mut payload = serde_json::Map::new();
payload.insert("jsonrpc".to_string(), serde_json::json!("2.0"));
payload.insert("id".to_string(), serde_json::json!("1"));
payload.insert("method".to_string(), serde_json::json!(method));
let mut p = vec![serde_json::json!(format!("token:{}", secret))];
if let serde_json::Value::Array(arr) = params {
p.extend(arr);
}
payload.insert("params".to_string(), serde_json::json!(p));
let client = reqwest::Client::new();
let res = client.post(&url)
.json(&payload)
.send()
.await
.map_err(|e| e.to_string())?;
let json: serde_json::Value = res.json().await.map_err(|e| e.to_string())?;
Ok(json)
}
#[tauri::command]
async fn start_download(
app_handle: tauri::AppHandle,
state: tauri::State<'_, AppState>,
id: String,
url: String,
@@ -392,7 +447,38 @@ async fn start_download(
max_tries: Option<i32>,
proxy: Option<String>,
) -> Result<(), String> {
println!("start_download called for id: {}", id);
let task = queue::DownloadTask {
is_media: false,
id, url, destination, filename, connections, speed_limit, username, password, headers, checksum, cookies, mirrors, user_agent, max_tries, proxy,
format_selector: None, cookie_source: None,
};
state.cmd_tx.send(queue::DownloadCommand::Enqueue(task)).map_err(|e| e.to_string())?;
Ok(())
}
pub(crate) async fn start_download_internal(
app_handle: tauri::AppHandle,
tasks_map: Arc<Mutex<HashMap<String, u32>>>,
cmd_tx: tokio::sync::mpsc::UnboundedSender<queue::DownloadCommand>,
task: queue::DownloadTask,
) -> Result<(), String> {
let id = task.id;
let url = task.url;
let destination = task.destination;
let filename = task.filename;
let connections = task.connections;
let speed_limit = task.speed_limit;
let username = task.username;
let password = task.password;
let headers = task.headers;
let checksum = task.checksum;
let cookies = task.cookies;
let mirrors = task.mirrors;
let user_agent = task.user_agent;
let max_tries = task.max_tries;
let proxy = task.proxy;
println!("start_download_internal called for id: {}", id);
let resource_dir = app_handle.path().resource_dir().map_err(|e| e.to_string())?;
let aria2c_path = resource_dir.join("binaries").join("aria2c");
@@ -414,11 +500,21 @@ async fn start_download(
}
}
let rpc_port = std::net::TcpListener::bind("127.0.0.1:0")
.and_then(|listener| listener.local_addr())
.map(|addr| addr.port())
.unwrap_or(6800);
let rpc_secret = format!("{:x}", std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos());
let mut cmd = AsyncCommand::new(&aria2c_path);
cmd.arg("--enable-rpc=false")
cmd.arg("--enable-rpc=true")
.arg(format!("--rpc-listen-port={}", rpc_port))
.arg("--rpc-listen-all=false")
.arg("--continue=true")
.arg("--allow-overwrite=false")
.arg("--summary-interval=1")
.arg("--console-log-level=warn")
.arg("--download-result=hide")
.arg("--check-certificate=false")
.arg(format!("--dir={}", resolved_dest.to_string_lossy()))
.arg(format!("--out={}", filename));
@@ -432,21 +528,44 @@ async fn start_download(
cmd.arg(format!("--max-overall-download-limit={}", limit));
}
}
let mut config_file = tempfile::Builder::new().prefix("aria2-").suffix(".conf").tempfile().map_err(|e| e.to_string())?;
let mut config_content = String::new();
if let Some(user) = username {
if !user.is_empty() {
cmd.arg(format!("--http-user={}", user));
cmd.arg(format!("--ftp-user={}", user));
config_content.push_str(&format!("http-user={}\nftp-user={}\n", user, user));
}
}
if let Some(pass) = password {
if !pass.is_empty() {
cmd.arg(format!("--http-passwd={}", pass));
cmd.arg(format!("--ftp-passwd={}", pass));
config_content.push_str(&format!("http-passwd={}\nftp-passwd={}\n", pass, pass));
}
}
if let Some(hdr) = headers {
if !hdr.is_empty() {
cmd.arg(format!("--header={}", hdr));
for header in hdr.lines().map(str::trim).filter(|h| !h.is_empty()) {
config_content.push_str(&format!("header={}\n", header));
}
}
if let Some(cks) = cookies {
if !cks.is_empty() {
config_content.push_str(&format!("header=Cookie: {}\n", cks));
}
}
if let Some(p) = proxy {
if !p.is_empty() {
config_content.push_str(&format!("all-proxy={}\n", p));
}
}
use std::io::Write;
config_file.write_all(config_content.as_bytes()).map_err(|e| e.to_string())?;
let config_path = config_file.into_temp_path();
if !config_content.is_empty() {
cmd.arg(format!("--conf-path={}", config_path.display()));
}
if let Some(ua) = user_agent {
if !ua.is_empty() {
cmd.arg(format!("--user-agent={}", ua));
}
}
if let Some(chk) = checksum {
@@ -454,24 +573,9 @@ async fn start_download(
cmd.arg(format!("--checksum={}", chk));
}
}
if let Some(cks) = cookies {
if !cks.is_empty() {
cmd.arg(format!("--header=Cookie: {}", cks));
}
}
if let Some(ua) = user_agent {
if !ua.is_empty() {
cmd.arg(format!("--user-agent={}", ua));
}
}
if let Some(tries) = max_tries {
cmd.arg(format!("--max-tries={}", tries));
}
if let Some(p) = proxy {
if !p.is_empty() {
cmd.arg(format!("--all-proxy={}", p));
}
}
for uri in collect_download_uris(&url, mirrors.as_deref()) {
cmd.arg(uri);
@@ -485,64 +589,94 @@ async fn start_download(
let mut child = cmd.spawn().map_err(|e| format!("Failed to spawn aria2c: {}", e))?;
let pid = child.id().unwrap_or(0);
state.tasks.lock().unwrap().insert(id.clone(), pid);
tasks_map.lock().unwrap().insert(id.clone(), pid);
let stdout = child.stdout.take().unwrap();
let app_handle_clone = app_handle.clone();
let id_clone = id.clone();
let id_clone_tx = id.clone();
tokio::spawn(async move {
let mut reader = BufReader::new(stdout).lines();
let percentage_re = Regex::new(r"\((\d+)%\)").unwrap();
let speed_re = Regex::new(r"DL:([^\s\]]+)").unwrap();
let eta_re = Regex::new(r"ETA:([^\s\]]+)").unwrap();
let _keep_alive = config_path; // Keep the temp file alive
loop {
tokio::select! {
line_result = reader.next_line() => {
match line_result {
Ok(Some(line)) => {
if line.contains("DL:") {
let fraction = percentage_re.captures(&line)
.and_then(|cap| cap.get(1))
.and_then(|m| m.as_str().parse::<f64>().ok())
.unwrap_or(0.0) / 100.0;
tokio::time::sleep(std::time::Duration::from_millis(1000)).await;
let speed = speed_re.captures(&line)
.and_then(|cap| cap.get(1))
.map(|m| m.as_str().to_string())
.unwrap_or_else(|| "-".to_string());
if let Ok(Some(status)) = child.try_wait() {
if status.success() {
let _ = app_handle_clone.emit("download-complete", id_clone.clone());
} else {
let _ = app_handle_clone.emit("download-failed", id_clone.clone());
}
break;
}
let eta = eta_re.captures(&line)
.and_then(|cap| cap.get(1))
.map(|m| m.as_str().to_string())
.unwrap_or_else(|| "-".to_string());
if let Ok(json) = rpc_call(rpc_port, &rpc_secret, "aria2.tellActive", serde_json::json!([["status", "completedLength", "totalLength", "downloadSpeed"]])).await {
if let Some(arr) = json.get("result").and_then(|r| r.as_array()) {
if let Some(item) = arr.first() {
let completed = item.get("completedLength").and_then(|v| v.as_str()).and_then(|s| s.parse::<f64>().ok()).unwrap_or(0.0);
let total = item.get("totalLength").and_then(|v| v.as_str()).and_then(|s| s.parse::<f64>().ok()).unwrap_or(1.0);
let speed_bytes = item.get("downloadSpeed").and_then(|v| v.as_str()).and_then(|s| s.parse::<f64>().ok()).unwrap_or(0.0);
let fraction = if total > 0.0 { completed / total } else { 0.0 };
let speed = if speed_bytes > 1024.0 * 1024.0 {
format!("{:.1} MB/s", speed_bytes / (1024.0 * 1024.0))
} else if speed_bytes > 1024.0 {
format!("{:.1} KB/s", speed_bytes / 1024.0)
} else {
format!("{:.0} B/s", speed_bytes)
};
let _ = app_handle_clone.emit("download-progress", DownloadProgressEvent {
id: id_clone.clone(),
fraction,
speed,
eta,
});
let eta = if speed_bytes > 0.0 && total > completed {
let seconds = (total - completed) / speed_bytes;
if seconds > 3600.0 {
format!("{:.0}h {:.0}m", seconds / 3600.0, (seconds % 3600.0) / 60.0)
} else if seconds > 60.0 {
format!("{:.0}m {:.0}s", seconds / 60.0, seconds % 60.0)
} else {
format!("{:.0}s", seconds)
}
}
_ => break, // EOF or error
} else {
"-".to_string()
};
let _ = app_handle_clone.emit("download-progress", DownloadProgressEvent {
id: id_clone.clone(),
fraction,
speed,
eta,
});
continue;
}
}
status = child.wait() => {
println!("child exit status: {:?}", status);
if let Ok(exit_status) = status {
if exit_status.success() {
let _ = app_handle_clone.emit("download-complete", id_clone.clone());
} else {
// If it exited with error, emit failed (7 means paused/aborted usually, but we emit failed anyway, UI can filter)
let _ = app_handle_clone.emit("download-failed", id_clone.clone());
}
}
if let Ok(json) = rpc_call(rpc_port, &rpc_secret, "aria2.tellStopped", serde_json::json!([0, 10, ["status", "completedLength", "totalLength"]])).await {
if let Some(arr) = json.get("result").and_then(|r| r.as_array()) {
if arr.iter().any(|item| item.get("status").and_then(|v| v.as_str()) == Some("complete")) {
let _ = app_handle_clone.emit("download-complete", id_clone.clone());
let _ = rpc_call(rpc_port, &rpc_secret, "aria2.shutdown", serde_json::json!([])).await;
break;
}
if arr.iter().any(|item| item.get("status").and_then(|v| v.as_str()) == Some("error")) {
// check if it actually completed length == total length
if let Some(item) = arr.iter().find(|i| i.get("status").and_then(|v| v.as_str()) == Some("error")) {
let comp = item.get("completedLength").and_then(|v| v.as_str()).and_then(|s| s.parse::<f64>().ok()).unwrap_or(0.0);
let tot = item.get("totalLength").and_then(|v| v.as_str()).and_then(|s| s.parse::<f64>().ok()).unwrap_or(1.0);
if comp > 0.0 && comp >= tot {
let _ = app_handle_clone.emit("download-complete", id_clone.clone());
let _ = rpc_call(rpc_port, &rpc_secret, "aria2.shutdown", serde_json::json!([])).await;
break;
}
}
let _ = app_handle_clone.emit("download-failed", id_clone.clone());
let _ = rpc_call(rpc_port, &rpc_secret, "aria2.shutdown", serde_json::json!([])).await;
break;
}
break;
}
}
}
let _ = cmd_tx.send(queue::DownloadCommand::TaskFinished(id_clone_tx));
});
Ok(())
@@ -550,7 +684,6 @@ async fn start_download(
#[tauri::command]
async fn start_media_download(
app_handle: tauri::AppHandle,
state: tauri::State<'_, AppState>,
id: String,
url: String,
@@ -562,7 +695,36 @@ async fn start_media_download(
username: Option<String>,
password: Option<String>,
headers: Option<String>,
proxy: Option<String>,
) -> Result<(), String> {
let task = queue::DownloadTask {
is_media: true,
id, url, destination, filename, speed_limit, username, password, headers, proxy,
connections: None, checksum: None, cookies: None, mirrors: None, user_agent: None, max_tries: None,
format_selector, cookie_source,
};
state.cmd_tx.send(queue::DownloadCommand::Enqueue(task)).map_err(|e| e.to_string())?;
Ok(())
}
pub(crate) async fn start_media_download_internal(
app_handle: tauri::AppHandle,
tasks_map: Arc<Mutex<HashMap<String, u32>>>,
cmd_tx: tokio::sync::mpsc::UnboundedSender<queue::DownloadCommand>,
task: queue::DownloadTask,
) -> Result<(), String> {
let id = task.id;
let url = task.url;
let destination = task.destination;
let filename = task.filename;
let format_selector = task.format_selector;
let cookie_source = task.cookie_source;
let speed_limit = task.speed_limit;
let username = task.username;
let password = task.password;
let headers = task.headers;
let proxy = task.proxy;
println!("start_media_download called for id: {}", id);
let resource_dir = app_handle.path().resource_dir().map_err(|e| e.to_string())?;
let ytdlp_path = resource_dir.join("binaries").join("yt-dlp");
@@ -599,6 +761,9 @@ async fn start_media_download(
.arg("--socket-timeout").arg("20")
.arg("--retries").arg("3")
.arg("--extractor-retries").arg("3")
.arg("--downloader").arg("aria2c")
.arg("--downloader-args").arg("aria2c:-x 16 -s 16 -k 1M")
.arg("--concurrent-fragments").arg("4")
.arg("-o").arg(out_path.to_string_lossy().to_string());
if let Some(limit) = speed_limit {
@@ -607,12 +772,42 @@ async fn start_media_download(
}
}
if let Some(p) = proxy {
if !p.is_empty() {
cmd.arg("--proxy").arg(p);
}
}
if let Some(cs) = cookie_source {
if !cs.is_empty() && cs != "none" {
cmd.arg("--cookies-from-browser").arg(cs);
}
}
let mut config_file = tempfile::Builder::new().prefix("ytdlp-").suffix(".conf").tempfile().map_err(|e| e.to_string())?;
let mut config_content = String::new();
if let Some(user) = username {
if !user.is_empty() {
config_content.push_str(&format!("--username\n{}\n", user));
}
}
if let Some(pass) = password {
if !pass.is_empty() {
config_content.push_str(&format!("--password\n{}\n", pass));
}
}
if let Some(headers) = headers {
for header in headers.lines().map(str::trim).filter(|header| !header.is_empty()) {
config_content.push_str(&format!("--add-header\n{}\n", header));
}
}
use std::io::Write;
config_file.write_all(config_content.as_bytes()).map_err(|e| e.to_string())?;
let config_path = config_file.into_temp_path();
if !config_content.is_empty() {
cmd.arg("--config-locations").arg(&config_path);
}
if let Some(format) = format_selector {
cmd.arg("-f").arg(format);
// If the filename implies an audio format, use it as audio output
@@ -634,33 +829,18 @@ async fn start_media_download(
}
}
if let Some(user) = username {
if !user.is_empty() {
cmd.arg("--username").arg(&user);
}
}
if let Some(pass) = password {
if !pass.is_empty() {
cmd.arg("--password").arg(&pass);
}
}
if let Some(headers) = headers {
for header in headers.lines().map(str::trim).filter(|header| !header.is_empty()) {
cmd.arg("--add-headers").arg(header);
}
}
cmd.arg(&url);
cmd.stdout(Stdio::piped());
cmd.stderr(Stdio::piped()); // Also pipe stderr for better error reporting
let mut child = cmd.spawn().map_err(|e| format!("Failed to spawn yt-dlp: {}", e))?;
let pid = child.id().unwrap_or(0);
state.tasks.lock().unwrap().insert(id.clone(), pid);
tasks_map.lock().unwrap().insert(id.clone(), pid);
let stdout = child.stdout.take().unwrap();
let app_handle_clone = app_handle.clone();
let id_clone = id.clone();
let id_clone_tx = id.clone();
// yt-dlp parsing regex
let pct_re = Regex::new(r"\[download\]\s+(\d+(?:\.\d+)?)%").unwrap();
@@ -668,6 +848,7 @@ async fn start_media_download(
let eta_re = Regex::new(r"ETA\s+([^\s]+)").unwrap();
tokio::spawn(async move {
let _keep_alive = config_path; // Keep the temp file alive
let mut reader = BufReader::new(stdout).lines();
let mut current_track: f64 = 0.0;
let mut last_fraction: f64 = 0.0;
@@ -712,10 +893,12 @@ async fn start_media_download(
}
}
status = child.wait() => {
println!("child exit status: {:?}", status);
if let Ok(exit_status) = status {
if exit_status.success() {
let _ = app_handle_clone.emit("download-complete", id_clone.clone());
} else {
// If it exited with error, emit failed (7 means paused/aborted usually, but we emit failed anyway, UI can filter)
let _ = app_handle_clone.emit("download-failed", id_clone.clone());
}
}
@@ -723,6 +906,7 @@ async fn start_media_download(
}
}
}
let _ = cmd_tx.send(queue::DownloadCommand::TaskFinished(id_clone_tx));
});
Ok(())
@@ -733,21 +917,39 @@ async fn pause_download(state: tauri::State<'_, AppState>, id: String) -> Result
println!("pause_download called for id: {}", id);
if let Some(pid) = state.tasks.lock().unwrap().remove(&id) {
if pid > 0 {
#[cfg(unix)]
{
let _ = std::process::Command::new("kill")
.arg("-15") // SIGTERM
.arg(pid.to_string())
.status();
println!("Sent SIGTERM to pid: {}", pid);
use sysinfo::{System, Pid};
let mut sys = System::new_all();
sys.refresh_processes(sysinfo::ProcessesToUpdate::All, true);
let parent_pid = Pid::from_u32(pid);
let mut to_kill = vec![parent_pid];
let mut idx = 0;
while idx < to_kill.len() {
let current = to_kill[idx];
for (proc_pid, proc) in sys.processes() {
if let Some(p) = proc.parent() {
if p == current {
to_kill.push(*proc_pid);
}
}
}
idx += 1;
}
#[cfg(windows)]
{
let _ = std::process::Command::new("taskkill")
.arg("/PID")
.arg(pid.to_string())
.status();
println!("Sent taskkill to pid: {}", pid);
for p in to_kill.into_iter().rev() {
if let Some(proc) = sys.process(p) {
#[cfg(unix)]
{
if proc.kill_with(sysinfo::Signal::Term).is_none() {
proc.kill(); // Fallback to SIGKILL
}
}
#[cfg(windows)]
proc.kill();
println!("Sent termination signal to pid: {}", p.as_u32());
}
}
}
}
@@ -1099,6 +1301,7 @@ pub fn run() {
let extension_frontend_ready = Arc::new(AtomicBool::new(false));
let server_frontend_ready = extension_frontend_ready.clone();
let (cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
tauri::Builder::default()
.plugin(tauri_plugin_single_instance::init(|app, _args, _cwd| {
if let Some(window) = app.get_webview_window("main") {
@@ -1108,11 +1311,13 @@ pub fn run() {
}))
.plugin(tauri_plugin_deep_link::init())
.manage(AppState {
tasks: Mutex::new(HashMap::new()),
tasks: Arc::new(Mutex::new(HashMap::new())),
cmd_tx,
extension_pairing_token,
extension_frontend_ready,
})
.setup(move |app| {
queue::setup_queue(app.handle().clone(), cmd_rx);
match extension_server::start_server(
app.handle().clone(),
server_pairing_token.clone(),
+70
View File
@@ -0,0 +1,70 @@
use tokio::sync::mpsc;
use tauri::Manager;
#[derive(Clone)]
pub struct DownloadTask {
pub is_media: bool,
pub id: String,
pub url: String,
pub destination: String,
pub filename: String,
pub connections: Option<i32>,
pub speed_limit: Option<String>,
pub username: Option<String>,
pub password: Option<String>,
pub headers: Option<String>,
pub checksum: Option<String>,
pub cookies: Option<String>,
pub mirrors: Option<String>,
pub user_agent: Option<String>,
pub max_tries: Option<i32>,
pub proxy: Option<String>,
pub format_selector: Option<String>,
pub cookie_source: Option<String>,
}
pub enum DownloadCommand {
Enqueue(DownloadTask),
SetLimit(usize),
TaskFinished(String), // id
}
pub fn setup_queue(app_handle: tauri::AppHandle, mut cmd_rx: mpsc::UnboundedReceiver<DownloadCommand>) {
tokio::spawn(async move {
let mut queue = std::collections::VecDeque::new();
let mut active_count = 0;
let mut limit = 3;
while let Some(cmd) = cmd_rx.recv().await {
match cmd {
DownloadCommand::Enqueue(task) => {
queue.push_back(task);
}
DownloadCommand::SetLimit(l) => {
limit = l;
}
DownloadCommand::TaskFinished(_) => {
if active_count > 0 {
active_count -= 1;
}
}
}
while active_count < limit && !queue.is_empty() {
if let Some(task) = queue.pop_front() {
active_count += 1;
let ah = app_handle.clone();
let state = ah.state::<crate::AppState>();
let tx = state.cmd_tx.clone();
let tasks_map = state.tasks.clone();
if task.is_media {
let _ = crate::start_media_download_internal(ah, tasks_map, tx, task).await;
} else {
let _ = crate::start_download_internal(ah, tasks_map, tx, task).await;
}
}
}
}
});
}