diff --git a/apps/desktop/src-tauri/Cargo.lock b/apps/desktop/src-tauri/Cargo.lock index a080d94..3c1e192 100644 --- a/apps/desktop/src-tauri/Cargo.lock +++ b/apps/desktop/src-tauri/Cargo.lock @@ -3909,6 +3909,7 @@ dependencies = [ "tauri-plugin-notification", "tauri-plugin-opener", "tauri-plugin-single-instance", + "tempfile", "tokio", ] diff --git a/apps/desktop/src-tauri/Cargo.toml b/apps/desktop/src-tauri/Cargo.toml index d232b48..b6b3cdd 100644 --- a/apps/desktop/src-tauri/Cargo.toml +++ b/apps/desktop/src-tauri/Cargo.toml @@ -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" diff --git a/apps/desktop/src-tauri/src/lib.rs b/apps/desktop/src-tauri/src/lib.rs index 76ff44d..ee54234 100644 --- a/apps/desktop/src-tauri/src/lib.rs +++ b/apps/desktop/src-tauri/src/lib.rs @@ -29,6 +29,26 @@ async fn fetch_metadata(url: String, user_agent: Option, 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>, - extension_pairing_token: extension_server::SharedExtensionToken, - extension_frontend_ready: extension_server::SharedFrontendReady, +mod queue; + +pub struct AppState { + pub tasks: Arc>>, + pub cmd_tx: tokio::sync::mpsc::UnboundedSender, + 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 { uris } +async fn rpc_call(port: u16, secret: &str, method: &str, params: serde_json::Value) -> Result { + 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, proxy: Option, ) -> 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>>, + cmd_tx: tokio::sync::mpsc::UnboundedSender, + 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::().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::().ok()).unwrap_or(0.0); + let total = item.get("totalLength").and_then(|v| v.as_str()).and_then(|s| s.parse::().ok()).unwrap_or(1.0); + let speed_bytes = item.get("downloadSpeed").and_then(|v| v.as_str()).and_then(|s| s.parse::().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::().ok()).unwrap_or(0.0); + let tot = item.get("totalLength").and_then(|v| v.as_str()).and_then(|s| s.parse::().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, password: Option, headers: Option, + proxy: Option, ) -> 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>>, + cmd_tx: tokio::sync::mpsc::UnboundedSender, + 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(), diff --git a/apps/desktop/src-tauri/src/queue.rs b/apps/desktop/src-tauri/src/queue.rs new file mode 100644 index 0000000..c98a5fb --- /dev/null +++ b/apps/desktop/src-tauri/src/queue.rs @@ -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, + pub speed_limit: Option, + pub username: Option, + pub password: Option, + pub headers: Option, + pub checksum: Option, + pub cookies: Option, + pub mirrors: Option, + pub user_agent: Option, + pub max_tries: Option, + pub proxy: Option, + pub format_selector: Option, + pub cookie_source: Option, +} + +pub enum DownloadCommand { + Enqueue(DownloadTask), + SetLimit(usize), + TaskFinished(String), // id +} + +pub fn setup_queue(app_handle: tauri::AppHandle, mut cmd_rx: mpsc::UnboundedReceiver) { + 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::(); + 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; + } + } + } + } + }); +}