diff --git a/src-tauri/src/commands.rs b/src-tauri/src/commands.rs index 3bfe312..853615e 100644 --- a/src-tauri/src/commands.rs +++ b/src-tauri/src/commands.rs @@ -168,24 +168,20 @@ mod tests { assert!( authorize_exact_path(Path::new("owned.bin"), std::slice::from_ref(&owned)).is_err() ); - assert!( - authorize_exact_path( - &root.path().join("sub/../owned.bin"), - std::slice::from_ref(&owned) - ) - .is_err() - ); + assert!(authorize_exact_path( + &root.path().join("sub/../owned.bin"), + std::slice::from_ref(&owned) + ) + .is_err()); assert!( authorize_exact_path(Path::new("/etc/hosts"), std::slice::from_ref(&owned)).is_err() ); if let Some(home) = std::env::var_os("HOME") { - assert!( - authorize_exact_path( - &PathBuf::from(home).join(".ssh"), - std::slice::from_ref(&owned) - ) - .is_err() - ); + assert!(authorize_exact_path( + &PathBuf::from(home).join(".ssh"), + std::slice::from_ref(&owned) + ) + .is_err()); } } diff --git a/src-tauri/src/db.rs b/src-tauri/src/db.rs index 85b7708..1016bf0 100644 --- a/src-tauri/src/db.rs +++ b/src-tauri/src/db.rs @@ -80,11 +80,7 @@ fn init_at_path_internal(app_data_dir: &Path) -> Result { // We no longer touch the keychain on backend startup. // Legacy imports will safely preserve any pairing token in the JSON payload. // The frontend will manually trigger migration to the keychain via IPC if access is granted. - import_legacy_data( - &mut connection, - app_data_dir, - false, - )?; + import_legacy_data(&mut connection, app_data_dir, false)?; Ok(DbState { conn: Mutex::new(connection), @@ -357,41 +353,50 @@ pub fn sanitize_current_settings_and_restore_token( let Some(settings) = load_settings(connection)? else { return Ok((false, false)); }; - let (sanitized, legacy_token, keychain_granted) = sanitize_settings_text(&settings, force_migrate)?; + let (sanitized, legacy_token, keychain_granted) = + sanitize_settings_text(&settings, force_migrate)?; if sanitized == settings { return Ok((false, keychain_granted)); } let should_migrate = force_migrate || keychain_granted; - if should_migrate - && get_keychain_password(PAIRING_TOKEN_KEYCHAIN_ID).is_err() { - if let Some(token) = legacy_token.filter(|token| !token.trim().is_empty()) { - if let Err(error) = set_keychain_password(PAIRING_TOKEN_KEYCHAIN_ID, &token) { - log::warn!( + if should_migrate && get_keychain_password(PAIRING_TOKEN_KEYCHAIN_ID).is_err() { + if let Some(token) = legacy_token.filter(|token| !token.trim().is_empty()) { + if let Err(error) = set_keychain_password(PAIRING_TOKEN_KEYCHAIN_ID, &token) { + log::warn!( "Persisted pairing token could not be migrated yet; original settings retained: {}", error ); - return Ok((true, keychain_granted)); - } + return Ok((true, keychain_granted)); } } + } save_settings(connection, &sanitized)?; Ok((false, keychain_granted)) } -fn sanitize_settings_value(value: &Value, force_migrate: bool) -> Result<(String, Option, bool), String> { +fn sanitize_settings_value( + value: &Value, + force_migrate: bool, +) -> Result<(String, Option, bool), String> { match value { Value::String(text) => sanitize_settings_text(text, force_migrate), _ => sanitize_settings_document(value.clone(), force_migrate), } } -fn sanitize_settings_text(text: &str, force_migrate: bool) -> Result<(String, Option, bool), String> { +fn sanitize_settings_text( + text: &str, + force_migrate: bool, +) -> Result<(String, Option, bool), String> { let document: Value = serde_json::from_str(text) .map_err(|error| format!("failed to decode persisted settings: {error}"))?; sanitize_settings_document(document, force_migrate) } -fn sanitize_settings_document(mut document: Value, force_migrate: bool) -> Result<(String, Option, bool), String> { +fn sanitize_settings_document( + mut document: Value, + force_migrate: bool, +) -> Result<(String, Option, bool), String> { let state_value = if document.get("state").is_some() { document .get_mut("state") @@ -402,14 +407,14 @@ fn sanitize_settings_document(mut document: Value, force_migrate: bool) -> Resul let state = state_value .as_object_mut() .ok_or_else(|| "persisted settings state must be an object".to_string())?; - + let keychain_granted = state .get("keychainAccessGranted") .and_then(|v| v.as_bool()) .unwrap_or(false); - + let should_migrate = force_migrate || keychain_granted; - + let token = if should_migrate { state .remove("extensionPairingToken") @@ -417,7 +422,7 @@ fn sanitize_settings_document(mut document: Value, force_migrate: bool) -> Resul } else { None }; - + let serialized = serde_json::to_string(&document) .map_err(|error| format!("failed to encode persisted settings: {error}"))?; Ok((serialized, token, keychain_granted)) @@ -755,7 +760,7 @@ pub fn hydrate_pairing_token( if skip_keychain { return Ok((generate_pairing_token(), false)); } - + let existing = get_keychain_password(PAIRING_TOKEN_KEYCHAIN_ID).ok(); let generated = generate_pairing_token(); let decision = decide_pairing_token( diff --git a/src-tauri/src/download.rs b/src-tauri/src/download.rs index 6155568..bc08f6d 100644 --- a/src-tauri/src/download.rs +++ b/src-tauri/src/download.rs @@ -118,7 +118,11 @@ impl DownloadCoordinator { .map_err(|_| "download coordinator is unavailable".to_string()) } - pub async fn pause_media_with_ack(&self, id: String, ack: tokio::sync::oneshot::Sender<()>) -> Result<(), String> { + pub async fn pause_media_with_ack( + &self, + id: String, + ack: tokio::sync::oneshot::Sender<()>, + ) -> Result<(), String> { self.media_tx .send(MediaCmd::PauseWithAck(id, ack)) .await @@ -474,14 +478,27 @@ async fn download_file( let mut attempts = 0_usize; loop { attempts += 1; - match download_attempt(&events, &client, &default_headers, &payload, url, &mut control_rx).await { + match download_attempt( + &events, + &client, + &default_headers, + &payload, + url, + &mut control_rx, + ) + .await + { Ok(()) => return DownloadOutcome::Completed, Err(AttemptError::Controlled(DownloadControl::Pause)) => { return DownloadOutcome::Paused; } Err(AttemptError::Controlled(DownloadControl::Cancel)) => { if let Err(e) = fs::remove_file(&payload.output_path).await { - log::warn!("Failed to remove cancelled file '{}': {}", payload.output_path.display(), e); + log::warn!( + "Failed to remove cancelled file '{}': {}", + payload.output_path.display(), + e + ); } return DownloadOutcome::Cancelled; } @@ -591,7 +608,9 @@ async fn download_attempt( .get(reqwest::header::CONTENT_RANGE) .and_then(|h| h.to_str().ok()); if !content_range.is_some_and(|r| r.starts_with(&format!("bytes {}-", existing_len))) { - return Err(AttemptError::Failed("Server returned invalid Content-Range for resume".to_string())); + return Err(AttemptError::Failed( + "Server returned invalid Content-Range for resume".to_string(), + )); } } let completed_at_start = if resumed { existing_len } else { 0 }; @@ -721,11 +740,17 @@ fn build_client(payload: &DownloadPayload) -> Result<(Client, HeaderMap), String if proxy == "none" { builder = builder.no_proxy(); } else { - builder = builder.proxy(reqwest::Proxy::all(proxy).map_err(|_| "Invalid proxy URL configured".to_string())?); + builder = builder.proxy( + reqwest::Proxy::all(proxy) + .map_err(|_| "Invalid proxy URL configured".to_string())?, + ); } } - builder.build().map_err(|error| error.to_string()).map(|c| (c, headers)) + builder + .build() + .map_err(|error| error.to_string()) + .map(|c| (c, headers)) } pub(crate) fn format_speed(bytes_per_second: f64) -> String { @@ -808,14 +833,16 @@ mod tests { let (coordinator, mut events) = DownloadCoordinator::spawn_headless(); coordinator .send(DownloadCmd::CaptureUrls(vec![ - "https://example.com/startup.zip".to_string() + "https://example.com/startup.zip".to_string(), ])) .await .unwrap(); - assert!(tokio::time::timeout(Duration::from_millis(20), events.recv()) - .await - .is_err()); + assert!( + tokio::time::timeout(Duration::from_millis(20), events.recv()) + .await + .is_err() + ); coordinator .send(DownloadCmd::FrontendReady(true)) @@ -827,9 +854,7 @@ mod tests { .await .unwrap() .unwrap(), - DownloadEvent::CapturedUrls( - "https://example.com/startup.zip".to_string() - ) + DownloadEvent::CapturedUrls("https://example.com/startup.zip".to_string()) ); } } diff --git a/src-tauri/src/download_ownership.rs b/src-tauri/src/download_ownership.rs index e080e3c..aa9ba61 100644 --- a/src-tauri/src/download_ownership.rs +++ b/src-tauri/src/download_ownership.rs @@ -18,7 +18,12 @@ pub fn canonical_download_filename(filename: &str) -> String { let sanitized = leaf .chars() .map(|character| { - if character.is_control() || matches!(character, '<' | '>' | ':' | '"' | '/' | '\\' | '|' | '?' | '*') { + if character.is_control() + || matches!( + character, + '<' | '>' | ':' | '"' | '/' | '\\' | '|' | '?' | '*' + ) + { '-' } else { character @@ -30,7 +35,10 @@ pub fn canonical_download_filename(filename: &str) -> String { "download".to_string() } else if crate::platform::is_windows_reserved_filename(sanitized) { let path = Path::new(sanitized); - let stem = path.file_stem().and_then(|value| value.to_str()).unwrap_or("download"); + let stem = path + .file_stem() + .and_then(|value| value.to_str()) + .unwrap_or("download"); match path.extension().and_then(|value| value.to_str()) { Some(extension) => format!("{stem}-.{extension}"), None => format!("{stem}-"), @@ -135,7 +143,7 @@ fn load_records(app_handle: &tauri::AppHandle) -> Result Result, String> { let settings = crate::settings::load_settings(app_handle).ok(); - + let database = app_handle.state::(); let connection = database.lock()?; let downloads = crate::db::load_downloads(&connection)? @@ -144,7 +152,6 @@ fn legacy_download_queue_paths(app_handle: &tauri::AppHandle) -> Result, _>>() .map_err(|error| format!("Invalid download queue ownership data: {error}"))?; - let mut paths = Vec::new(); for download in downloads { let category = format!("{:?}", download.category); @@ -214,7 +221,10 @@ mod tests { #[test] fn canonicalizes_untrusted_download_filenames() { - assert_eq!(canonical_download_filename("../folder/video?.mp4"), "video-.mp4"); + assert_eq!( + canonical_download_filename("../folder/video?.mp4"), + "video-.mp4" + ); assert_eq!(canonical_download_filename(" report. "), "report"); assert_eq!(canonical_download_filename(".."), "download"); assert_eq!(canonical_download_filename("CON.txt"), "CON-.txt"); diff --git a/src-tauri/src/engines.rs b/src-tauri/src/engines.rs index 9e98f83..99a7b02 100644 --- a/src-tauri/src/engines.rs +++ b/src-tauri/src/engines.rs @@ -72,7 +72,10 @@ fn executable_relative_candidates( .join("engine-dist") .join(target) .join(binary_name), - executable_dir.join("engines").join(target).join(binary_name), + executable_dir + .join("engines") + .join(target) + .join(binary_name), ]; if cfg!(target_os = "macos") { @@ -99,11 +102,7 @@ fn development_candidates(cwd: &Path, target: &str, binary_name: &str) -> Vec, next: Next) -> Response { let mut response = next.run(request).await; - response - .headers_mut() - .insert(SERVER_HEADER, HeaderValue::from_static("1")); - response - .headers_mut() - .insert(PROTOCOL_VERSION_HEADER, HeaderValue::from_static(PROTOCOL_VERSION)); - response + response + .headers_mut() + .insert(SERVER_HEADER, HeaderValue::from_static("1")); + response.headers_mut().insert( + PROTOCOL_VERSION_HEADER, + HeaderValue::from_static(PROTOCOL_VERSION), + ); + response } async fn bind_extension_listener() -> Result<(u16, tokio::net::TcpListener), String> { @@ -382,7 +383,7 @@ fn is_allowed_origin(origin: &str) -> bool { #[cfg(test)] mod tests { - use super::{add_server_identity, PROTOCOL_VERSION_HEADER, SERVER_HEADER}; + use super::{add_server_identity, PROTOCOL_VERSION_HEADER, SERVER_HEADER}; use axum::{http::StatusCode, middleware, routing::get, Router}; #[tokio::test] @@ -398,13 +399,15 @@ mod tests { axum::serve(listener, app).await.unwrap(); }); - let response = reqwest::get(format!("http://{address}/ping")).await.unwrap(); - assert_eq!(response.status(), StatusCode::FORBIDDEN); - assert_eq!(response.headers().get(SERVER_HEADER).unwrap(), "1"); - assert_eq!( - response.headers().get(PROTOCOL_VERSION_HEADER).unwrap(), - "2" - ); + let response = reqwest::get(format!("http://{address}/ping")) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::FORBIDDEN); + assert_eq!(response.headers().get(SERVER_HEADER).unwrap(), "1"); + assert_eq!( + response.headers().get(PROTOCOL_VERSION_HEADER).unwrap(), + "2" + ); server.abort(); } diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index e886ddb..4f8dcd1 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -1,17 +1,17 @@ #![allow(unexpected_cfgs)] // Learn more about Tauri commands at https://tauri.app/develop/calling-rust/ -use tauri::{Manager, Emitter}; use regex::Regex; use serde::Serialize; -use ts_rs::TS; -use uuid::Uuid; -use tauri_plugin_deep_link::DeepLinkExt; use std::collections::HashMap; use std::hash::{Hash, Hasher}; use std::path::PathBuf; use std::sync::OnceLock; use std::time::{Duration, Instant}; +use tauri::{Emitter, Manager}; +use tauri_plugin_deep_link::DeepLinkExt; +use ts_rs::TS; +use uuid::Uuid; #[derive(Serialize, TS)] #[ts(export, export_to = "../../src/bindings/")] @@ -60,7 +60,11 @@ fn is_media_processing_line(line: &str) -> bool { } fn json_str<'a>(value: &'a serde_json::Value, key: &str) -> Option<&'a str> { - value.get(key).and_then(|v| v.as_str()).map(str::trim).filter(|v| !v.is_empty()) + value + .get(key) + .and_then(|v| v.as_str()) + .map(str::trim) + .filter(|v| !v.is_empty()) } fn json_lower(value: &serde_json::Value, key: &str) -> String { @@ -68,11 +72,16 @@ fn json_lower(value: &serde_json::Value, key: &str) -> String { } fn json_u64(value: &serde_json::Value, key: &str) -> Option { - value.get(key).and_then(|v| v.as_u64().or_else(|| v.as_f64().map(|f| f as u64))) + value + .get(key) + .and_then(|v| v.as_u64().or_else(|| v.as_f64().map(|f| f as u64))) } fn json_f64(value: &serde_json::Value, key: &str) -> Option { - value.get(key).and_then(|v| v.as_f64().or_else(|| v.as_str().and_then(|s| s.parse::().ok()))) + value.get(key).and_then(|v| { + v.as_f64() + .or_else(|| v.as_str().and_then(|s| s.parse::().ok())) + }) } fn media_filesize(value: &serde_json::Value) -> Option { @@ -84,7 +93,10 @@ fn media_exact_filesize(value: &serde_json::Value) -> Option { } fn media_approx_filesize(value: &serde_json::Value) -> Option { - media_exact_filesize(value).is_none().then(|| json_u64(value, "filesize_approx")).flatten() + media_exact_filesize(value) + .is_none() + .then(|| json_u64(value, "filesize_approx")) + .flatten() } fn media_sized_bytes(value: &serde_json::Value) -> Option<(u64, bool)> { @@ -126,7 +138,11 @@ fn split_size_estimate(bytes: Option<(u64, bool)>) -> (Option, Option) } fn codec_is_present(codec: Option<&str>) -> bool { - codec.map(str::trim).filter(|codec| !codec.is_empty()).map(|codec| codec.to_lowercase() != "none").unwrap_or(false) + codec + .map(str::trim) + .filter(|codec| !codec.is_empty()) + .map(|codec| codec.to_lowercase() != "none") + .unwrap_or(false) } fn has_video_stream(value: &serde_json::Value) -> bool { @@ -146,7 +162,11 @@ fn is_excluded_yt_dlp_format(value: &serde_json::Value) -> bool { for key in ["format_note", "format", "format_id", "protocol"] { let text = json_lower(value, key); - if text.contains("storyboard") || text.contains("thumbnail") || text.contains("subtitle") || text.contains("subtitles") { + if text.contains("storyboard") + || text.contains("thumbnail") + || text.contains("subtitle") + || text.contains("subtitles") + { return true; } } @@ -197,19 +217,30 @@ fn format_score(value: &serde_json::Value) -> u64 { height_score + bitrate_score + size_score } -fn best_matching_format<'a, F>(formats: &'a [&'a serde_json::Value], predicate: F) -> Option<&'a serde_json::Value> +fn best_matching_format<'a, F>( + formats: &'a [&'a serde_json::Value], + predicate: F, +) -> Option<&'a serde_json::Value> where F: Fn(&serde_json::Value) -> bool, { - formats.iter().copied().filter(|format| predicate(format)).max_by_key(|format| format_score(format)) + formats + .iter() + .copied() + .filter(|format| predicate(format)) + .max_by_key(|format| format_score(format)) } -fn best_audio_format<'a>(formats: &'a [&'a serde_json::Value], ext: Option<&str>) -> Option<&'a serde_json::Value> { +fn best_audio_format<'a>( + formats: &'a [&'a serde_json::Value], + ext: Option<&str>, +) -> Option<&'a serde_json::Value> { best_matching_format(formats, |format| { if !has_audio_stream(format) || has_video_stream(format) { return false; } - ext.map(|wanted| json_lower(format, "ext") == wanted).unwrap_or(true) + ext.map(|wanted| json_lower(format, "ext") == wanted) + .unwrap_or(true) }) } @@ -231,7 +262,11 @@ fn display_codec(codec: Option<&str>, fallback: &str) -> String { "VP9".to_string() } else if lower.starts_with("vp8") { "VP8".to_string() - } else if lower.starts_with("hev1") || lower.starts_with("hvc1") || lower.contains("h265") || lower.contains("hevc") { + } else if lower.starts_with("hev1") + || lower.starts_with("hvc1") + || lower.contains("h265") + || lower.contains("hevc") + { "H.265".to_string() } else if lower.starts_with("mp4a") || lower.contains("aac") { "AAC".to_string() @@ -246,7 +281,11 @@ fn display_codec(codec: Option<&str>, fallback: &str) -> String { } } -fn joined_format_label(container: &str, video_codec: Option<&str>, audio_codec: Option<&str>) -> String { +fn joined_format_label( + container: &str, + video_codec: Option<&str>, + audio_codec: Option<&str>, +) -> String { let mut codecs = Vec::new(); if video_codec.is_some() { codecs.push(display_codec(video_codec, "Video")); @@ -320,7 +359,10 @@ fn raw_media_format( split_size_estimate(estimated_stream_bytes(value, duration_seconds)); let resolution = if has_video_stream(value) { - format_height(value).map(|height| format!("{height}p")).or_else(|| json_str(value, "resolution").map(ToOwned::to_owned)).unwrap_or_else(|| "Video".to_string()) + format_height(value) + .map(|height| format!("{height}p")) + .or_else(|| json_str(value, "resolution").map(ToOwned::to_owned)) + .unwrap_or_else(|| "Video".to_string()) } else { "Audio only".to_string() }; @@ -333,14 +375,25 @@ fn raw_media_format( joined_format_label(&ext, None, json_str(value, "acodec")) }; - Some(MediaFormat { format_id, resolution, ext, format_label, fps, filesize, filesize_approx }) + Some(MediaFormat { + format_id, + resolution, + ext, + format_label, + fps, + filesize, + filesize_approx, + }) } fn build_media_format_options( formats_arr: &[serde_json::Value], duration_seconds: Option, ) -> Vec { - let clean_formats: Vec<&serde_json::Value> = formats_arr.iter().filter(|format| !is_excluded_yt_dlp_format(format)).collect(); + let clean_formats: Vec<&serde_json::Value> = formats_arr + .iter() + .filter(|format| !is_excluded_yt_dlp_format(format)) + .collect(); let has_video = clean_formats.iter().any(|format| has_video_stream(format)); let has_audio = clean_formats.iter().any(|format| has_audio_stream(format)); let mut options = Vec::new(); @@ -348,31 +401,51 @@ fn build_media_format_options( if has_video { let available_heights: Vec = [2160_u64, 1440, 1080, 720, 480, 360] .into_iter() - .filter(|height| clean_formats.iter().any(|format| matches_media_height(format, *height))) + .filter(|height| { + clean_formats + .iter() + .any(|format| matches_media_height(format, *height)) + }) .collect(); for height in available_heights { - if let Some(video) = best_matching_format(&clean_formats, |format| has_video_stream(format) && matches_media_height(format, height)) { - let audio = if has_audio_stream(video) { None } else { best_audio_format(&clean_formats, None) }; + if let Some(video) = best_matching_format(&clean_formats, |format| { + has_video_stream(format) && matches_media_height(format, height) + }) { + let audio = if has_audio_stream(video) { + None + } else { + best_audio_format(&clean_formats, None) + }; let (filesize, filesize_approx) = estimated_merged_size(Some(video), audio, duration_seconds); options.push(MediaFormat { - format_id: selected_format_id(Some(video), audio) - .unwrap_or_else(|| format!("bestvideo[height<={height}]+bestaudio/best[height<={height}]")), + format_id: selected_format_id(Some(video), audio).unwrap_or_else(|| { + format!("bestvideo[height<={height}]+bestaudio/best[height<={height}]") + }), resolution: format!("{height}p"), ext: "mkv".to_string(), - format_label: joined_format_label("mkv", json_str(video, "vcodec"), audio.and_then(|format| json_str(format, "acodec"))), + format_label: joined_format_label( + "mkv", + json_str(video, "vcodec"), + audio.and_then(|format| json_str(format, "acodec")), + ), fps: json_f64(video, "fps"), filesize, filesize_approx, }); } - if let Some(video) = best_matching_format(&clean_formats, |format| has_video_stream(format) && matches_media_height(format, height) && json_lower(format, "ext") == "mp4") { + if let Some(video) = best_matching_format(&clean_formats, |format| { + has_video_stream(format) + && matches_media_height(format, height) + && json_lower(format, "ext") == "mp4" + }) { let audio = if has_audio_stream(video) { None } else { - best_audio_format(&clean_formats, Some("m4a")).or_else(|| best_audio_format(&clean_formats, None)) + best_audio_format(&clean_formats, Some("m4a")) + .or_else(|| best_audio_format(&clean_formats, None)) }; let (filesize, filesize_approx) = estimated_merged_size(Some(video), audio, duration_seconds); @@ -389,11 +462,17 @@ fn build_media_format_options( }); } - if let Some(video) = best_matching_format(&clean_formats, |format| has_video_stream(format) && matches_media_height(format, height) && json_lower(format, "ext") == "webm") { + if let Some(video) = best_matching_format(&clean_formats, |format| { + has_video_stream(format) + && matches_media_height(format, height) + && json_lower(format, "ext") == "webm" + }) { let audio = if has_audio_stream(video) { None } else { - best_audio_format(&clean_formats, Some("webm")).or_else(|| best_audio_format(&clean_formats, Some("opus"))).or_else(|| best_audio_format(&clean_formats, None)) + best_audio_format(&clean_formats, Some("webm")) + .or_else(|| best_audio_format(&clean_formats, Some("opus"))) + .or_else(|| best_audio_format(&clean_formats, None)) }; let (filesize, filesize_approx) = estimated_merged_size(Some(video), audio, duration_seconds); @@ -429,7 +508,9 @@ fn build_media_format_options( }); } - if let Some(audio) = best_audio_format(&clean_formats, Some("webm")).or_else(|| best_audio_format(&clean_formats, Some("opus"))) { + if let Some(audio) = best_audio_format(&clean_formats, Some("webm")) + .or_else(|| best_audio_format(&clean_formats, Some("opus"))) + { let (filesize, filesize_approx) = split_size_estimate(estimated_stream_bytes(audio, duration_seconds)); options.push(MediaFormat { @@ -449,7 +530,9 @@ fn build_media_format_options( let (filesize, filesize_approx) = split_size_estimate(estimated_stream_bytes(audio, duration_seconds)); options.push(MediaFormat { - format_id: json_str(audio, "format_id").unwrap_or("bestaudio/best").to_string(), + format_id: json_str(audio, "format_id") + .unwrap_or("bestaudio/best") + .to_string(), resolution: "Audio only".to_string(), ext: "mp3".to_string(), format_label: "MP3 • Best audio".to_string(), @@ -522,13 +605,7 @@ fn parse_media_progress_line(line: &str) -> Option { downloaded / total } else { progress_json_string(&progress, "_percent_str") - .and_then(|percent| { - percent - .trim_end_matches('%') - .trim() - .parse::() - .ok() - }) + .and_then(|percent| percent.trim_end_matches('%').trim().parse::().ok()) .unwrap_or(0.0) / 100.0 }; @@ -547,9 +624,7 @@ fn parse_media_progress_line(line: &str) -> Option { .unwrap_or_else(|| "-".to_string()); let size = progress_json_string(&progress, "_total_bytes_str") .or_else(|| progress_json_string(&progress, "_total_bytes_estimate_str")) - .or_else(|| { - (total > 0.0).then(|| crate::download::format_size(total)) - }); + .or_else(|| (total > 0.0).then(|| crate::download::format_size(total))); return Some(MediaProgress { fraction: fraction.clamp(0.0, 1.0), @@ -591,14 +666,12 @@ fn parse_media_progress_line(line: &str) -> Option { }); } - let percent_re = YTDLP_PCT_RE - .get_or_init(|| Regex::new(r"\[download\]\s+~?\s*(\d+(?:\.\d+)?)%").unwrap()); + let percent_re = + YTDLP_PCT_RE.get_or_init(|| Regex::new(r"\[download\]\s+~?\s*(\d+(?:\.\d+)?)%").unwrap()); let captures = percent_re.captures(line)?; let fraction = captures.get(1)?.as_str().parse::().ok()? / 100.0; - let speed_re = - YTDLP_SPD_RE.get_or_init(|| Regex::new(r"\bat\s+([^\s]+)").unwrap()); - let eta_re = - YTDLP_ETA_RE.get_or_init(|| Regex::new(r"\bETA\s+([^\s]+)").unwrap()); + let speed_re = YTDLP_SPD_RE.get_or_init(|| Regex::new(r"\bat\s+([^\s]+)").unwrap()); + let eta_re = YTDLP_ETA_RE.get_or_init(|| Regex::new(r"\bETA\s+([^\s]+)").unwrap()); Some(MediaProgress { fraction: fraction.clamp(0.0, 1.0), @@ -617,7 +690,11 @@ fn parse_media_progress_line(line: &str) -> Option { }) } -fn media_progress_speed(progress: &MediaProgress, now: Instant, last_sample: &mut Option<(Instant, f64)>) -> String { +fn media_progress_speed( + progress: &MediaProgress, + now: Instant, + last_sample: &mut Option<(Instant, f64)>, +) -> String { let Some(downloaded_bytes) = progress.downloaded_bytes else { return progress.speed.clone(); }; @@ -714,9 +791,6 @@ async fn cleanup_media_artifacts(out_path: &std::path::Path, remove_primary: boo } } - - - async fn validate_url_ssrf(url: &str) -> Result, String> { let parsed = reqwest::Url::parse(url).map_err(|_| "SSRF blocked: Invalid URL")?; if parsed.scheme() != "http" && parsed.scheme() != "https" { @@ -724,14 +798,14 @@ async fn validate_url_ssrf(url: &str) -> Result Result { - if (ipv6.segments()[0] & 0xfe00) == 0xfc00 { // ULA check + if (ipv6.segments()[0] & 0xfe00) == 0xfc00 { + // ULA check return Err("SSRF blocked: Private/local IP not allowed".to_string()); } - if (ipv6.segments()[0] & 0xffc0) == 0xfe80 { // Link-local check + if (ipv6.segments()[0] & 0xffc0) == 0xfe80 { + // Link-local check return Err("SSRF blocked: Private/local IP not allowed".to_string()); } } @@ -754,18 +830,23 @@ async fn validate_url_ssrf(url: &str) -> Result, username: Option, password: Option) -> Result { +async fn fetch_metadata( + url: String, + user_agent: Option, + username: Option, + password: Option, +) -> Result { let mut current_url = url.clone(); let mut redirects = 0; let res; - + loop { if redirects >= 5 { return Err("Too many redirects".to_string()); } - + let mut builder = reqwest::Client::builder().redirect(reqwest::redirect::Policy::none()); - + if let Some(ref ua) = user_agent { if !ua.is_empty() { builder = builder.user_agent(ua); @@ -873,15 +954,23 @@ async fn fetch_metadata(url: String, user_agent: Option, username: Optio } } - Ok(MetadataResponse { url: current_url, filename, size: size_str, size_bytes }) + Ok(MetadataResponse { + url: current_url, + filename, + size: size_str, + size_bytes, + }) } const MEDIA_METADATA_CACHE_TTL: Duration = Duration::from_secs(60); const MEDIA_METADATA_TIMEOUT: Duration = Duration::from_secs(55); const MEDIA_METADATA_CACHE_MAX_ENTRIES: usize = 128; -static MEDIA_METADATA_CACHE: OnceLock>> = OnceLock::new(); -static MEDIA_METADATA_LOCKS: OnceLock>>>> = OnceLock::new(); +static MEDIA_METADATA_CACHE: OnceLock>> = + OnceLock::new(); +static MEDIA_METADATA_LOCKS: OnceLock< + tokio::sync::Mutex>>>, +> = OnceLock::new(); fn media_metadata_cache_key( url: &str, @@ -912,7 +1001,9 @@ async fn release_media_metadata_lock( } } -fn resolve_metadata_ytdlp_path(app_handle: &tauri::AppHandle) -> Result<(PathBuf, &'static str), String> { +fn resolve_metadata_ytdlp_path( + app_handle: &tauri::AppHandle, +) -> Result<(PathBuf, &'static str), String> { resolve_bundled_binary_path(app_handle, "yt-dlp") .map(|path| (path, "bundled")) .map_err(|e| format!("failed to find bundled yt-dlp: {e}")) @@ -958,14 +1049,8 @@ async fn fetch_media_metadata( } drop(cache_guard); - let result = fetch_media_metadata_uncached( - app_handle, - url, - cookie_browser, - username, - password, - ) - .await; + let result = + fetch_media_metadata_uncached(app_handle, url, cookie_browser, username, password).await; let result = match result { Ok(metadata) if metadata.formats.is_empty() => { @@ -996,36 +1081,58 @@ async fn fetch_media_metadata( result } -async fn fetch_media_metadata_uncached(app_handle: tauri::AppHandle, url: String, cookie_browser: Option, username: Option, password: Option) -> Result { +async fn fetch_media_metadata_uncached( + app_handle: tauri::AppHandle, + url: String, + cookie_browser: Option, + username: Option, + password: Option, +) -> Result { // Pass bundled tools by absolute path so extraction never depends on // system Python, a user-managed PATH, or auto-detection heuristics. - let deno_path = resolve_bundled_binary_path(&app_handle, "deno").map_err(|e| format!("failed to find bundled deno: {e}"))?; - let ffmpeg_path = resolve_bundled_binary_path(&app_handle, "ffmpeg").map_err(|e| format!("failed to find bundled ffmpeg: {e}"))?; + let deno_path = resolve_bundled_binary_path(&app_handle, "deno") + .map_err(|e| format!("failed to find bundled deno: {e}"))?; + let ffmpeg_path = resolve_bundled_binary_path(&app_handle, "ffmpeg") + .map_err(|e| format!("failed to find bundled ffmpeg: {e}"))?; let deno_runtime = format!("deno:{}", deno_path.to_string_lossy()); let trusted_path = crate::platform::trusted_system_path()?; use tauri_plugin_shell::ShellExt; let (ytdlp_path, _) = resolve_metadata_ytdlp_path(&app_handle)?; - let mut cmd = app_handle.shell().command(ytdlp_path.to_string_lossy().to_string()); - cmd = cmd.env("PATH", trusted_path) - .arg("--ffmpeg-location").arg(&ffmpeg_path) - .arg("--js-runtimes").arg(&deno_runtime) + let mut cmd = app_handle + .shell() + .command(ytdlp_path.to_string_lossy().to_string()); + cmd = cmd + .env("PATH", trusted_path) + .arg("--ffmpeg-location") + .arg(&ffmpeg_path) + .arg("--js-runtimes") + .arg(&deno_runtime) .arg("--no-warnings") .arg("--no-playlist") .arg("--skip-download") - .arg("--socket-timeout").arg("20") - .arg("--retries").arg("3") - .arg("--extractor-retries").arg("3") - .arg("--compat-options").arg("no-youtube-unavailable-videos") - .arg("--print").arg("%(.{title,duration,thumbnail,formats})j"); + .arg("--socket-timeout") + .arg("20") + .arg("--retries") + .arg("3") + .arg("--extractor-retries") + .arg("3") + .arg("--compat-options") + .arg("no-youtube-unavailable-videos") + .arg("--print") + .arg("%(.{title,duration,thumbnail,formats})j"); if let Some(browser) = cookie_browser { if !browser.is_empty() { cmd = 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_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() { @@ -1040,7 +1147,9 @@ async fn fetch_media_metadata_uncached(app_handle: tauri::AppHandle, url: String } } use std::io::Write; - config_file.write_all(config_content.as_bytes()).map_err(|e| e.to_string())?; + 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 = cmd.arg("--config-location").arg(&config_path); @@ -1058,25 +1167,44 @@ async fn fetch_media_metadata_uncached(app_handle: tauri::AppHandle, url: String })? .map_err(|e| format!("Failed to execute yt-dlp: {}", e))?; if output.status.success() { - let value: serde_json::Value = serde_json::from_slice(&output.stdout).map_err(|e| format!("Failed to parse JSON: {}", e))?; - - let title = value.get("title").and_then(|v| v.as_str()).unwrap_or("Unknown Title").to_string(); + let value: serde_json::Value = serde_json::from_slice(&output.stdout) + .map_err(|e| format!("Failed to parse JSON: {}", e))?; + + let title = value + .get("title") + .and_then(|v| v.as_str()) + .unwrap_or("Unknown Title") + .to_string(); let duration_seconds = value.get("duration").and_then(|v| v.as_f64()); let duration = duration_seconds.map(|v| v as u64); - let thumbnail = value.get("thumbnail").and_then(|v| v.as_str()).map(|s| s.to_string()); + let thumbnail = value + .get("thumbnail") + .and_then(|v| v.as_str()) + .map(|s| s.to_string()); let formats = value .get("formats") .and_then(|v| v.as_array()) .map(|formats_arr| build_media_format_options(formats_arr, duration_seconds)) .unwrap_or_default(); - Ok(MediaMetadata { title, duration, thumbnail, formats }) + Ok(MediaMetadata { + title, + duration, + thumbnail, + formats, + }) } else { let err = String::from_utf8_lossy(&output.stderr).trim().to_string(); if err.is_empty() { - Err(format!("yt-dlp failed while fetching media metadata (exit status: {:?})", output.status.code())) + Err(format!( + "yt-dlp failed while fetching media metadata (exit status: {:?})", + output.status.code() + )) } else { - Err(format!("yt-dlp failed while fetching media metadata: {}", err)) + Err(format!( + "yt-dlp failed while fetching media metadata: {}", + err + )) } } } @@ -1096,7 +1224,9 @@ async fn test_ffmpeg(app_handle: tauri::AppHandle) -> Result { use tauri_plugin_shell::ShellExt; let binary_path = resolve_bundled_binary_path(&app_handle, "ffmpeg")?; - let output = app_handle.shell().command(&binary_path) + let output = app_handle + .shell() + .command(&binary_path) .arg("-version") .output() .await @@ -1106,12 +1236,19 @@ async fn test_ffmpeg(app_handle: tauri::AppHandle) -> Result { let text = String::from_utf8_lossy(&output.stdout).trim().to_string(); let first_line = text.lines().next().unwrap_or(""); let re = regex::Regex::new(r"(?i)version\s+([\d\.]+)").unwrap(); - let clean = re.captures(first_line) + let clean = re + .captures(first_line) .and_then(|c| c.get(1)) .map(|m| m.as_str().to_string()) .unwrap_or_else(|| { let parts: Vec<&str> = first_line.split_whitespace().collect(); - parts.get(2).unwrap_or(&first_line).split('-').next().unwrap_or("").to_string() + parts + .get(2) + .unwrap_or(&first_line) + .split('-') + .next() + .unwrap_or("") + .to_string() }); Ok(clean) } else { @@ -1125,7 +1262,9 @@ async fn test_deno(app_handle: tauri::AppHandle) -> Result { use tauri_plugin_shell::ShellExt; let binary_path = resolve_bundled_binary_path(&app_handle, "deno")?; - let output = app_handle.shell().command(&binary_path) + let output = app_handle + .shell() + .command(&binary_path) .arg("--version") .output() .await @@ -1134,7 +1273,12 @@ async fn test_deno(app_handle: tauri::AppHandle) -> Result { if output.status.success() { let text = String::from_utf8_lossy(&output.stdout).trim().to_string(); let re = regex::Regex::new(r"deno\s+(\d+\.\d+\.\d+)").unwrap(); - let clean = re.captures(&text).and_then(|c| c.get(1)).map(|m| m.as_str()).unwrap_or(&text).to_string(); + let clean = re + .captures(&text) + .and_then(|c| c.get(1)) + .map(|m| m.as_str()) + .unwrap_or(&text) + .to_string(); Ok(clean) } else { let err = String::from_utf8_lossy(&output.stderr); @@ -1213,10 +1357,7 @@ fn approved_download_roots(app_handle: &tauri::AppHandle) -> Vec Result { +fn approve_download_root(app_handle: tauri::AppHandle, path: String) -> Result { let resolved = resolve_path(path.trim(), &app_handle); if !resolved.is_absolute() { return Err("Download root must be an absolute path".to_string()); @@ -1235,13 +1376,19 @@ fn approve_download_root( if !roots.is_array() { *roots = serde_json::Value::Array(Vec::new()); } - let values = roots.as_array_mut().expect("approved roots must be an array"); - if !values.iter().filter_map(serde_json::Value::as_str).any(|root| { - crate::platform::paths_equal( - std::path::Path::new(root), - std::path::Path::new(&canonical_text), - ) - }) { + let values = roots + .as_array_mut() + .expect("approved roots must be an array"); + if !values + .iter() + .filter_map(serde_json::Value::as_str) + .any(|root| { + crate::platform::paths_equal( + std::path::Path::new(root), + std::path::Path::new(&canonical_text), + ) + }) + { values.push(serde_json::Value::String(canonical_text.clone())); } })?; @@ -1291,19 +1438,17 @@ impl Drop for Aria2DaemonGuard { } } - - +pub mod commands; pub mod download; -pub mod queue; +pub mod download_ownership; +mod engines; +pub mod error; #[allow(dead_code)] pub mod ipc; mod parity; -pub mod error; -pub mod commands; -pub mod download_ownership; -pub mod retry; -mod engines; mod platform; +pub mod queue; +pub mod retry; mod settings; pub use error::AppError; @@ -1366,21 +1511,25 @@ pub struct EngineStatusResult { pub engines: Vec, } - pub(crate) fn resolve_path(path: &str, app_handle: &tauri::AppHandle) -> std::path::PathBuf { use tauri::Manager; let mut resolved = std::path::PathBuf::from(path); - if let Some(stripped) = path - .strip_prefix("~/") - .or_else(|| path.strip_prefix("~\\")) - { - if let Some(home) = app_handle.path().home_dir().ok().or_else(|| std::env::var("USERPROFILE").ok().map(std::path::PathBuf::from)) { + if let Some(stripped) = path.strip_prefix("~/").or_else(|| path.strip_prefix("~\\")) { + if let Some(home) = app_handle.path().home_dir().ok().or_else(|| { + std::env::var("USERPROFILE") + .ok() + .map(std::path::PathBuf::from) + }) { resolved = home.join(stripped); } else { log::warn!("Failed to resolve home directory for ~ expansion"); } } else if path == "~" { - if let Some(home) = app_handle.path().home_dir().ok().or_else(|| std::env::var("USERPROFILE").ok().map(std::path::PathBuf::from)) { + if let Some(home) = app_handle.path().home_dir().ok().or_else(|| { + std::env::var("USERPROFILE") + .ok() + .map(std::path::PathBuf::from) + }) { resolved = home; } else { log::warn!("Failed to resolve home directory for ~ expansion"); @@ -1516,13 +1665,18 @@ fn dispatch_deep_links(app_handle: tauri::AppHandle, deep_links: Vec) }); } -pub(crate) async fn rpc_call(port: u16, secret: &str, method: &str, params: serde_json::Value) -> Result { +pub(crate) 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); @@ -1534,12 +1688,13 @@ pub(crate) async fn rpc_call(port: u16, secret: &str, method: &str, params: serd .build() .map_err(|e| e.to_string())?; - let res = client.post(&url) + 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())?; if let Some(error) = json.get("error") { return Err(error.to_string()); @@ -1550,9 +1705,16 @@ pub(crate) async fn rpc_call(port: u16, secret: &str, method: &str, params: serd } #[tauri::command] -async fn test_aria2c(app_handle: tauri::AppHandle, state: tauri::State<'_, AppState>) -> Result { +async fn test_aria2c( + app_handle: tauri::AppHandle, + state: tauri::State<'_, AppState>, +) -> Result { let guard = app_handle.state::(); - let startup_err = guard.startup_error.lock().unwrap_or_else(|e| e.into_inner()).clone(); + let startup_err = guard + .startup_error + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(); if let Some(err) = startup_err { return Err(format!("aria2 daemon unavailable: {err}")); } @@ -1582,7 +1744,13 @@ async fn run_sidecar_version( ) -> (Option, Option, Option) { let binary_path = match resolve_bundled_binary_path(app_handle, sidecar_name) { Ok(p) => p, - Err(e) => return (None, Some(format!("Missing bundled binary '{}': {}", sidecar_name, e)), None), + Err(e) => { + return ( + None, + Some(format!("Missing bundled binary '{}': {}", sidecar_name, e)), + None, + ) + } }; if let Err(error) = validate_bundled_binary(&binary_path) { @@ -1608,7 +1776,13 @@ async fn run_sidecar_version( let output = match result { Ok(Ok(output)) => output, - Ok(Err(e)) => return (None, Some(format!("Failed to execute '{}': {}", sidecar_name, e)), None), + Ok(Err(e)) => { + return ( + None, + Some(format!("Failed to execute '{}': {}", sidecar_name, e)), + None, + ) + } Err(timeout) => { return ( None, @@ -1624,7 +1798,11 @@ async fn run_sidecar_version( }; let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string(); - let stderr_tail = if stderr.is_empty() { None } else { Some(stderr.clone()) }; + let stderr_tail = if stderr.is_empty() { + None + } else { + Some(stderr.clone()) + }; let result = if output.status.success() { let stdout = String::from_utf8_lossy(&output.stdout).trim().to_string(); @@ -1636,7 +1814,15 @@ async fn run_sidecar_version( } else { format!("Exited with code {:?}", output.status.code()) }; - (if stdout.is_empty() { None } else { Some(stdout) }, Some(err), stderr_tail) + ( + if stdout.is_empty() { + None + } else { + Some(stdout) + }, + Some(err), + stderr_tail, + ) }; if result.1.is_none() { @@ -1747,18 +1933,34 @@ async fn check_aria2(app_handle: &tauri::AppHandle, port: u16, secret: &str) -> let expected_sidecar = crate::platform::engine_binary_name(sidecar_name); let resolved = resolve_bundled_binary_path(app_handle, sidecar_name); - let resolved_path = resolved.as_ref().ok().map(|p| p.to_string_lossy().to_string()); + let resolved_path = resolved + .as_ref() + .ok() + .map(|p| p.to_string_lossy().to_string()); let (startup_err, daemon_stderr) = { let guard = app_handle.state::(); - let se = guard.startup_error.lock().unwrap_or_else(|e| e.into_inner()).clone(); - let stderr = guard.last_stderr.lock().unwrap_or_else(|e| e.into_inner()).clone(); + let se = guard + .startup_error + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(); + let stderr = guard + .last_stderr + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(); (se, stderr) }; let daemon_alive = startup_err.is_none(); - let last_stderr_tail = if daemon_stderr.is_empty() { None } else { Some(daemon_stderr) }; + let last_stderr_tail = if daemon_stderr.is_empty() { + None + } else { + Some(daemon_stderr) + }; - let (version_raw, run_error, stderr_tail) = run_sidecar_version(app_handle, sidecar_name, &["--version"]).await; + let (version_raw, run_error, stderr_tail) = + run_sidecar_version(app_handle, sidecar_name, &["--version"]).await; let version = version_raw.and_then(|v| v.lines().next().map(|l| l.trim().to_string())); let rpc_ready = if daemon_alive { @@ -1771,7 +1973,9 @@ async fn check_aria2(app_handle: &tauri::AppHandle, port: u16, secret: &str) -> let error = startup_err.or(run_error); let ready = daemon_alive && rpc_ready && version.is_some(); - let remediation_hint = error.as_ref().and_then(|e| generate_remediation_hint(e, sidecar_name)); + let remediation_hint = error + .as_ref() + .and_then(|e| generate_remediation_hint(e, sidecar_name)); EngineStatusItem { name: "Aria2".to_string(), @@ -1798,7 +2002,10 @@ async fn check_ytdlp(app_handle: &tauri::AppHandle) -> EngineStatusItem { let expected_sidecar = crate::platform::engine_binary_name(sidecar_name); let resolved = resolve_bundled_binary_path(app_handle, sidecar_name); - let resolved_path = resolved.as_ref().ok().map(|p| p.to_string_lossy().to_string()); + let resolved_path = resolved + .as_ref() + .ok() + .map(|p| p.to_string_lossy().to_string()); let (has_internal_dir, has_python_framework) = if let Some(ref path) = resolved_path { let parent = std::path::Path::new(path).parent().map(|p| p.to_path_buf()); @@ -1819,7 +2026,8 @@ async fn check_ytdlp(app_handle: &tauri::AppHandle) -> EngineStatusItem { (false, false) }; - let (version_raw, run_error, stderr_tail) = run_sidecar_version(app_handle, sidecar_name, &["--version"]).await; + let (version_raw, run_error, stderr_tail) = + run_sidecar_version(app_handle, sidecar_name, &["--version"]).await; let version = version_raw.and_then(|v| v.lines().next().map(|l| l.trim().to_string())); let mut error = run_error; @@ -1827,11 +2035,16 @@ async fn check_ytdlp(app_handle: &tauri::AppHandle) -> EngineStatusItem { if error.is_none() && has_internal_dir && !has_python_framework { error = Some("_internal/Python.framework was not found beside yt-dlp sidecar".to_string()); - remediation_hint = Some("The yt-dlp distribution is missing its embedded Python runtime. Reinstall Firelink.".to_string()); + remediation_hint = Some( + "The yt-dlp distribution is missing its embedded Python runtime. Reinstall Firelink." + .to_string(), + ); } if remediation_hint.is_none() { - remediation_hint = error.as_ref().and_then(|e| generate_remediation_hint(e, sidecar_name)); + remediation_hint = error + .as_ref() + .and_then(|e| generate_remediation_hint(e, sidecar_name)); } EngineStatusItem { @@ -1859,9 +2072,13 @@ async fn check_ffmpeg(app_handle: &tauri::AppHandle) -> EngineStatusItem { let expected_sidecar = crate::platform::engine_binary_name(sidecar_name); let resolved = resolve_bundled_binary_path(app_handle, sidecar_name); - let resolved_path = resolved.as_ref().ok().map(|p| p.to_string_lossy().to_string()); + let resolved_path = resolved + .as_ref() + .ok() + .map(|p| p.to_string_lossy().to_string()); - let (version_raw, run_error, stderr_tail) = run_sidecar_version(app_handle, sidecar_name, &["-version"]).await; + let (version_raw, run_error, stderr_tail) = + run_sidecar_version(app_handle, sidecar_name, &["-version"]).await; let version = version_raw.as_ref().and_then(|text| { text.lines().next().and_then(|first| { let re = regex::Regex::new(r"(?i)version\s+([\d\.]+)").unwrap(); @@ -1869,13 +2086,17 @@ async fn check_ffmpeg(app_handle: &tauri::AppHandle) -> EngineStatusItem { caps.get(1).map(|m| m.as_str().to_string()) } else { let parts: Vec<&str> = first.split_whitespace().collect(); - parts.get(2).map(|v| v.split('-').next().unwrap_or(v).to_string()) + parts + .get(2) + .map(|v| v.split('-').next().unwrap_or(v).to_string()) } }) }); let error = run_error; - let remediation_hint = error.as_ref().and_then(|e| generate_remediation_hint(e, sidecar_name)); + let remediation_hint = error + .as_ref() + .and_then(|e| generate_remediation_hint(e, sidecar_name)); EngineStatusItem { name: "FFmpeg".to_string(), @@ -1902,16 +2123,27 @@ async fn check_deno(app_handle: &tauri::AppHandle) -> EngineStatusItem { let expected_sidecar = crate::platform::engine_binary_name(sidecar_name); let resolved = resolve_bundled_binary_path(app_handle, sidecar_name); - let resolved_path = resolved.as_ref().ok().map(|p| p.to_string_lossy().to_string()); + let resolved_path = resolved + .as_ref() + .ok() + .map(|p| p.to_string_lossy().to_string()); - let (version_raw, run_error, stderr_tail) = run_sidecar_version(app_handle, sidecar_name, &["--version"]).await; - let version = version_raw.as_ref().and_then(|text| { - let re = regex::Regex::new(r"deno\s+(\d+\.\d+\.\d+)").ok()?; - re.captures(text).and_then(|c| c.get(1)).map(|m| m.as_str().to_string()) - }).or(version_raw); + let (version_raw, run_error, stderr_tail) = + run_sidecar_version(app_handle, sidecar_name, &["--version"]).await; + let version = version_raw + .as_ref() + .and_then(|text| { + let re = regex::Regex::new(r"deno\s+(\d+\.\d+\.\d+)").ok()?; + re.captures(text) + .and_then(|c| c.get(1)) + .map(|m| m.as_str().to_string()) + }) + .or(version_raw); let error = run_error; - let remediation_hint = error.as_ref().and_then(|e| generate_remediation_hint(e, sidecar_name)); + let remediation_hint = error + .as_ref() + .and_then(|e| generate_remediation_hint(e, sidecar_name)); EngineStatusItem { name: "Deno".to_string(), @@ -1958,7 +2190,12 @@ async fn get_aria2_engine_status( app_handle: tauri::AppHandle, state: tauri::State<'_, AppState>, ) -> Result { - Ok(check_aria2(&app_handle, state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret).await) + Ok(check_aria2( + &app_handle, + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret, + ) + .await) } #[tauri::command] @@ -1967,7 +2204,9 @@ async fn get_ytdlp_engine_status(app_handle: tauri::AppHandle) -> Result Result { +async fn get_ffmpeg_engine_status( + app_handle: tauri::AppHandle, +) -> Result { Ok(check_ffmpeg(&app_handle).await) } @@ -1976,8 +2215,10 @@ async fn get_deno_engine_status(app_handle: tauri::AppHandle) -> Result Result { +fn resolve_bundled_binary_path( + app_handle: &tauri::AppHandle, + binary_name: &str, +) -> Result { crate::engines::resolve_bundled_binary_path(app_handle, binary_name) } @@ -2001,7 +2242,6 @@ pub(crate) async fn start_media_download_internal( ) -> Result<(), String> { let safe_filename = crate::download_ownership::canonical_download_filename(&filename); - let resolved_dest = resolve_path(&destination, &app_handle); if !is_safe_path(&resolved_dest, &app_handle) { @@ -2015,14 +2255,22 @@ pub(crate) async fn start_media_download_internal( let out_path = resolved_dest.join(&safe_filename); let total_tracks: f64 = if let Some(ref format) = format_selector { - if format.contains('+') { 2.0 } else { 1.0 } + if format.contains('+') { + 2.0 + } else { + 1.0 + } } else { 1.0 }; use tauri_plugin_shell::ShellExt; - let mut config_file = tempfile::Builder::new().prefix("ytdlp-").suffix(".conf").tempfile().map_err(|e| e.to_string())?; + 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() { @@ -2035,16 +2283,24 @@ pub(crate) async fn start_media_download_internal( } } if let Some(headers) = headers { - for header in headers.lines().map(str::trim).filter(|header| !header.is_empty()) { + 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())?; + config_file + .write_all(config_content.as_bytes()) + .map_err(|e| e.to_string())?; let config_path = config_file.into_temp_path(); use crate::ipc::DownloadStateEvent; - use crate::retry::{BackoffOutcome, MAX_RETRIES, backoff_and_emit_cancel, is_transient_network_error}; + use crate::retry::{ + backoff_and_emit_cancel, is_transient_network_error, BackoffOutcome, MAX_RETRIES, + }; const STDERR_TAIL: usize = 2048; @@ -2079,25 +2335,37 @@ pub(crate) async fn start_media_download_internal( while strike <= MAX_RETRIES { let ytdlp_path = resolve_bundled_binary_path(&app_handle, "yt-dlp")?; - let mut cmd = app_handle.shell().command(&ytdlp_path) - .arg("--newline") - .arg("--progress-delta").arg("0.2") - .arg("--progress-template") - .arg(format!("download:{MEDIA_PROGRESS_PREFIX}%(progress)j")) - - .arg("--socket-timeout").arg("20") - .arg("--retries").arg("3") - .arg("--extractor-retries").arg("3") - .arg("--downloader").arg(&aria2c_path) - .arg("--downloader-args").arg("aria2c:-c -x 16 -s 16 -k 1M --summary-interval=1") - .arg("--ffmpeg-location").arg(&ffmpeg_path) - .arg("--js-runtimes").arg(format!("deno:{}", deno_path.to_string_lossy())) - .arg("--concurrent-fragments").arg("4") - .arg("--no-warnings") - .arg("--continue") - .arg("--compat-options").arg("no-youtube-unavailable-videos") - .arg("-o").arg(out_path.to_string_lossy().to_string()) - .env("PATH", &trusted_path); + let mut cmd = app_handle + .shell() + .command(&ytdlp_path) + .arg("--newline") + .arg("--progress-delta") + .arg("0.2") + .arg("--progress-template") + .arg(format!("download:{MEDIA_PROGRESS_PREFIX}%(progress)j")) + .arg("--socket-timeout") + .arg("20") + .arg("--retries") + .arg("3") + .arg("--extractor-retries") + .arg("3") + .arg("--downloader") + .arg(&aria2c_path) + .arg("--downloader-args") + .arg("aria2c:-c -x 16 -s 16 -k 1M --summary-interval=1") + .arg("--ffmpeg-location") + .arg(&ffmpeg_path) + .arg("--js-runtimes") + .arg(format!("deno:{}", deno_path.to_string_lossy())) + .arg("--concurrent-fragments") + .arg("4") + .arg("--no-warnings") + .arg("--continue") + .arg("--compat-options") + .arg("no-youtube-unavailable-videos") + .arg("-o") + .arg(out_path.to_string_lossy().to_string()) + .env("PATH", &trusted_path); if let Some(limit) = speed_limit.as_ref() { if !limit.is_empty() { @@ -2116,7 +2384,9 @@ pub(crate) async fn start_media_download_internal( if let Some(cs) = cookie_source.as_ref() { let mut cs = cs.clone(); if !cs.is_empty() && cs != "none" { - if cs == "safari" { cs = "safari:".to_string() } + if cs == "safari" { + cs = "safari:".to_string() + } cmd = cmd.arg("--cookies-from-browser").arg(cs); } } @@ -2154,7 +2424,9 @@ pub(crate) async fn start_media_download_internal( cmd = cmd.arg("--").arg(&url); - let (mut rx, child) = cmd.spawn().map_err(|e| format!("Failed to spawn yt-dlp: {}", e))?; + let (mut rx, child) = cmd + .spawn() + .map_err(|e| format!("Failed to spawn yt-dlp: {}", e))?; log::info!("yt-dlp spawned for id: {} (strike {})", id, strike); let mut stderr_tail = String::new(); @@ -2303,17 +2575,12 @@ pub(crate) async fn start_media_download_internal( } let reason = failure_reason.clone(); - let outcome = backoff_and_emit_cancel( - strike, - reason, - cancel_rx, - |retry_reason| { - let _ = app_handle.emit( - "download-state", - DownloadStateEvent::retrying(id, retry_reason), - ); - }, - ) + let outcome = backoff_and_emit_cancel(strike, reason, cancel_rx, |retry_reason| { + let _ = app_handle.emit( + "download-state", + DownloadStateEvent::retrying(id, retry_reason), + ); + }) .await; if outcome == BackoffOutcome::Aborted { @@ -2416,7 +2683,10 @@ async fn resume_download( id: String, ) -> Result { let Some(gid) = state.queue_manager.aria2_gid_for_download(&id) else { - log::info!("aria2 resume [{}]: no mapped gid; re-enqueue is permitted", id); + log::info!( + "aria2 resume [{}]: no mapped gid; re-enqueue is permitted", + id + ); state.queue_manager.release_registered_id(&id).await; return Ok(false); }; @@ -2427,7 +2697,12 @@ async fn resume_download( return Ok(false); } - let status = aria2_download_status(state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret, &gid).await?; + let status = aria2_download_status( + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret, + &gid, + ) + .await?; match status.as_str() { "paused" => { use tauri::Emitter; @@ -2435,13 +2710,13 @@ async fn resume_download( "download-state", crate::ipc::DownloadStateEvent::new(&id, crate::ipc::DownloadStatus::Queued), ); - + let queue_manager = state.queue_manager.clone(); let aria2_port = state.aria2_port.load(std::sync::atomic::Ordering::Relaxed); let aria2_secret = state.aria2_secret.clone(); let id_clone = id.clone(); let gid_clone = gid.clone(); - + let app_handle_clone = app_handle.clone(); tauri::async_runtime::spawn(async move { let acquired = queue_manager.ensure_aria2_permit(&id_clone).await; @@ -2540,8 +2815,18 @@ async fn remove_download( let gid = state.queue_manager.aria2_gid_for_download(&id); if let Some(gid) = gid.as_deref().filter(|gid| !gid.starts_with("native:")) { let removal_result = async { - force_remove_aria2_gid(state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret, gid).await?; - wait_for_aria2_stopped(state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret, gid).await + force_remove_aria2_gid( + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret, + gid, + ) + .await?; + wait_for_aria2_stopped( + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret, + gid, + ) + .await } .await; if let Err(error) = removal_result { @@ -2553,10 +2838,12 @@ async fn remove_download( state.queue_manager.release_permit(&id).await; log::info!("aria2 remove [{}]: gid {} stopped and forgotten", id, gid); } else { - let (tx, rx) = tokio::sync::oneshot::channel(); if matches!(active_kind, Some(crate::queue::TaskKind::Media)) { - state.download_coordinator.pause_media_with_ack(id.clone(), tx).await?; + state + .download_coordinator + .pause_media_with_ack(id.clone(), tx) + .await?; } else if let Ok(download_id) = Uuid::parse_str(&id) { let command = if delete_assets { download::DownloadCmd::CancelWithAck(download_id, tx) @@ -2643,7 +2930,12 @@ async fn detach_download_for_reconfigure( serde_json::json!([gid]), ) .await?; - wait_for_aria2_stopped(state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret, gid).await + wait_for_aria2_stopped( + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret, + gid, + ) + .await } .await; if let Err(error) = removal_result { @@ -2656,12 +2948,17 @@ async fn detach_download_for_reconfigure( state.queue_manager.release_registered_id(&id).await; log::info!("aria2 detach [{}]: gid {} stopped and forgotten", id, gid); } else { - let (tx, rx) = tokio::sync::oneshot::channel(); if matches!(active_kind, Some(crate::queue::TaskKind::Media)) { - state.download_coordinator.pause_media_with_ack(id.clone(), tx).await?; + state + .download_coordinator + .pause_media_with_ack(id.clone(), tx) + .await?; } else if let Ok(download_id) = Uuid::parse_str(&id) { - state.download_coordinator.send(crate::download::DownloadCmd::PauseWithAck(download_id, tx)).await?; + state + .download_coordinator + .send(crate::download::DownloadCmd::PauseWithAck(download_id, tx)) + .await?; } else { let _ = tx.send(()); // Fallback if no task exists } @@ -2720,14 +3017,7 @@ fn aria2_gid_not_found(error: &str) -> bool { } async fn force_remove_aria2_gid(port: u16, secret: &str, gid: &str) -> Result<(), String> { - match rpc_call( - port, - secret, - "aria2.forceRemove", - serde_json::json!([gid]), - ) - .await - { + match rpc_call(port, secret, "aria2.forceRemove", serde_json::json!([gid])).await { Ok(result) => ensure_aria2_gid_result("forceRemove", gid, &result), Err(error) if aria2_gid_not_found(&error) => { log::info!("aria2 forceRemove: gid {} was already absent", gid); @@ -2769,14 +3059,18 @@ fn update_dock_badge(_app_handle: tauri::AppHandle, count: i32) { #[cfg(target_os = "macos")] { use cocoa::appkit::NSApp; - use cocoa::base::{nil, id}; + use cocoa::base::{id, nil}; use cocoa::foundation::NSString; use objc::{msg_send, sel, sel_impl}; - + unsafe { let app = NSApp(); let dock_tile: id = msg_send![app, dockTile]; - let label = if count > 0 { count.to_string() } else { "".to_string() }; + let label = if count > 0 { + count.to_string() + } else { + "".to_string() + }; let ns_label = NSString::alloc(nil).init_str(&label); let _: () = msg_send![dock_tile, setBadgeLabel: ns_label]; } @@ -2794,7 +3088,10 @@ fn get_platform_info() -> crate::ipc::PlatformInfo { #[tauri::command] fn set_prevent_sleep(state: tauri::State<'_, AppState>, prevent: bool) -> Result<(), String> { - let mut current_preventer = state.sleep_preventer.lock().unwrap_or_else(|e| e.into_inner()); + let mut current_preventer = state + .sleep_preventer + .lock() + .unwrap_or_else(|e| e.into_inner()); if prevent { if current_preventer.is_none() { let keepawake = keepawake::Builder::default() @@ -2818,9 +3115,7 @@ pub(crate) fn execute_system_action(action: crate::ipc::PostQueueAction) -> Resu crate::ipc::PostQueueAction::Restart => { system_shutdown::reboot().map_err(|e| e.to_string()) } - crate::ipc::PostQueueAction::Sleep => { - system_shutdown::sleep().map_err(|e| e.to_string()) - } + crate::ipc::PostQueueAction::Sleep => system_shutdown::sleep().map_err(|e| e.to_string()), crate::ipc::PostQueueAction::None => Err("Invalid action".to_string()), } } @@ -2914,7 +3209,10 @@ async fn enqueue_many( &item.filename, )?; } - let tasks = items.into_iter().map(queue::EnqueueItem::into_task).collect(); + let tasks = items + .into_iter() + .map(queue::EnqueueItem::into_task) + .collect(); let results = state.queue_manager.enqueue_many(tasks).await; for result in &results { @@ -2955,7 +3253,10 @@ async fn remove_from_queue( } #[tauri::command] -async fn set_concurrent_limit(state: tauri::State<'_, AppState>, limit: usize) -> Result<(), String> { +async fn set_concurrent_limit( + state: tauri::State<'_, AppState>, + limit: usize, +) -> Result<(), String> { state.queue_manager.set_capacity(limit); Ok(()) } @@ -2973,7 +3274,10 @@ fn normalize_speed_limit_for_aria2(limit: &str) -> Option { return None; } - let unit = captures.get(2).map(|m| m.as_str().to_ascii_uppercase()).unwrap_or_default(); + let unit = captures + .get(2) + .map(|m| m.as_str().to_ascii_uppercase()) + .unwrap_or_default(); Some(if unit.is_empty() { format!("{amount}K") } else { @@ -2982,7 +3286,10 @@ fn normalize_speed_limit_for_aria2(limit: &str) -> Option { } #[tauri::command] -async fn set_global_speed_limit(state: tauri::State<'_, AppState>, limit: Option) -> Result<(), String> { +async fn set_global_speed_limit( + state: tauri::State<'_, AppState>, + limit: Option, +) -> Result<(), String> { let limit_str = limit .as_deref() .and_then(normalize_speed_limit_for_aria2) @@ -2991,8 +3298,11 @@ async fn set_global_speed_limit(state: tauri::State<'_, AppState>, limit: Option state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret, "aria2.changeGlobalOption", - serde_json::json!([{"max-overall-download-limit": limit_str}]) - ).await.map(|_| ()).map_err(|e| { + serde_json::json!([{"max-overall-download-limit": limit_str}]), + ) + .await + .map(|_| ()) + .map_err(|e| { eprintln!("Failed to set global speed limit: {}", e); e }) @@ -3043,7 +3353,12 @@ fn open_automation_settings(app_handle: tauri::AppHandle) -> Result<(), String> #[cfg(target_os = "macos")] { use tauri_plugin_opener::OpenerExt; - app_handle.opener().open_url("x-apple.systempreferences:com.apple.preference.security?Privacy_Automation", None::) + app_handle + .opener() + .open_url( + "x-apple.systempreferences:com.apple.preference.security?Privacy_Automation", + None::, + ) .map_err(|e| format!("Failed to open Automation settings: {}", e))?; Ok(()) } @@ -3155,10 +3470,10 @@ fn grant_keychain_access( app_state: tauri::State<'_, AppState>, ) -> Result { let mut connection = database.lock()?; - + // Explicitly force migration of any legacy token to the keychain let _ = crate::db::sanitize_current_settings_and_restore_token(&connection, true); - + match crate::db::hydrate_pairing_token(&mut connection, false) { Ok((token, token_changed)) => { // Update the extension server's token in memory @@ -3214,9 +3529,7 @@ fn db_save_settings( } #[tauri::command] -fn db_load_settings( - state: tauri::State<'_, crate::db::DbState>, -) -> Result, String> { +fn db_load_settings(state: tauri::State<'_, crate::db::DbState>) -> Result, String> { let connection = state.lock()?; crate::db::load_settings(&connection) } @@ -3239,9 +3552,7 @@ fn db_replace_downloads( } #[tauri::command] -fn db_get_all_queues( - state: tauri::State<'_, crate::db::DbState>, -) -> Result, String> { +fn db_get_all_queues(state: tauri::State<'_, crate::db::DbState>) -> Result, String> { let connection = state.lock()?; crate::db::load_queues(&connection) } @@ -3290,18 +3601,15 @@ fn redact_log_line(line: &str) -> String { static HEADER: OnceLock = OnceLock::new(); static QUERY: OnceLock = OnceLock::new(); let secret = SECRET.get_or_init(|| { - regex::Regex::new( - r"(?i)(authorization|cookie|password|token|secret)\s*[:=]\s*([^\s,;]+)", - ) - .expect("valid secret redaction regex") + regex::Regex::new(r"(?i)(authorization|cookie|password|token|secret)\s*[:=]\s*([^\s,;]+)") + .expect("valid secret redaction regex") }); let header = HEADER.get_or_init(|| { regex::Regex::new(r"(?i)(authorization|cookie)\s*:\s*[^\r\n]+") .expect("valid sensitive header redaction regex") }); let query = QUERY.get_or_init(|| { - regex::Regex::new(r"(https?://[^\s?]+)\?[^\s]+") - .expect("valid URL query redaction regex") + regex::Regex::new(r"(https?://[^\s?]+)\?[^\s]+").expect("valid URL query redaction regex") }); let redacted = header.replace_all(line, "$1: [redacted]"); let redacted = secret.replace_all(&redacted, "$1=[redacted]"); @@ -3345,7 +3653,9 @@ async fn read_logs(app_handle: tauri::AppHandle, limit: usize) -> Result Result<(), String> { for file in log_files(&app_handle).await? { - tokio::fs::write(&file, "").await.map_err(|e| format!("Failed to clear log file {:?}: {}", file, e))?; + tokio::fs::write(&file, "") + .await + .map_err(|e| format!("Failed to clear log file {:?}: {}", file, e))?; } Ok(()) } @@ -3364,7 +3674,11 @@ async fn export_logs( chrono::Utc::now().to_rfc3339(), ); let (aria2, ytdlp, ffmpeg, deno) = tokio::join!( - check_aria2(&app_handle, state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret), + check_aria2( + &app_handle, + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret + ), check_ytdlp(&app_handle), check_ffmpeg(&app_handle), check_deno(&app_handle), @@ -3437,15 +3751,13 @@ fn build_main_tray(app_handle: &tauri::AppHandle) -> Result<(), String> { .map_err(|e| e.to_string())?; let pause_all_i = MenuItem::with_id(app_handle, "pause_all", "Pause All", true, None::<&str>) .map_err(|e| e.to_string())?; - let resume_all_i = MenuItem::with_id(app_handle, "resume_all", "Resume All", true, None::<&str>) - .map_err(|e| e.to_string())?; + let resume_all_i = + MenuItem::with_id(app_handle, "resume_all", "Resume All", true, None::<&str>) + .map_err(|e| e.to_string())?; let quit_i = MenuItem::with_id(app_handle, "quit", "Quit", true, None::<&str>) .map_err(|e| e.to_string())?; - let menu = Menu::with_items( - app_handle, - &[&show_i, &pause_all_i, &resume_all_i, &quit_i], - ) - .map_err(|e| e.to_string())?; + let menu = Menu::with_items(app_handle, &[&show_i, &pause_all_i, &resume_all_i, &quit_i]) + .map_err(|e| e.to_string())?; #[cfg(target_os = "macos")] let tray_icon_bytes = include_bytes!("../icons/trayTemplate.png").as_slice(); @@ -3486,8 +3798,7 @@ fn build_main_tray(app_handle: &tauri::AppHandle) -> Result<(), String> { { tray = tray.icon_as_template(true); } - tray.build(app_handle) - .map_err(|e| e.to_string())?; + tray.build(app_handle).map_err(|e| e.to_string())?; Ok(()) } @@ -3519,10 +3830,7 @@ fn get_extension_server_port(state: tauri::State<'_, AppState>) -> Option { } #[tauri::command] -fn set_extension_frontend_ready( - state: tauri::State<'_, AppState>, - ready: bool, -) { +fn set_extension_frontend_ready(state: tauri::State<'_, AppState>, ready: bool) { state .extension_frontend_ready .store(ready, Ordering::Release); @@ -3540,24 +3848,33 @@ mod tests { aggregate_media_fraction, build_media_format_options, collect_download_uris, is_excluded_yt_dlp_format, json_lower, media_progress_speed, normalize_speed_limit_for_aria2, parse_firelink_deep_link, parse_media_progress_line, - FirelinkDeepLink, - redact_log_line, MediaProgress, MEDIA_PROGRESS_PREFIX, + redact_log_line, FirelinkDeepLink, MediaProgress, MEDIA_PROGRESS_PREFIX, }; use serde_json::json; use std::time::{Duration, Instant}; #[test] fn normalizes_bare_global_speed_limits_as_kib_per_second() { - assert_eq!(normalize_speed_limit_for_aria2("1024"), Some("1024K".to_string())); - assert_eq!(normalize_speed_limit_for_aria2("512K"), Some("512K".to_string())); - assert_eq!(normalize_speed_limit_for_aria2("1.5 MB/s"), Some("1.5M".to_string())); + assert_eq!( + normalize_speed_limit_for_aria2("1024"), + Some("1024K".to_string()) + ); + assert_eq!( + normalize_speed_limit_for_aria2("512K"), + Some("512K".to_string()) + ); + assert_eq!( + normalize_speed_limit_for_aria2("1.5 MB/s"), + Some("1.5M".to_string()) + ); assert_eq!(normalize_speed_limit_for_aria2("0"), None); assert_eq!(normalize_speed_limit_for_aria2("bad"), None); } #[test] fn redacts_secrets_and_signed_url_queries_from_support_logs() { - let line = "Authorization: bearer-secret Cookie=session=abc https://example.com/file?token=secret"; + let line = + "Authorization: bearer-secret Cookie=session=abc https://example.com/file?token=secret"; let redacted = redact_log_line(line); assert!(!redacted.contains("bearer-secret")); assert!(!redacted.contains("session=abc")); @@ -3668,25 +3985,25 @@ mod tests { "vcodec": "none", "acodec": "mp4a.40.2", "filesize": 10_000_000_u64 - }) + }), ]; let options = build_media_format_options(&formats, Some(600.0)); - assert!(!options.iter().any(|format| format.ext == "mhtml")); - assert!(!options.iter().any(|format| format.resolution == "Best")); - assert!(!options.iter().any(|format| format.resolution == "1440p")); - assert!(options.iter().any(|format| { - format.resolution == "1080p" - && format.ext == "mkv" - && format.format_label == "MKV • H.264 + AAC" - && format.filesize == Some(110_000_000) - })); - assert!(options.iter().any(|format| { - format.resolution == "1080p" - && format.ext == "mp4" - && format.format_label == "MP4 • H.264 + AAC" - && format.filesize == Some(110_000_000) + assert!(!options.iter().any(|format| format.ext == "mhtml")); + assert!(!options.iter().any(|format| format.resolution == "Best")); + assert!(!options.iter().any(|format| format.resolution == "1440p")); + assert!(options.iter().any(|format| { + format.resolution == "1080p" + && format.ext == "mkv" + && format.format_label == "MKV • H.264 + AAC" + && format.filesize == Some(110_000_000) + })); + assert!(options.iter().any(|format| { + format.resolution == "1080p" + && format.ext == "mp4" + && format.format_label == "MP4 • H.264 + AAC" + && format.filesize == Some(110_000_000) })); assert!(options.iter().any(|format| { format.resolution == "Audio only" && format.format_label == "M4A • AAC" @@ -3711,7 +4028,7 @@ mod tests { "vcodec": "none", "acodec": "opus", "filesize": 28_000_000_u64 - }) + }), ]; let options = build_media_format_options(&formats, Some(1_800.0)); let option = options @@ -3913,7 +4230,8 @@ pub fn run() { let server_frontend_ready = extension_frontend_ready.clone(); let extension_server_port = Arc::new(RwLock::new(None)); let server_extension_port = extension_server_port.clone(); - let (extension_server_shutdown_tx, extension_server_shutdown_rx) = tokio::sync::watch::channel(false); + let (extension_server_shutdown_tx, extension_server_shutdown_rx) = + tokio::sync::watch::channel(false); let initial_aria2_port = 6800; // Will be determined dynamically in background let aria2_port = Arc::new(std::sync::atomic::AtomicU16::new(initial_aria2_port)); @@ -3947,7 +4265,7 @@ pub fn run() { .map_err(|error| format!("failed to initialize persistence: {error}"))?; let initial_pairing_token = { // Generate a temporary session token for the extension server on startup. - // The frontend will hydrate the real token via IPC once it mounts, + // The frontend will hydrate the real token via IPC once it mounts, // avoiding any macOS system prompts before the UI is fully visible. format!( "{}{}", @@ -4057,13 +4375,13 @@ pub fn run() { let mut success = false; for attempt_port in 6800..6900 { let mut cmd = std::process::Command::new(&binary_path); - + let mut config_file = tempfile::Builder::new().prefix("aria2-").suffix(".conf").tempfile().expect("failed to create aria2 config file"); use std::io::Write; let config_content = format!("rpc-secret={}\n", aria2_secret_clone); config_file.write_all(config_content.as_bytes()).expect("failed to write aria2 config file"); let config_path = config_file.into_temp_path(); - + cmd.arg("--enable-rpc=true") .arg(format!("--conf-path={}", config_path.display())) .arg(format!("--rpc-listen-port={}", attempt_port)) @@ -4076,14 +4394,14 @@ pub fn run() { .arg("--download-result=hide") .arg("--max-concurrent-downloads=9999") .arg("--check-certificate=true"); - + if let Some(limit) = normalize_speed_limit_for_aria2(&global_speed_limit) { cmd.arg(format!("--max-overall-download-limit={}", limit)); } - + cmd.stdout(std::process::Stdio::null()); cmd.stderr(std::process::Stdio::piped()); - + match cmd.spawn() { Ok(mut child) => { // Give it a moment to fail if port is in use @@ -4098,7 +4416,7 @@ pub fn run() { aria2_port_clone.store(attempt_port, std::sync::atomic::Ordering::Relaxed); ws_port = attempt_port; success = true; - + let daemon_app = app_handle_bg.clone(); if let Some(stderr) = child.stderr.take() { std::thread::spawn(move || { @@ -4121,11 +4439,11 @@ pub fn run() { } }); } - + let guard = app_handle_bg.state::(); *guard.child.lock().unwrap() = Some(child); *guard.config_path.lock().unwrap() = Some(config_path); - + let mut last_err = String::new(); let start = std::time::Instant::now(); let mut ready = false; @@ -4237,7 +4555,7 @@ pub fn run() { let total = status_info.get("totalLength").and_then(|s| s.as_str()).unwrap_or("0").parse::().unwrap_or(0); let completed = status_info.get("completedLength").and_then(|s| s.as_str()).unwrap_or("0").parse::().unwrap_or(0); let speed_bytes = status_info.get("downloadSpeed").and_then(|s| s.as_str()).unwrap_or("0").parse::().unwrap_or(0.0); - + let fraction = if total > 0 { completed as f64 / total as f64 } else { 0.0 }; let speed = crate::download::format_speed(speed_bytes); let eta = if speed_bytes > 0.0 && total > completed { @@ -4250,7 +4568,7 @@ pub fn run() { } else { None }; - + use tauri::Emitter; let _ = app_handle_poll.emit("download-progress", DownloadProgressEvent { id, @@ -4348,6 +4666,6 @@ pub fn run() { } }); } -mod extension_server; mod db; +mod extension_server; mod scheduler; diff --git a/src-tauri/src/main.rs b/src-tauri/src/main.rs index 532c4a1..9835479 100644 --- a/src-tauri/src/main.rs +++ b/src-tauri/src/main.rs @@ -1,5 +1,8 @@ // Prevents additional console window on Windows in release, DO NOT REMOVE!! -#![cfg_attr(all(not(debug_assertions), target_os = "windows"), windows_subsystem = "windows")] +#![cfg_attr( + all(not(debug_assertions), target_os = "windows"), + windows_subsystem = "windows" +)] fn main() { firelink_lib::run() diff --git a/src-tauri/src/parity.rs b/src-tauri/src/parity.rs index 261df91..b4896bc 100644 --- a/src-tauri/src/parity.rs +++ b/src-tauri/src/parity.rs @@ -8,7 +8,11 @@ pub async fn get_system_proxy() -> Result, String> { match sysproxy::Sysproxy::get_system_proxy() { Ok(proxy) => { if proxy.enable { - let protocol = if proxy.host.contains("://") { "" } else { "http://" }; + let protocol = if proxy.host.contains("://") { + "" + } else { + "http://" + }; Ok(Some(format!("{}{}:{}", protocol, proxy.host, proxy.port))) } else { Ok(None) @@ -26,12 +30,28 @@ pub fn get_file_category(filename: String) -> DownloadCategory { .map(|s| s.to_lowercase()) .unwrap_or_default(); - let music_exts = ["mp3", "wav", "aac", "flac", "ogg", "m4a", "wma", "alac", "ape", "mid", "midi"]; - let movie_exts = ["mp4", "mkv", "avi", "mov", "wmv", "flv", "webm", "m4v", "mpeg", "mpg", "3gp", "ts", "vob"]; - let compressed_exts = ["zip", "rar", "7z", "tar", "gz", "xz", "bz2", "lz", "lzma", "zst", "iso", "cab", "tgz", "tbz", "z", "sit", "sitx"]; - let picture_exts = ["jpg", "jpeg", "png", "gif", "webp", "bmp", "tiff", "svg", "ico", "heic", "raw", "psd", "ai"]; - let document_exts = ["pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx", "txt", "rtf", "csv", "md", "epub", "mobi", "azw3"]; - let app_exts = ["exe", "msi", "bat", "cmd", "app", "dmg", "pkg", "apk", "appx", "deb", "rpm", "appimage", "run", "sh", "bin", "jar"]; + let music_exts = [ + "mp3", "wav", "aac", "flac", "ogg", "m4a", "wma", "alac", "ape", "mid", "midi", + ]; + let movie_exts = [ + "mp4", "mkv", "avi", "mov", "wmv", "flv", "webm", "m4v", "mpeg", "mpg", "3gp", "ts", "vob", + ]; + let compressed_exts = [ + "zip", "rar", "7z", "tar", "gz", "xz", "bz2", "lz", "lzma", "zst", "iso", "cab", "tgz", + "tbz", "z", "sit", "sitx", + ]; + let picture_exts = [ + "jpg", "jpeg", "png", "gif", "webp", "bmp", "tiff", "svg", "ico", "heic", "raw", "psd", + "ai", + ]; + let document_exts = [ + "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx", "txt", "rtf", "csv", "md", "epub", + "mobi", "azw3", + ]; + let app_exts = [ + "exe", "msi", "bat", "cmd", "app", "dmg", "pkg", "apk", "appx", "deb", "rpm", "appimage", + "run", "sh", "bin", "jar", + ]; if music_exts.contains(&ext.as_str()) { DownloadCategory::Musics @@ -65,8 +85,13 @@ pub struct AvailableReleaseUpdate { #[serde(tag = "type")] #[ts(export, export_to = "../../src/bindings/")] pub enum ReleaseCheckOutcome { - UpdateAvailable { update: AvailableReleaseUpdate }, - UpToDate { latest_version: String, local_version: String }, + UpdateAvailable { + update: AvailableReleaseUpdate, + }, + UpToDate { + latest_version: String, + local_version: String, + }, } #[derive(Deserialize)] @@ -81,11 +106,14 @@ struct GitHubRelease { } #[tauri::command] -pub async fn check_for_updates(app_handle: tauri::AppHandle) -> Result { +pub async fn check_for_updates( + app_handle: tauri::AppHandle, +) -> Result { let current_version = app_handle.package_info().version.to_string(); - + let client = reqwest::Client::new(); - let res = client.get("https://api.github.com/repos/nimbold/Firelink/releases?per_page=30") + let res = client + .get("https://api.github.com/repos/nimbold/Firelink/releases?per_page=30") .header("User-Agent", "Firelink") .header("Accept", "application/vnd.github+json") .send() @@ -97,11 +125,12 @@ pub async fn check_for_updates(app_handle: tauri::AppHandle) -> Result = res.json().await.map_err(|e| e.to_string())?; - - let latest_stable = releases.into_iter() + + let latest_stable = releases + .into_iter() .filter(|r| !r.draft && !r.prerelease) .max_by(|a, b| cmp_versions(&a.tag_name, &b.tag_name)); - + let release = match latest_stable { Some(r) => r, None => return Err("No stable release was found.".to_string()), @@ -115,10 +144,12 @@ pub async fn check_for_updates(app_handle: tauri::AppHandle) -> Result Result std::cmp::Ordering { use semver::Version; - + let a_clean = a.trim_start_matches(['v', 'V']); let b_clean = b.trim_start_matches(['v', 'V']); - + let a_ver = Version::parse(a_clean).unwrap_or_else(|_| Version::new(0, 0, 0)); let b_ver = Version::parse(b_clean).unwrap_or_else(|_| Version::new(0, 0, 0)); - + a_ver.cmp(&b_ver) } @@ -181,14 +212,18 @@ pub async fn create_category_directories( } pub static SUPPORTED_DOMAINS: &[&str] = &[ - "youtube.com", "youtu.be", - "twitter.com", "x.com", + "youtube.com", + "youtu.be", + "twitter.com", + "x.com", "vimeo.com", "twitch.tv", "instagram.com", "tiktok.com", - "facebook.com", "fb.watch", - "reddit.com", "v.redd.it", + "facebook.com", + "fb.watch", + "reddit.com", + "v.redd.it", "soundcloud.com", ]; diff --git a/src-tauri/src/platform.rs b/src-tauri/src/platform.rs index 68ee8ca..4bbdd23 100644 --- a/src-tauri/src/platform.rs +++ b/src-tauri/src/platform.rs @@ -97,8 +97,10 @@ pub fn is_windows_reserved_filename(filename: &str) -> bool { .unwrap_or(filename) .trim_end_matches(['.', ' ']) .to_ascii_uppercase(); - matches!(stem.as_str(), "CON" | "PRN" | "AUX" | "NUL" | "CLOCK$" | "CONIN$" | "CONOUT$") - || numbered_windows_device(&stem, "COM") + matches!( + stem.as_str(), + "CON" | "PRN" | "AUX" | "NUL" | "CLOCK$" | "CONIN$" | "CONOUT$" + ) || numbered_windows_device(&stem, "COM") || numbered_windows_device(&stem, "LPT") } @@ -125,10 +127,18 @@ mod tests { #[test] fn recognizes_windows_reserved_device_names() { - for filename in ["CON", "con.txt", "PRN.", "aux.mp4", "NUL", "COM1.zip", "lpt9"] { + for filename in [ + "CON", "con.txt", "PRN.", "aux.mp4", "NUL", "COM1.zip", "lpt9", + ] { assert!(is_windows_reserved_filename(filename), "{filename}"); } - for filename in ["console.txt", "com0.zip", "com10.zip", "lpt.txt", "movie.mp4"] { + for filename in [ + "console.txt", + "com0.zip", + "com10.zip", + "lpt.txt", + "movie.mp4", + ] { assert!(!is_windows_reserved_filename(filename), "{filename}"); } } diff --git a/src-tauri/src/queue.rs b/src-tauri/src/queue.rs index 9bf83a6..10c34bb 100644 --- a/src-tauri/src/queue.rs +++ b/src-tauri/src/queue.rs @@ -1,14 +1,14 @@ use crate::ipc::{DownloadStateEvent, DownloadStatus, QueueDirection}; -use crate::retry::{BackoffOutcome, MAX_RETRIES, backoff_and_emit, is_transient_network_error}; +use crate::retry::{backoff_and_emit, is_transient_network_error, BackoffOutcome, MAX_RETRIES}; +use log; use serde::Deserialize; +use serde_json; use std::collections::{HashMap, HashSet, VecDeque}; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use tauri::{AppHandle, Manager}; use tokio::sync::{Mutex, Notify, OwnedSemaphorePermit, Semaphore}; use ts_rs::TS; -use serde_json; -use log; /// Default capacity when no setting is read yet. pub const DEFAULT_MAX_CONCURRENT: usize = 3; @@ -120,14 +120,12 @@ pub struct QueueManager { impl QueueManager { /// Production constructor. Wired up in lib.rs setup(). pub fn new(app_handle: AppHandle, capacity: usize) -> Self { - let spawner: Arc = - Arc::new(ProductionSpawner::new(app_handle.clone())); + let spawner: Arc = Arc::new(ProductionSpawner::new(app_handle.clone())); Self::test_new(app_handle, capacity, spawner) } } impl QueueManager { - /// Test-only constructor injecting a fake spawner. pub fn test_new( app_handle: AppHandle, @@ -193,11 +191,7 @@ impl QueueManager { /// Acquire a permit from the semaphore (blocks until one is available). pub async fn acquire_permit(&self) -> Option { - self.semaphore - .clone() - .acquire_owned() - .await - .ok() + self.semaphore.clone().acquire_owned().await.ok() } /// Park an already-acquired permit under `id`. @@ -258,9 +252,13 @@ impl QueueManager { .map(|(id, _)| id.clone()) .collect() }; - + for id in ids_to_fail { - self.apply_completion(&id, PendingOutcome::Error("Aria2 WebSocket connection lost".to_string())).await; + self.apply_completion( + &id, + PendingOutcome::Error("Aria2 WebSocket connection lost".to_string()), + ) + .await; } } @@ -271,10 +269,9 @@ impl QueueManager { fn emit_state(&self, id: impl Into, status: DownloadStatus) { use tauri::Emitter; - let _ = self.app_handle.emit( - "download-state", - DownloadStateEvent::new(id, status), - ); + let _ = self + .app_handle + .emit("download-state", DownloadStateEvent::new(id, status)); } /// Resize the global concurrency limit. Grow adds permits immediately; @@ -328,12 +325,7 @@ impl QueueManager { continue; } // (2) Acquire a slot. - let permit_opt = self - .semaphore - .clone() - .acquire_owned() - .await - .ok(); + let permit_opt = self.semaphore.clone().acquire_owned().await.ok(); let permit = match permit_opt { Some(p) => p, None => break, // Semaphore closed, exit dispatcher @@ -380,7 +372,10 @@ impl QueueManager { // aria2's RPC returns instantly, so the permit must outlive the // dispatch_one call. Media/Native runners release on exit. self.park_permit(&id, permit).await; - self.active_kinds.lock().await.insert(id.clone(), task.kind.clone()); + self.active_kinds + .lock() + .await + .insert(id.clone(), task.kind.clone()); self.emit_state(&id, DownloadStatus::Downloading); match task.kind { @@ -513,8 +508,6 @@ impl QueueManager { self.aria2_retry_cancelled.lock().await.remove(id); } - - pub fn aria2_gid_for_download(&self, id: &str) -> Option { self.aria2_gids .read() @@ -592,10 +585,10 @@ impl QueueManager { let id = match id { Some(id) => id, None => { - self.pending_completion - .lock() - .await - .insert(gid.to_string(), (String::new(), PendingOutcome::Error(error))); + self.pending_completion.lock().await.insert( + gid.to_string(), + (String::new(), PendingOutcome::Error(error)), + ); return; } }; @@ -617,13 +610,15 @@ impl QueueManager { let transient = is_transient_network_error(&error); let strikes_left = strike < MAX_RETRIES; if !(transient && strikes_left) { - self.apply_completion(&id, PendingOutcome::Error(error)).await; + self.apply_completion(&id, PendingOutcome::Error(error)) + .await; return; } let payload = self.aria2_payloads.lock().await.get(&id).cloned(); if payload.is_none() { - self.apply_completion(&id, PendingOutcome::Error(error)).await; + self.apply_completion(&id, PendingOutcome::Error(error)) + .await; return; } let payload = payload.unwrap(); @@ -702,11 +697,8 @@ impl QueueManager { this.emit_state(&id_for_task, DownloadStatus::Downloading); } Err(retry_error) => { - this.apply_completion( - &id_for_task, - PendingOutcome::Error(retry_error), - ) - .await; + this.apply_completion(&id_for_task, PendingOutcome::Error(retry_error)) + .await; } } }); @@ -853,7 +845,8 @@ impl SidecarSpawner for ProductionSpawner { "dir".to_string(), serde_json::json!(resolved_dest.to_string_lossy().to_string()), ); - let safe_filename = crate::download_ownership::canonical_download_filename(&payload.filename); + let safe_filename = + crate::download_ownership::canonical_download_filename(&payload.filename); options.insert("out".to_string(), serde_json::json!(safe_filename)); let conn = payload.connections.unwrap_or(1); options.insert("split".to_string(), serde_json::json!(conn.to_string())); @@ -904,7 +897,14 @@ impl SidecarSpawner for ProductionSpawner { let uris = crate::collect_download_uris(&payload.url, payload.mirrors.as_deref()); let params = serde_json::json!([uris, options]); - match crate::rpc_call(state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), &state.aria2_secret, "aria2.addUri", params).await { + match crate::rpc_call( + state.aria2_port.load(std::sync::atomic::Ordering::Relaxed), + &state.aria2_secret, + "aria2.addUri", + params, + ) + .await + { Ok(result) => { let gid = result.as_str().unwrap_or("").to_string(); if gid.is_empty() { @@ -919,13 +919,17 @@ impl SidecarSpawner for ProductionSpawner { log::warn!("aria2 addUri failed, falling back to native: {}", e); let download_id = uuid::Uuid::parse_str(id).map_err(|e| e.to_string())?; let mt = payload.max_tries.unwrap_or(1).max(1) as u32; - let safe_filename = crate::download_ownership::canonical_download_filename(&payload.filename); + let safe_filename = + crate::download_ownership::canonical_download_filename(&payload.filename); state .download_coordinator .send(crate::download::DownloadCmd::Start(Box::new( crate::download::DownloadPayload { id: download_id, - urls: crate::collect_download_uris(&payload.url, payload.mirrors.as_deref()), + urls: crate::collect_download_uris( + &payload.url, + payload.mirrors.as_deref(), + ), output_path: resolved_dest.join(safe_filename), speed_limit: payload.speed_limit.clone(), username: payload.username.clone(), @@ -995,7 +999,10 @@ impl SidecarSpawner for ProductionSpawner { let _ = crate::download_ownership::set_primary_path(&self.app_handle, id, &path); } } - let _ = state.download_coordinator.finish_media(id.to_string()).await; + let _ = state + .download_coordinator + .finish_media(id.to_string()) + .await; outcome } @@ -1004,7 +1011,8 @@ impl SidecarSpawner for ProductionSpawner { let download_id = uuid::Uuid::parse_str(id).map_err(|e| e.to_string())?; let mt = payload.max_tries.unwrap_or(1).max(1) as u32; let resolved_dest = crate::resolve_path(&payload.destination, &self.app_handle); - let safe_filename = crate::download_ownership::canonical_download_filename(&payload.filename); + let safe_filename = + crate::download_ownership::canonical_download_filename(&payload.filename); let output_path = resolved_dest.join(safe_filename); let _ = crate::download_ownership::set_primary_path(&self.app_handle, id, &output_path); state diff --git a/src-tauri/src/retry.rs b/src-tauri/src/retry.rs index e480b8c..6b894c1 100644 --- a/src-tauri/src/retry.rs +++ b/src-tauri/src/retry.rs @@ -177,7 +177,14 @@ mod tests { #[test] fn schedule_is_three_strike_exponential() { - assert_eq!(BACKOFF_SCHEDULE, [Duration::from_secs(2), Duration::from_secs(5), Duration::from_secs(10)]); + assert_eq!( + BACKOFF_SCHEDULE, + [ + Duration::from_secs(2), + Duration::from_secs(5), + Duration::from_secs(10) + ] + ); assert_eq!(MAX_RETRIES, 3); } @@ -196,10 +203,16 @@ mod tests { #[test] fn classifies_reqwest_timeouts_as_transient() { assert!(is_transient_network_error("operation timed out")); - assert!(is_transient_network_error("error sending request: operation timed out")); + assert!(is_transient_network_error( + "error sending request: operation timed out" + )); assert!(is_transient_network_error("connection reset by peer")); - assert!(is_transient_network_error("connection refused (os error 61)")); - assert!(is_transient_network_error("dns error: failed to lookup address")); + assert!(is_transient_network_error( + "connection refused (os error 61)" + )); + assert!(is_transient_network_error( + "dns error: failed to lookup address" + )); } #[test] @@ -221,7 +234,9 @@ mod tests { assert!(is_transient_network_error( "ERROR: unable to download video: Connection timed out" )); - assert!(is_transient_network_error("Connection was closed by server")); + assert!(is_transient_network_error( + "Connection was closed by server" + )); assert!(is_transient_network_error("Timeout.")); assert!(is_transient_network_error("network is unreachable")); } @@ -234,13 +249,17 @@ mod tests { assert!(!is_transient_network_error("HTTP 403 Forbidden")); assert!(!is_transient_network_error("HTTP 410 Gone")); assert!(!is_transient_network_error("HTTP 401 Unauthorized")); - assert!(!is_transient_network_error("HTTP 451 Unavailable For Legal Reasons")); + assert!(!is_transient_network_error( + "HTTP 451 Unavailable For Legal Reasons" + )); } #[test] fn refuses_to_retry_permanent_fs_errors() { assert!(!is_transient_network_error("No space left on device")); - assert!(!is_transient_network_error("Permission denied (os error 13)")); + assert!(!is_transient_network_error( + "Permission denied (os error 13)" + )); } #[test] diff --git a/src-tauri/src/scheduler.rs b/src-tauri/src/scheduler.rs index 69d8171..a1d6fe8 100644 --- a/src-tauri/src/scheduler.rs +++ b/src-tauri/src/scheduler.rs @@ -36,18 +36,15 @@ pub fn spawn_scheduler( loop { interval.tick().await; - let settings = settings_cache - .read() - .ok() - .and_then(|settings| { - settings.as_ref().map(|settings| { - ( - settings.scheduler.clone(), - settings.scheduler_last_start_key.clone(), - settings.scheduler_last_stop_key.clone(), - ) - }) - }); + let settings = settings_cache.read().ok().and_then(|settings| { + settings.as_ref().map(|settings| { + ( + settings.scheduler.clone(), + settings.scheduler_last_start_key.clone(), + settings.scheduler_last_stop_key.clone(), + ) + }) + }); if let Some((scheduler, scheduler_last_start_key, scheduler_last_stop_key)) = settings { if !scheduler.enabled { continue; @@ -78,10 +75,13 @@ pub fn spawn_scheduler( .get("start") .is_none_or(|instant| instant.elapsed() >= Duration::from_secs(5)) { - let _ = app_handle.emit("schedule-trigger", serde_json::json!({ - "action": "start", - "key": start_key - })); + let _ = app_handle.emit( + "schedule-trigger", + serde_json::json!({ + "action": "start", + "key": start_key + }), + ); last_emit.insert("start", std::time::Instant::now()); } @@ -93,15 +93,17 @@ pub fn spawn_scheduler( &start_key, &scheduler_last_stop_key, &stop_key, - ) - && last_emit - .get("stop") - .is_none_or(|instant| instant.elapsed() >= Duration::from_secs(5)) + ) && last_emit + .get("stop") + .is_none_or(|instant| instant.elapsed() >= Duration::from_secs(5)) { - let _ = app_handle.emit("schedule-trigger", serde_json::json!({ - "action": "stop", - "key": stop_key - })); + let _ = app_handle.emit( + "schedule-trigger", + serde_json::json!({ + "action": "stop", + "key": stop_key + }), + ); last_emit.insert("stop", std::time::Instant::now()); } } diff --git a/src-tauri/src/settings.rs b/src-tauri/src/settings.rs index d3351f2..da82e4d 100644 --- a/src-tauri/src/settings.rs +++ b/src-tauri/src/settings.rs @@ -53,7 +53,10 @@ pub fn preserve_scheduler_runtime_keys( let existing_document = match decode_document(&Value::String(existing.to_string())) { Ok(doc) => doc, Err(e) => { - log::warn!("Failed to decode existing settings, dropping runtime keys: {}", e); + log::warn!( + "Failed to decode existing settings, dropping runtime keys: {}", + e + ); return Ok(incoming.to_string()); } }; @@ -147,9 +150,7 @@ fn default_category_subfolders() -> HashMap { fn normalize_category_subfolder(value: &str, fallback: &str) -> String { let parts = value .split(['/', '\\']) - .filter(|part| { - !part.is_empty() && *part != "." && *part != ".." && !part.ends_with(':') - }) + .filter(|part| !part.is_empty() && *part != "." && *part != ".." && !part.ends_with(':')) .collect::>(); if parts.is_empty() { fallback.to_string() @@ -329,14 +330,8 @@ mod tests { let merged = preserve_scheduler_runtime_keys(Some(&existing), &incoming).unwrap(); let merged: Value = serde_json::from_str(&merged).unwrap(); - assert_eq!( - merged["state"]["schedulerLastStartKey"], - "2026-06-22-start" - ); - assert_eq!( - merged["state"]["schedulerLastStopKey"], - "2026-06-22-stop" - ); + assert_eq!(merged["state"]["schedulerLastStartKey"], "2026-06-22-start"); + assert_eq!(merged["state"]["schedulerLastStopKey"], "2026-06-22-stop"); } #[test] diff --git a/src-tauri/tests/download_engine.rs b/src-tauri/tests/download_engine.rs index c3e5eb0..06ec819 100644 --- a/src-tauri/tests/download_engine.rs +++ b/src-tauri/tests/download_engine.rs @@ -9,6 +9,7 @@ use axum::{ routing::get, Router, }; +use firelink_lib::download::{DownloadCmd, DownloadCoordinator, DownloadEvent, DownloadPayload}; use futures_util::stream; use sha2::{Digest, Sha256}; use std::{ @@ -19,9 +20,12 @@ use std::{ }, time::{Duration, Instant}, }; -use firelink_lib::download::{DownloadCmd, DownloadCoordinator, DownloadEvent, DownloadPayload}; use tempfile::TempDir; -use tokio::{net::TcpListener, sync::{mpsc, oneshot}, task::JoinHandle}; +use tokio::{ + net::TcpListener, + sync::{mpsc, oneshot}, + task::JoinHandle, +}; use uuid::Uuid; const TEST_TIMEOUT: Duration = Duration::from_secs(10); diff --git a/src-tauri/tests/queue_manager.rs b/src-tauri/tests/queue_manager.rs index ba4cb4a..a45c8a6 100644 --- a/src-tauri/tests/queue_manager.rs +++ b/src-tauri/tests/queue_manager.rs @@ -1,8 +1,8 @@ use firelink_lib::queue::{ QueueManager, QueuedTask, SidecarSpawner, SpawnPayload, TaskKind, MEDIA_RUN_CANCELLED, }; -use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; use std::time::Duration; use tauri::test::{mock_builder, mock_context, noop_assets}; use tauri::Listener; @@ -84,7 +84,10 @@ async fn release_permit_is_idempotent() { mgr.release_permit("a").await; // second release: no-op let avail_after_second = mgr.available_permits(); assert_eq!(avail_after_first - avail_before, 1); - assert_eq!(avail_after_second, avail_after_first, "second release must not free another slot"); + assert_eq!( + avail_after_second, avail_after_first, + "second release must not free another slot" + ); } #[tokio::test] @@ -140,7 +143,8 @@ async fn dispatcher_parks_when_idle_no_busy_spin() { // No permit should have been acquired while idle. assert_eq!( - mgr_arc.available_permits(), 3, + mgr_arc.available_permits(), + 3, "dispatcher must not acquire permits when pending is empty" ); @@ -190,13 +194,19 @@ async fn grow_releases_immediately_and_dispatches_waiting_tasks() { // Give dispatcher time to dispatch 2 (capacity) of the 4. tokio::time::sleep(Duration::from_millis(100)).await; let native_after_initial = spawner.native_calls.load(Ordering::SeqCst); - assert_eq!(native_after_initial, 2, "only capacity-many tasks dispatch initially"); + assert_eq!( + native_after_initial, 2, + "only capacity-many tasks dispatch initially" + ); // Grow to 4; the remaining 2 should dispatch. mgr_arc.set_capacity(4); tokio::time::sleep(Duration::from_millis(100)).await; let native_after_grow = spawner.native_calls.load(Ordering::SeqCst); - assert_eq!(native_after_grow, 4, "grow must allow the waiting tasks to dispatch"); + assert_eq!( + native_after_grow, 4, + "grow must allow the waiting tasks to dispatch" + ); handle.abort(); } @@ -325,8 +335,7 @@ async fn media_terminal_error_emits_failed_without_completed() { #[tokio::test] async fn media_cancellation_does_not_emit_completed() { - let (manager, event_rx) = - make_media_manager(Err(MEDIA_RUN_CANCELLED.to_string())); + let (manager, event_rx) = make_media_manager(Err(MEDIA_RUN_CANCELLED.to_string())); let manager = Arc::new(manager); manager.push(media_task("media-cancelled")).await.unwrap(); let dispatcher = { @@ -358,14 +367,19 @@ async fn aria2_permit_survives_rpc_return() { tokio::time::sleep(Duration::from_millis(100)).await; assert_eq!(spawner.add_uri_calls.load(Ordering::SeqCst), 1); assert_eq!( - mgr_arc.available_permits(), 0, + mgr_arc.available_permits(), + 0, "permit must remain parked while aria2 download is notionally running" ); // Now simulate aria2 completion: release_permit frees the slot. mgr_arc.release_permit("a").await; tokio::time::sleep(Duration::from_millis(50)).await; - assert_eq!(mgr_arc.available_permits(), 1, "release frees the parked permit"); + assert_eq!( + mgr_arc.available_permits(), + 1, + "release frees the parked permit" + ); handle.abort(); } @@ -432,13 +446,17 @@ async fn move_up_down_reorders_pending() { mgr_arc.push(sample_task("b")).await.unwrap(); mgr_arc.push(sample_task("c")).await.unwrap(); - mgr_arc.move_in_queue("c", "main", QueueDirection::Down).await; + mgr_arc + .move_in_queue("c", "main", QueueDirection::Down) + .await; assert_eq!(mgr_arc.pending_order(None).await, vec!["a", "b", "c"]); mgr_arc.move_in_queue("c", "main", QueueDirection::Up).await; assert_eq!(mgr_arc.pending_order(None).await, vec!["a", "c", "b"]); - mgr_arc.move_in_queue("a", "main", QueueDirection::Down).await; + mgr_arc + .move_in_queue("a", "main", QueueDirection::Down) + .await; assert_eq!(mgr_arc.pending_order(None).await, vec!["c", "a", "b"]); mgr_arc.move_in_queue("c", "main", QueueDirection::Up).await; @@ -463,7 +481,10 @@ async fn moving_one_queue_does_not_reorder_another_queue() { mgr.push(a2).await.unwrap(); mgr.push(b2).await.unwrap(); - assert_eq!(mgr.move_in_queue("a2", "a", QueueDirection::Up).await, vec!["a2", "a1"]); + assert_eq!( + mgr.move_in_queue("a2", "a", QueueDirection::Up).await, + vec!["a2", "a1"] + ); assert_eq!(mgr.pending_order(Some("b")).await, vec!["b1", "b2"]); } @@ -481,17 +502,15 @@ async fn notify_fires_on_push_and_release() { }; mgr_arc.push(sample_task("x")).await.unwrap(); - let dispatched = timeout( - Duration::from_millis(150), - async { - loop { - if mgr_arc.available_permits() == 0 { - return; - } - tokio::time::sleep(Duration::from_millis(5)).await; + let dispatched = timeout(Duration::from_millis(150), async { + loop { + if mgr_arc.available_permits() == 0 { + return; } - }, - ).await; + tokio::time::sleep(Duration::from_millis(5)).await; + } + }) + .await; assert!(dispatched.is_ok(), "push must wake the idle dispatcher"); handle.abort();