diff --git a/src-tauri/src/lib_commands/cache/offline.rs b/src-tauri/src/lib_commands/cache/offline.rs index a8a41ff3..2e823212 100644 --- a/src-tauri/src/lib_commands/cache/offline.rs +++ b/src-tauri/src/lib_commands/cache/offline.rs @@ -2,24 +2,6 @@ use super::*; // ─── Offline Track Cache ────────────────────────────────────────────────────── -/// Streams an HTTP response body directly to `dest_path` in small chunks. -/// Never buffers the full file in memory — keeps RAM flat regardless of file size. -pub(crate) async fn stream_to_file(response: reqwest::Response, dest_path: &std::path::Path) -> Result<(), String> { - use futures_util::StreamExt; - use tokio::io::AsyncWriteExt; - - let mut file = tokio::fs::File::create(dest_path) - .await - .map_err(|e| e.to_string())?; - let mut stream = response.bytes_stream(); - while let Some(chunk) = stream.next().await { - let chunk = chunk.map_err(|e| e.to_string())?; - file.write_all(&chunk).await.map_err(|e| e.to_string())?; - } - file.flush().await.map_err(|e| e.to_string())?; - Ok(()) -} - pub(crate) async fn enqueue_analysis_seed(app: &tauri::AppHandle, track_id: &str, bytes: &[u8]) -> Result { if let Some(cache) = app.try_state::() { if cache.cpu_seed_redundant_for_track(track_id).unwrap_or(false) { @@ -116,26 +98,15 @@ pub(crate) async fn download_track_offline( // and released automatically when this function returns (success or error). let _permit = dl_sem.acquire().await.map_err(|e| e.to_string())?; - let client = reqwest::Client::builder() - .user_agent(subsonic_wire_user_agent()) - .timeout(std::time::Duration::from_secs(120)) - .build() - .map_err(|e| e.to_string())?; + let client = subsonic_http_client(std::time::Duration::from_secs(120))?; let response = client.get(&url).send().await.map_err(|e| e.to_string())?; if !response.status().is_success() { return Err(format!("HTTP {}", response.status().as_u16())); } - // Stream directly to a .part file; rename on success to avoid partial files. let part_path = file_path.with_extension(format!("{suffix}.part")); - if let Err(e) = stream_to_file(response, &part_path).await { - let _ = tokio::fs::remove_file(&part_path).await; - return Err(e); - } - tokio::fs::rename(&part_path, &file_path) - .await - .map_err(|e| e.to_string())?; + finalize_streamed_download(response, &file_path, &part_path).await?; enqueue_analysis_seed_from_file(&app, &track_id, &file_path).await; diff --git a/src-tauri/src/lib_commands/file_transfer.rs b/src-tauri/src/lib_commands/file_transfer.rs new file mode 100644 index 00000000..e0084ca1 --- /dev/null +++ b/src-tauri/src/lib_commands/file_transfer.rs @@ -0,0 +1,57 @@ +use std::path::Path; +use std::time::Duration; + +/// Build a reqwest client with the standard Subsonic UA and a single overall timeout. +/// For flows that need separate connect + read timeouts (long-running update/zip +/// downloads with progress events), build the client inline. +pub(crate) fn subsonic_http_client(timeout: Duration) -> Result { + reqwest::Client::builder() + .user_agent(super::subsonic_wire_user_agent()) + .timeout(timeout) + .build() + .map_err(|e| e.to_string()) +} + +/// Streams an HTTP response body to `dest_path` in chunks. Never buffers the full +/// file in memory — keeps RAM flat regardless of file size. +pub(crate) async fn stream_to_file( + response: reqwest::Response, + dest_path: &Path, +) -> Result<(), String> { + use futures_util::StreamExt; + use tokio::io::AsyncWriteExt; + + let mut file = tokio::fs::File::create(dest_path) + .await + .map_err(|e| e.to_string())?; + let mut stream = response.bytes_stream(); + while let Some(chunk) = stream.next().await { + let chunk = chunk.map_err(|e| e.to_string())?; + file.write_all(&chunk).await.map_err(|e| e.to_string())?; + } + file.flush().await.map_err(|e| e.to_string())?; + Ok(()) +} + +/// Streams `response` to `part_path`, then renames `part_path` → `dest_path`. +/// On any failure the partial `.part` file is best-effort removed so it does +/// not linger on disk. Caller must ensure `dest_path.parent()` exists. +/// +/// Note vs. previous inline implementations: the offline/device single-track +/// flows used to leave a `.part` orphan if the final rename failed. This helper +/// always cleans up, matching the batch-sync flow that already did. +pub(crate) async fn finalize_streamed_download( + response: reqwest::Response, + dest_path: &Path, + part_path: &Path, +) -> Result<(), String> { + if let Err(e) = stream_to_file(response, part_path).await { + let _ = tokio::fs::remove_file(part_path).await; + return Err(e); + } + if let Err(e) = tokio::fs::rename(part_path, dest_path).await { + let _ = tokio::fs::remove_file(part_path).await; + return Err(e.to_string()); + } + Ok(()) +} diff --git a/src-tauri/src/lib_commands/mod.rs b/src-tauri/src/lib_commands/mod.rs index 06dac471..55912f33 100644 --- a/src-tauri/src/lib_commands/mod.rs +++ b/src-tauri/src/lib_commands/mod.rs @@ -2,10 +2,12 @@ use super::*; mod app_api; mod cache; +mod file_transfer; mod sync; mod ui; pub(crate) use app_api::*; pub(crate) use cache::*; +pub(crate) use file_transfer::*; pub(crate) use sync::*; pub(crate) use ui::*; diff --git a/src-tauri/src/lib_commands/sync/batch.rs b/src-tauri/src/lib_commands/sync/batch.rs index 43406b0e..7a0433bb 100644 --- a/src-tauri/src/lib_commands/sync/batch.rs +++ b/src-tauri/src/lib_commands/sync/batch.rs @@ -137,11 +137,7 @@ pub(crate) async fn calculate_sync_payload( auth: SubsonicAuthPayload, target_dir: String, ) -> Result { - let client = reqwest::Client::builder() - .user_agent(subsonic_wire_user_agent()) - .timeout(std::time::Duration::from_secs(30)) - .build() - .map_err(|e| e.to_string())?; + let client = subsonic_http_client(std::time::Duration::from_secs(30))?; let mut add_bytes = 0; let mut add_count = 0; @@ -365,11 +361,7 @@ pub(crate) async fn sync_batch_to_device( } // Shared reqwest client — reused across all downloads. - let client = reqwest::Client::builder() - .user_agent(subsonic_wire_user_agent()) - .timeout(Duration::from_secs(300)) - .build() - .map_err(|e| e.to_string())?; + let client = subsonic_http_client(Duration::from_secs(300))?; // Concurrency limiter: max 2 parallel USB writes. let semaphore = std::sync::Arc::new(tokio::sync::Semaphore::new(2)); @@ -446,8 +438,7 @@ pub(crate) async fn sync_batch_to_device( }; let part_path = dest_path.with_extension(format!("{}.part", track.suffix)); - if let Err(e) = stream_to_file(response, &part_path).await { - let _ = tokio::fs::remove_file(&part_path).await; + if let Err(e) = finalize_streamed_download(response, &dest_path, &part_path).await { f.fetch_add(1, Ordering::Relaxed); let _ = app2.emit("device:sync:progress", serde_json::json!({ "jobId": job, "trackId": track.id, "status": "error", @@ -455,15 +446,6 @@ pub(crate) async fn sync_batch_to_device( })); return; } - if let Err(e) = tokio::fs::rename(&part_path, &dest_path).await { - let _ = tokio::fs::remove_file(&part_path).await; - f.fetch_add(1, Ordering::Relaxed); - let _ = app2.emit("device:sync:progress", serde_json::json!({ - "jobId": job, "trackId": track.id, "status": "error", - "error": e.to_string(), - })); - return; - } d.fetch_add(1, Ordering::Relaxed); status = "done"; diff --git a/src-tauri/src/lib_commands/sync/device.rs b/src-tauri/src/lib_commands/sync/device.rs index fca83df4..69ddb4f6 100644 --- a/src-tauri/src/lib_commands/sync/device.rs +++ b/src-tauri/src/lib_commands/sync/device.rs @@ -345,11 +345,7 @@ pub(crate) async fn sync_track_to_device( .map_err(|e| e.to_string())?; } - let client = reqwest::Client::builder() - .user_agent(subsonic_wire_user_agent()) - .timeout(std::time::Duration::from_secs(300)) - .build() - .map_err(|e| e.to_string())?; + let client = subsonic_http_client(std::time::Duration::from_secs(300))?; let response = client.get(&track.url).send().await.map_err(|e| e.to_string())?; if !response.status().is_success() { @@ -361,16 +357,12 @@ pub(crate) async fn sync_track_to_device( } let part_path = dest_path.with_extension(format!("{}.part", track.suffix)); - if let Err(e) = stream_to_file(response, &part_path).await { - let _ = tokio::fs::remove_file(&part_path).await; + if let Err(e) = finalize_streamed_download(response, &dest_path, &part_path).await { let _ = app.emit("device:sync:progress", serde_json::json!({ "jobId": job_id, "trackId": track.id, "status": "error", "error": e, })); return Err(e); } - tokio::fs::rename(&part_path, &dest_path) - .await - .map_err(|e| e.to_string())?; let _ = app.emit("device:sync:progress", serde_json::json!({ "jobId": job_id, "trackId": track.id, "status": "done", "path": path_str,