From e07182fbf26acc80756ba9622f0d1e18873f2daf Mon Sep 17 00:00:00 2001 From: NimBold Date: Tue, 14 Jul 2026 18:29:14 +0330 Subject: [PATCH] fix(downloads): harden queue and lifecycle synchronization Serialize queue controls, make multi-item moves atomic, await stale enqueue cleanup, and guard late media and progress events. Add deterministic table sorting and regression coverage for worst-case lifecycle races. --- src-tauri/src/download.rs | 118 +++++++-- src-tauri/src/lib.rs | 37 ++- src-tauri/src/queue.rs | 211 +++++++++++++-- src-tauri/tests/queue_manager.rs | 34 ++- src/App.tsx | 10 +- src/components/DownloadTable.tsx | 253 +++++++++--------- src/components/Sidebar.tsx | 27 +- src/ipc.ts | 1 + src/store/downloadStore.test.ts | 38 +++ src/store/downloadStore.ts | 11 +- src/store/useDownloadStore.test.ts | 137 +++++++++- src/store/useDownloadStore.ts | 343 ++++++++++++++++--------- src/utils/downloadTableSorting.test.ts | 60 +++++ src/utils/downloadTableSorting.ts | 110 ++++++++ src/utils/downloads.ts | 5 + 15 files changed, 1074 insertions(+), 321 deletions(-) create mode 100644 src/utils/downloadTableSorting.test.ts create mode 100644 src/utils/downloadTableSorting.ts diff --git a/src-tauri/src/download.rs b/src-tauri/src/download.rs index f27eeb6..683dc4f 100644 --- a/src-tauri/src/download.rs +++ b/src-tauri/src/download.rs @@ -46,18 +46,29 @@ impl DownloadCoordinator { .map_err(|_| "download coordinator is unavailable".to_string()) } - pub async fn register_media(&self, id: String) -> Result, String> { + pub async fn register_media( + &self, + id: String, + lifecycle_generation: u64, + ) -> Result, String> { let (cancel_tx, cancel_rx) = watch::channel(false); self.media_tx - .send(MediaCmd::Register { id, cancel_tx }) + .send(MediaCmd::Register { + id, + lifecycle_generation, + cancel_tx, + }) .await .map_err(|_| "download coordinator is unavailable".to_string())?; Ok(cancel_rx) } - pub async fn pause_media(&self, id: String) -> Result<(), String> { + pub async fn pause_media(&self, id: String, lifecycle_generation: u64) -> Result<(), String> { self.media_tx - .send(MediaCmd::Pause(id)) + .send(MediaCmd::Pause { + id, + lifecycle_generation, + }) .await .map_err(|_| "download coordinator is unavailable".to_string()) } @@ -65,16 +76,27 @@ impl DownloadCoordinator { pub async fn pause_media_with_ack( &self, id: String, + lifecycle_generation: u64, ack: tokio::sync::oneshot::Sender<()>, ) -> Result<(), String> { self.media_tx - .send(MediaCmd::PauseWithAck(id, ack)) + .send(MediaCmd::PauseWithAck { + id, + lifecycle_generation, + ack, + }) .await .map_err(|_| "download coordinator is unavailable".to_string()) } - pub async fn finish_media(&self, id: String) { - let _ = self.media_tx.send(MediaCmd::Finished(id)).await; + pub async fn finish_media(&self, id: String, lifecycle_generation: u64) { + let _ = self + .media_tx + .send(MediaCmd::Finished { + id, + lifecycle_generation, + }) + .await; } } @@ -96,11 +118,22 @@ impl CoordinatorEventSink { enum MediaCmd { Register { id: String, + lifecycle_generation: u64, cancel_tx: watch::Sender, }, - Pause(String), - PauseWithAck(String, tokio::sync::oneshot::Sender<()>), - Finished(String), + Pause { + id: String, + lifecycle_generation: u64, + }, + PauseWithAck { + id: String, + lifecycle_generation: u64, + ack: tokio::sync::oneshot::Sender<()>, + }, + Finished { + id: String, + lifecycle_generation: u64, + }, } async fn run_coordinator( @@ -108,8 +141,8 @@ async fn run_coordinator( mut command_rx: mpsc::Receiver, mut media_rx: mpsc::Receiver, ) { - let mut active_media = HashMap::>::new(); - let mut pending_media_acks = HashMap::>::new(); + let mut active_media = HashMap::)>::new(); + let mut pending_media_acks = HashMap::)>::new(); let mut pending_captured_urls = Vec::::new(); let mut frontend_ready = false; @@ -146,28 +179,36 @@ async fn run_coordinator( continue; }; match command { - MediaCmd::Register { id, cancel_tx } => { - if let Some(previous) = active_media.insert(id, cancel_tx) { + MediaCmd::Register { id, lifecycle_generation, cancel_tx } => { + if let Some((_, previous)) = active_media.insert(id, (lifecycle_generation, cancel_tx)) { let _ = previous.send(true); } } - MediaCmd::Pause(id) => { - if let Some(cancel_tx) = active_media.remove(&id) { - let _ = cancel_tx.send(true); + MediaCmd::Pause { id, lifecycle_generation } => { + if active_media.get(&id).is_some_and(|(generation, _)| *generation == lifecycle_generation) { + if let Some((_, cancel_tx)) = active_media.remove(&id) { + let _ = cancel_tx.send(true); + } } } - MediaCmd::PauseWithAck(id, ack) => { - if let Some(cancel_tx) = active_media.remove(&id) { - let _ = cancel_tx.send(true); - pending_media_acks.insert(id, ack); + MediaCmd::PauseWithAck { id, lifecycle_generation, ack } => { + if active_media.get(&id).is_some_and(|(generation, _)| *generation == lifecycle_generation) { + if let Some((_, cancel_tx)) = active_media.remove(&id) { + let _ = cancel_tx.send(true); + pending_media_acks.insert(id, (lifecycle_generation, ack)); + } } else { let _ = ack.send(()); } } - MediaCmd::Finished(id) => { - active_media.remove(&id); - if let Some(ack) = pending_media_acks.remove(&id) { - let _ = ack.send(()); + MediaCmd::Finished { id, lifecycle_generation } => { + if active_media.get(&id).is_some_and(|(generation, _)| *generation == lifecycle_generation) { + active_media.remove(&id); + } + if pending_media_acks.get(&id).is_some_and(|(generation, _)| *generation == lifecycle_generation) { + if let Some((_, ack)) = pending_media_acks.remove(&id) { + let _ = ack.send(()); + } } } } @@ -175,7 +216,7 @@ async fn run_coordinator( } } - for (_, cancel_tx) in active_media { + for (_, (_, cancel_tx)) in active_media { let _ = cancel_tx.send(true); } } @@ -251,4 +292,29 @@ mod tests { DownloadEvent::CapturedUrls("https://example.com/startup.zip".to_string()) ); } + + #[tokio::test] + async fn stale_media_finish_cannot_remove_a_newer_lifecycle() { + let (coordinator, _events) = DownloadCoordinator::spawn_headless(); + let mut old_cancel = coordinator.register_media("same-id".to_string(), 1).await.unwrap(); + let mut new_cancel = coordinator.register_media("same-id".to_string(), 2).await.unwrap(); + + tokio::time::timeout(Duration::from_secs(1), old_cancel.changed()) + .await + .unwrap() + .unwrap(); + coordinator.finish_media("same-id".to_string(), 1).await; + + let (ack_tx, ack_rx) = tokio::sync::oneshot::channel(); + coordinator + .pause_media_with_ack("same-id".to_string(), 2, ack_tx) + .await + .unwrap(); + coordinator.finish_media("same-id".to_string(), 2).await; + tokio::time::timeout(Duration::from_secs(1), ack_rx) + .await + .unwrap() + .unwrap(); + assert!(*new_cancel.borrow_and_update()); + } } diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index f680f74..2f24580 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -3357,6 +3357,11 @@ async fn pause_download( let _control_guard = state.queue_manager.acquire_aria2_control(&id).await; let active_kind = state.queue_manager.active_kind(&id).await; + let media_lifecycle_generation = state + .queue_manager + .registered_lifecycle_generation(&id) + .await + .unwrap_or_default(); let removed_pending = state.queue_manager.remove_from_pending(&id).await; let gid = state.queue_manager.aria2_gid_for_download(&id); @@ -3447,7 +3452,7 @@ async fn pause_download( if matches!(active_kind, Some(crate::queue::TaskKind::Media)) { state .download_coordinator - .pause_media_with_ack(id.clone(), tx) + .pause_media_with_ack(id.clone(), media_lifecycle_generation, tx) .await?; } else { let _ = tx.send(()); @@ -3498,6 +3503,7 @@ async fn resume_download( "paused" => { let control_epoch = state.queue_manager.next_aria2_control_epoch(&id).await; state.queue_manager.allow_aria2_retries(&id).await; + state.queue_manager.reset_aria2_retry_strikes(&id).await; use tauri::Emitter; let _ = app_handle.emit( "download-state", @@ -3713,6 +3719,11 @@ async fn remove_download( let _control_guard = state.queue_manager.acquire_aria2_control(&id).await; let active_kind = state.queue_manager.active_kind(&id).await; + let media_lifecycle_generation = state + .queue_manager + .registered_lifecycle_generation(&id) + .await + .unwrap_or_default(); state.queue_manager.remove_from_pending(&id).await; state.queue_manager.next_aria2_control_epoch(&id).await; @@ -3748,7 +3759,7 @@ async fn remove_download( if matches!(active_kind, Some(crate::queue::TaskKind::Media)) { state .download_coordinator - .pause_media_with_ack(id.clone(), tx) + .pause_media_with_ack(id.clone(), media_lifecycle_generation, tx) .await?; } else { let _ = tx.send(()); @@ -3864,6 +3875,11 @@ async fn detach_download_for_reconfigure( log::info!("detach_download_for_reconfigure called for id: {}", id); let _control_guard = state.queue_manager.acquire_aria2_control(&id).await; let active_kind = state.queue_manager.active_kind(&id).await; + let media_lifecycle_generation = state + .queue_manager + .registered_lifecycle_generation(&id) + .await + .unwrap_or_default(); state.queue_manager.remove_from_pending(&id).await; state.queue_manager.next_aria2_control_epoch(&id).await; state.queue_manager.cancel_aria2_retries(&id).await; @@ -3908,7 +3924,7 @@ async fn detach_download_for_reconfigure( if matches!(active_kind, Some(crate::queue::TaskKind::Media)) { state .download_coordinator - .pause_media_with_ack(id.clone(), tx) + .pause_media_with_ack(id.clone(), media_lifecycle_generation, tx) .await?; } else { let _ = tx.send(()); // Fallback if no task exists @@ -4461,6 +4477,19 @@ async fn move_in_queue( .await) } +#[tauri::command] +async fn move_many_in_queue( + state: tauri::State<'_, AppState>, + ids: Vec, + queue_id: String, + direction: crate::ipc::QueueDirection, +) -> Result, AppError> { + Ok(state + .queue_manager + .move_many_in_queue(&ids, &queue_id, direction) + .await) +} + #[tauri::command] async fn remove_from_queue( app_handle: tauri::AppHandle, @@ -6717,7 +6746,7 @@ pub fn run() { check_file_exists, toggle_tray_icon, set_extension_pairing_token, get_extension_server_port, set_extension_frontend_ready, set_concurrent_limit, set_global_speed_limit, remove_download, detach_download_for_reconfigure, - enqueue_download, enqueue_many, cancel_enqueue_generation, move_in_queue, remove_from_queue, get_pending_order, + enqueue_download, enqueue_many, cancel_enqueue_generation, move_in_queue, move_many_in_queue, remove_from_queue, get_pending_order, commands::reveal_in_file_manager, commands::open_downloaded_file, parity::get_system_proxy, parity::get_file_category, parity::check_for_updates, parity::is_supported_media, parity::get_supported_media_domains, parity::create_category_directories, diff --git a/src-tauri/src/queue.rs b/src-tauri/src/queue.rs index f0466e9..faf1718 100644 --- a/src-tauri/src/queue.rs +++ b/src-tauri/src/queue.rs @@ -73,6 +73,9 @@ pub struct QueuedTask { pub queue_id: String, pub kind: TaskKind, pub payload: SpawnPayload, + /// Frontend lifecycle generation that owns this sidecar and its permit. + /// Manual internal/test tasks use generation 0. + pub lifecycle_generation: u64, } /// Args mirroring start_download / start_media_download. Kept untyped-loose @@ -120,17 +123,24 @@ pub trait SidecarSpawner: Send + Sync + 'static { /// Run a media download to completion. The permit is parked for the full /// duration; release is handled by QueueManager on the runner's exit. - async fn run_media(&self, id: &str, payload: &SpawnPayload) -> Result<(), String>; + async fn run_media( + &self, + id: &str, + payload: &SpawnPayload, + lifecycle_generation: u64, + ) -> Result<(), String>; } /// The centralized concurrency gatekeeper. One instance lives in AppState. pub struct QueueManager { registered_ids: Mutex>, + registered_lifecycle_generations: Mutex>, enqueue_cancellations: Mutex>, enqueue_generations: Mutex>, pending: Mutex>, semaphore: Arc, active_permits: Mutex>, + active_permit_generations: Mutex>, active_kinds: Mutex>, target_capacity: AtomicUsize, slots_to_retire: AtomicUsize, @@ -197,11 +207,13 @@ impl QueueManager { ) -> Self { Self { registered_ids: Mutex::new(HashSet::new()), + registered_lifecycle_generations: Mutex::new(HashMap::new()), enqueue_cancellations: Mutex::new(HashMap::new()), enqueue_generations: Mutex::new(HashMap::new()), pending: Mutex::new(VecDeque::new()), semaphore: Arc::new(Semaphore::new(capacity)), active_permits: Mutex::new(HashMap::new()), + active_permit_generations: Mutex::new(HashMap::new()), active_kinds: Mutex::new(HashMap::new()), target_capacity: AtomicUsize::new(capacity), slots_to_retire: AtomicUsize::new(0), @@ -237,6 +249,7 @@ impl QueueManager { /// Explicitly release a backend registry id (e.g. on un-resumable false paths, removals, or detach). pub async fn release_registered_id(&self, id: &str) { self.registered_ids.lock().await.remove(id); + self.registered_lifecycle_generations.lock().await.remove(id); // A released lifecycle cannot be resumed by a delayed retry worker. // Epoch checks remain the authoritative guard; removing this marker // prevents terminal downloads from accumulating cancellation entries. @@ -247,6 +260,40 @@ impl QueueManager { self.registered_ids.lock().await.contains(id) } + pub async fn registered_lifecycle_generation(&self, id: &str) -> Option { + self.registered_lifecycle_generations + .lock() + .await + .get(id) + .copied() + } + + async fn is_registered_generation(&self, id: &str, generation: u64) -> bool { + self.registered_lifecycle_generations + .lock() + .await + .get(id) + .copied() + == Some(generation) + } + + async fn release_registered_id_for_generation(&self, id: &str, generation: u64) { + let released = { + let mut registered = self.registered_ids.lock().await; + let mut generations = self.registered_lifecycle_generations.lock().await; + if generations.get(id).copied() == Some(generation) { + registered.remove(id); + generations.remove(id); + true + } else { + false + } + }; + if released { + self.aria2_retry_cancelled.lock().await.remove(id); + } + } + /// Reject an in-flight enqueue generation if a newer UI action supersedes it. pub async fn cancel_enqueue_generation(&self, id: &str, generation: u64) { let mut cancellations = self.enqueue_cancellations.lock().await; @@ -282,6 +329,10 @@ impl QueueManager { return Err("Duplicate task".to_string()); } registered.insert(id.to_string()); + self.registered_lifecycle_generations + .lock() + .await + .insert(id.to_string(), generation); generations.insert(id.to_string(), generation); Ok(previous_generation) } @@ -298,6 +349,7 @@ impl QueueManager { return; } registered.remove(id); + self.registered_lifecycle_generations.lock().await.remove(id); match previous_generation { Some(previous) => { generations.insert(id.to_string(), previous); @@ -459,6 +511,13 @@ impl QueueManager { .insert(id.to_string(), permit); } + async fn tag_permit_generation(&self, id: &str, generation: u64) { + self.active_permit_generations + .lock() + .await + .insert(id.to_string(), generation); + } + pub async fn active_kind(&self, id: &str) -> Option { self.active_kinds.lock().await.get(id).cloned() } @@ -492,12 +551,30 @@ impl QueueManager { pub async fn release_permit(&self, id: &str) { let removed = self.active_permits.lock().await.remove(id).is_some(); + self.active_permit_generations.lock().await.remove(id); self.active_kinds.lock().await.remove(id); if removed { self.notify.notify_one(); } } + async fn release_permit_for_generation(&self, id: &str, generation: u64) { + let removed = { + let mut permits = self.active_permits.lock().await; + let mut generations = self.active_permit_generations.lock().await; + if generations.get(id).copied() == Some(generation) { + generations.remove(id); + permits.remove(id).is_some() + } else { + false + } + }; + if removed { + self.active_kinds.lock().await.remove(id); + self.notify.notify_one(); + } + } + pub async fn has_active_permit(&self, id: &str) -> bool { self.active_permits.lock().await.contains_key(id) } @@ -607,10 +684,12 @@ impl QueueManager { async fn dispatch_one(self: Arc, permit: OwnedSemaphorePermit, task: QueuedTask) { let id = task.id.clone(); + let lifecycle_generation = task.lifecycle_generation; // Park the permit BEFORE spawning. Uniform parking: // aria2's RPC returns instantly, so the permit must outlive the // dispatch_one call. Media runners release on exit. self.park_permit(&id, permit).await; + self.tag_permit_generation(&id, lifecycle_generation).await; self.active_kinds .lock() .await @@ -693,30 +772,55 @@ impl QueueManager { let payload = task.payload.clone(); let id_for_task = id.clone(); tauri::async_runtime::spawn(async move { - let outcome = this.spawner.run_media(&id_for_task, &payload).await; - this.finish_runner(&id_for_task, outcome).await; + let outcome = this + .spawner + .run_media(&id_for_task, &payload, lifecycle_generation) + .await; + this.finish_runner(&id_for_task, lifecycle_generation, outcome) + .await; }); } } } /// Terminal handler for non-aria2 transfers. Emits state and frees the permit. - /// Does not emit or release anything on intentional MEDIA_RUN_CANCELLED. + /// Intentional cancellation is silent, but still releases backend ownership. /// Note: `id` is the frontend download UUID, which survives indefinitely as /// the terminal state. - async fn finish_runner(self: Arc, id: &str, outcome: Result<(), String>) { + async fn finish_runner( + self: Arc, + id: &str, + lifecycle_generation: u64, + outcome: Result<(), String>, + ) { + if !self.is_registered_generation(id, lifecycle_generation).await { + log::info!( + "media runner [{}]: ignoring stale terminal outcome for lifecycle {}", + id, + lifecycle_generation + ); + self.release_permit_for_generation(id, lifecycle_generation).await; + return; + } + match outcome { Ok(()) => { self.emit_state(id, DownloadStatus::Completed); - self.release_registered_id(id).await; + self.release_registered_id_for_generation(id, lifecycle_generation) + .await; + } + Err(error) if error == MEDIA_RUN_CANCELLED => { + self.release_registered_id_for_generation(id, lifecycle_generation) + .await; } - Err(error) if error == MEDIA_RUN_CANCELLED => {} Err(error) => { self.emit_failed(id, error); - self.release_registered_id(id).await; + self.release_registered_id_for_generation(id, lifecycle_generation) + .await; } } - self.release_permit(id).await; + self.release_permit_for_generation(id, lifecycle_generation) + .await; } fn emit_failed(&self, id: &str, error: String) { @@ -852,6 +956,13 @@ impl QueueManager { self.aria2_retry_cancelled.lock().await.remove(id); } + /// A user-initiated resume starts a new retry lifecycle while retaining + /// the paused Aria2 payload/GID. Reset only the strike budget here; the + /// payload is still needed for the same-GID resume path. + pub async fn reset_aria2_retry_strikes(&self, id: &str) { + self.aria2_retry_strikes.lock().await.remove(id); + } + async fn finish_aria2_retry(&self, id: &str, gid: &str, retry_epoch: u64) { self.release_aria2_retry_inflight(id, retry_epoch).await; self.aria2_retrying_gids.lock().await.remove(gid); @@ -1276,6 +1387,20 @@ impl QueueManager { id: &str, queue_id: &str, direction: QueueDirection, + ) -> Vec { + self.move_many_in_queue(&[id.to_string()], queue_id, direction) + .await + } + + /// Atomically move a selected block of pending tasks up or down. The + /// frontend uses the same block semantics for multi-selection, so keeping + /// the operation under one pending-list lock prevents partial RPC moves + /// from leaving the backend in a different order than the UI. + pub async fn move_many_in_queue( + &self, + ids: &[String], + queue_id: &str, + direction: QueueDirection, ) -> Vec { let mut pending = self.pending.lock().await; let queue_positions = pending @@ -1283,22 +1408,46 @@ impl QueueManager { .enumerate() .filter_map(|(index, task)| (task.queue_id == queue_id).then_some(index)) .collect::>(); - let queue_pos = queue_positions + let selected_positions = queue_positions .iter() - .position(|index| pending[*index].id == id); - if let Some(queue_pos) = queue_pos { - let target = match direction { - QueueDirection::Up => queue_pos.checked_sub(1), - QueueDirection::Down => { - if queue_pos + 1 < queue_positions.len() { - Some(queue_pos + 1) - } else { - None - } - } + .enumerate() + .filter_map(|(position, index)| ids.iter().any(|id| id == &pending[*index].id).then_some(position)) + .collect::>(); + + if !selected_positions.is_empty() { + let queue_tasks = queue_positions + .iter() + .map(|index| pending[*index].clone()) + .collect::>(); + let selected_ids = ids.iter().collect::>(); + let selected_tasks = queue_tasks + .iter() + .filter(|task| selected_ids.contains(&task.id)) + .cloned() + .collect::>(); + let unselected_tasks = queue_tasks + .iter() + .filter(|task| !selected_ids.contains(&task.id)) + .cloned() + .collect::>(); + + let first_selected = *selected_positions.first().unwrap(); + let last_selected = *selected_positions.last().unwrap(); + let selected_count = selected_tasks.len(); + let insert_index = match direction { + QueueDirection::Up => first_selected.saturating_sub(1), + QueueDirection::Down => (last_selected + 1) + .saturating_sub(selected_count) + .saturating_add(1) + .min(unselected_tasks.len()), }; - if let Some(target) = target { - pending.swap(queue_positions[queue_pos], queue_positions[target]); + let mut reordered = Vec::with_capacity(queue_tasks.len()); + reordered.extend_from_slice(&unselected_tasks[..insert_index]); + reordered.extend(selected_tasks); + reordered.extend_from_slice(&unselected_tasks[insert_index..]); + + for (queue_index, pending_index) in queue_positions.iter().enumerate() { + pending[*pending_index] = reordered[queue_index].clone(); } } pending @@ -1740,11 +1889,16 @@ impl SidecarSpawner for ProductionSpawner { crate::ensure_aria2_gid_result("unpause", gid, &resumed) } - async fn run_media(&self, id: &str, payload: &SpawnPayload) -> Result<(), String> { + async fn run_media( + &self, + id: &str, + payload: &SpawnPayload, + lifecycle_generation: u64, + ) -> Result<(), String> { let state = self.app_handle.state::(); let mut cancel_rx = state .download_coordinator - .register_media(id.to_string()) + .register_media(id.to_string(), lifecycle_generation) .await?; let outcome = crate::start_media_download_internal( self.app_handle.clone(), @@ -1777,7 +1931,7 @@ impl SidecarSpawner for ProductionSpawner { } let _ = state .download_coordinator - .finish_media(id.to_string()) + .finish_media(id.to_string(), lifecycle_generation) .await; outcome.map(|_| ()) } @@ -1823,6 +1977,11 @@ impl EnqueueItem { id, queue_id: self.queue_id, kind, + lifecycle_generation: self + .lifecycle_generation + .as_deref() + .and_then(|generation| generation.parse().ok()) + .unwrap_or_default(), payload: SpawnPayload { url: self.url, destination: self.destination, diff --git a/src-tauri/tests/queue_manager.rs b/src-tauri/tests/queue_manager.rs index 6c2eb40..3426a91 100644 --- a/src-tauri/tests/queue_manager.rs +++ b/src-tauri/tests/queue_manager.rs @@ -63,7 +63,7 @@ impl SidecarSpawner for DelayedAria2Spawner { Ok(()) } - async fn run_media(&self, _id: &str, _payload: &SpawnPayload) -> Result<(), String> { + async fn run_media(&self, _id: &str, _payload: &SpawnPayload, _generation: u64) -> Result<(), String> { unreachable!("media is not used by delayed aria2 tests") } } @@ -86,7 +86,7 @@ impl SidecarSpawner for FailFirstAria2Spawner { Ok(()) } - async fn run_media(&self, _id: &str, _payload: &SpawnPayload) -> Result<(), String> { + async fn run_media(&self, _id: &str, _payload: &SpawnPayload, _generation: u64) -> Result<(), String> { unreachable!("media is not used by fail-first aria2 tests") } } @@ -109,7 +109,7 @@ impl firelink_lib::queue::SidecarSpawner for CountingSpawner { async fn remove_uri(&self, _gid: &str) -> Result<(), String> { Ok(()) } - async fn run_media(&self, _id: &str, _payload: &SpawnPayload) -> Result<(), String> { + async fn run_media(&self, _id: &str, _payload: &SpawnPayload, _generation: u64) -> Result<(), String> { self.media_calls.fetch_add(1, Ordering::SeqCst); Ok(()) } @@ -131,6 +131,7 @@ fn sample_task(id: &str) -> QueuedTask { id: id.to_string(), queue_id: "main".to_string(), kind: TaskKind::Aria2, + lifecycle_generation: 0, payload: SpawnPayload::default(), } } @@ -498,6 +499,7 @@ fn aria2_task(id: &str) -> QueuedTask { id: id.to_string(), queue_id: "main".to_string(), kind: TaskKind::Aria2, + lifecycle_generation: 0, payload: SpawnPayload::default(), } } @@ -507,6 +509,7 @@ fn media_task(id: &str) -> QueuedTask { id: id.to_string(), queue_id: "main".to_string(), kind: TaskKind::Media, + lifecycle_generation: 0, payload: SpawnPayload::default(), } } @@ -525,7 +528,7 @@ impl SidecarSpawner for FixedMediaSpawner { unreachable!("aria2 is not used by media terminal-state tests") } - async fn run_media(&self, _id: &str, _payload: &SpawnPayload) -> Result<(), String> { + async fn run_media(&self, _id: &str, _payload: &SpawnPayload, _generation: u64) -> Result<(), String> { self.outcome.clone() } } @@ -590,6 +593,7 @@ async fn media_cancellation_does_not_emit_completed() { assert!(!statuses.iter().any(|status| status == "failed")); assert!(!statuses.iter().any(|status| status == "completed")); assert_eq!(manager.available_permits(), 1); + assert!(!manager.is_registered("media-cancelled").await); dispatcher.abort(); } @@ -1079,6 +1083,28 @@ async fn move_up_down_reorders_pending() { assert_eq!(mgr_arc.pending_order(None).await, vec!["c", "a", "b"]); } +#[tokio::test] +async fn multi_move_reorders_selected_items_as_one_atomic_block() { + use firelink_lib::ipc::QueueDirection; + + let (mgr, _spawner) = make_manager(3); + for id in ["a", "b", "c", "d", "e"] { + mgr.push(sample_task(id)).await.unwrap(); + } + + let selected = vec!["b".to_string(), "d".to_string()]; + assert_eq!( + mgr.move_many_in_queue(&selected, "main", QueueDirection::Up) + .await, + vec!["b", "d", "a", "c", "e"] + ); + assert_eq!( + mgr.move_many_in_queue(&selected, "main", QueueDirection::Down) + .await, + vec!["a", "b", "d", "c", "e"] + ); +} + #[tokio::test] async fn moving_one_queue_does_not_reorder_another_queue() { use firelink_lib::ipc::QueueDirection; diff --git a/src/App.tsx b/src/App.tsx index 53c919d..91d44e9 100644 --- a/src/App.tsx +++ b/src/App.tsx @@ -1,4 +1,4 @@ -import { initMediaDomains, isActiveDownloadStatus, normalizeSpeedLimitForBackend } from './utils/downloads'; +import { initMediaDomains, isActiveDownloadStatus, isTransferActiveStatus, normalizeSpeedLimitForBackend } from './utils/downloads'; import { schedulerCompletionState } from './utils/schedulerCompletion'; import { useCallback, useEffect, useRef, useState } from "react"; import { Sidebar, SidebarFilter } from "./components/Sidebar"; @@ -112,7 +112,7 @@ function App() { const logsEnabled = useSettingsStore(state => state.logsEnabled); const extensionPairingToken = useSettingsStore(state => state.extensionPairingToken); const downloads = useDownloadStore(state => state.downloads); - const activeDownloadCount = downloads.filter(download => download.status === 'downloading').length; + const activeDownloadCount = downloads.filter(download => isTransferActiveStatus(download.status)).length; const queuedCount = downloads.filter(download => download.status === 'queued' || download.status === 'staged' ).length; @@ -124,11 +124,7 @@ function App() { const pendingPostActionTimer = useRef(null); const maxConcurrentDownloads = useSettingsStore(state => state.maxConcurrentDownloads); const preventsSleepWhileDownloading = useSettingsStore(state => state.preventsSleepWhileDownloading); - const activeTransferCount = downloads.filter(download => - download.status === 'downloading' || - download.status === 'processing' || - download.status === 'retrying' - ).length; + const activeTransferCount = downloads.filter(download => isTransferActiveStatus(download.status)).length; const { addToast } = useToast(); const isMacUserAgent = navigator.userAgent.includes('Mac'); const usesCustomWindowControls = !isMacUserAgent && platform.os !== 'macos'; diff --git a/src/components/DownloadTable.tsx b/src/components/DownloadTable.tsx index d6e82bc..247409f 100644 --- a/src/components/DownloadTable.tsx +++ b/src/components/DownloadTable.tsx @@ -1,4 +1,4 @@ -import React, { useState, useEffect, useRef } from 'react'; +import React, { useState, useEffect, useMemo, useRef } from 'react'; import { useDownloadStore, DownloadItem } from '../store/useDownloadStore'; import { useToast } from '../contexts/ToastContext'; import { useSettingsStore } from '../store/useSettingsStore'; @@ -17,8 +17,13 @@ import { canStartDownload, startActionLabel } from '../utils/downloadActions'; -import { isActiveDownloadStatus } from '../utils/downloads'; +import { isActiveDownloadStatus, isTransferActiveStatus } from '../utils/downloads'; import { readClipboardDownloadUrls } from '../utils/clipboard'; +import { + sortDownloads, + type DownloadSortColumn, + type DownloadSortConfig +} from '../utils/downloadTableSorting'; interface DownloadTableProps { filter: SidebarFilter; @@ -46,7 +51,11 @@ export const DownloadTable: React.FC = ({ filter }) => { const [animationParent] = useAutoAnimate(); const [selectedIds, setSelectedIds] = useState>(new Set()); const [lastSelectedId, setLastSelectedId] = useState(null); - const [sortConfig, setSortConfig] = useState<{ column: string; direction: 'asc' | 'desc' } | null>({ column: 'Date Added', direction: 'desc' }); + const [sortConfig, setSortConfig] = useState({ column: 'Date Added', direction: 'desc' }); + const [queueSortConfig, setQueueSortConfig] = useState(null); + const selectedIdsRef = useRef(selectedIds); + const sortedDownloadsRef = useRef([]); + selectedIdsRef.current = selectedIds; const [columnWidths, setColumnWidths] = useState(() => { try { const stored = JSON.parse(window.localStorage.getItem(COLUMN_WIDTHS_STORAGE_KEY) || 'null'); @@ -110,15 +119,15 @@ export const DownloadTable: React.FC = ({ filter }) => { if (!isInput) { if ((e.key === 'a' || e.key === 'A') && (e.metaKey || e.ctrlKey)) { e.preventDefault(); - const allIds = sortedDownloads.map(d => d.id); + const allIds = sortedDownloadsRef.current.map(d => d.id); setSelectedIds(new Set(allIds)); return; } if (e.key === 'Delete' || e.key === 'Backspace') { if (!activeEl || !activeEl.closest('.sidebar-inner')) { - if (selectedIds.size > 0) { - handleDelete(Array.from(selectedIds)); + if (selectedIdsRef.current.size > 0) { + handleDelete(Array.from(selectedIdsRef.current)); } } } @@ -126,7 +135,7 @@ export const DownloadTable: React.FC = ({ filter }) => { }; window.addEventListener('keydown', handleKeyDown); return () => window.removeEventListener('keydown', handleKeyDown); - }, [selectedIds]); + }, []); const showInteractionError = (message: string, error: unknown) => { @@ -192,43 +201,24 @@ export const DownloadTable: React.FC = ({ filter }) => { openProperties(item.id); }; - const filteredDownloads = downloads.filter((d: DownloadItem) => { - if (filter.startsWith('queue:')) { + const isQueueFilter = filter.startsWith('queue:'); + const filteredDownloads = useMemo(() => downloads.filter((d: DownloadItem) => { + if (isQueueFilter) { return d.queueId === filter.replace('queue:', '') && d.status !== 'completed'; } switch (filter) { case 'all': return true; - case 'active': return d.status === 'downloading'; + case 'active': return isTransferActiveStatus(d.status); case 'completed': return d.status === 'completed'; case 'unfinished': return d.status !== 'completed'; default: return d.category === filter; } - }); + }), [downloads, filter, isQueueFilter]); - const parseSpeed = (speedStr?: string) => { - if (!speedStr || speedStr === '-') return 0; - const val = parseFloat(speedStr); - if (speedStr.includes('KB/s')) return val * 1024; - if (speedStr.includes('MB/s')) return val * 1024 * 1024; - if (speedStr.includes('GB/s')) return val * 1024 * 1024 * 1024; - return val; - }; - - const parseEta = (etaStr?: string) => { - if (!etaStr || etaStr === '-') return Infinity; - let seconds = 0; - const hours = etaStr.match(/(\d+)h/); - const minutes = etaStr.match(/(\d+)m/); - const secs = etaStr.match(/(\d+)s/); - if (hours) seconds += parseInt(hours[1]) * 3600; - if (minutes) seconds += parseInt(minutes[1]) * 60; - if (secs) seconds += parseInt(secs[1]); - return seconds; - }; - - // Sort by queue position when viewing a specific queue so the visual - // order matches the queue order and move-up/down buttons reflect reality. - const sortedDownloads = filter.startsWith('queue:') + // Queue views use the persisted queue order until the user explicitly sorts + // a column. This keeps move-up/down controls truthful while still making + // every header a working sort target. + const sortedDownloads = useMemo(() => isQueueFilter && !queueSortConfig ? [...filteredDownloads].sort((left, right) => { const leftActive = isActiveDownloadStatus(left.status) && left.status !== 'queued'; const rightActive = isActiveDownloadStatus(right.status) && right.status !== 'queued'; @@ -236,40 +226,26 @@ export const DownloadTable: React.FC = ({ filter }) => { if (leftActive && !rightActive) return -1; if (!leftActive && rightActive) return 1; - return (left.queuePosition ?? 0) - (right.queuePosition ?? 0); + const positionComparison = (left.queuePosition ?? Number.MAX_SAFE_INTEGER) - + (right.queuePosition ?? Number.MAX_SAFE_INTEGER); + return positionComparison || left.id.localeCompare(right.id); }) - : sortConfig - ? [...filteredDownloads].sort((a, b) => { - const aUnfinished = a.status !== 'completed'; - const bUnfinished = b.status !== 'completed'; + : sortDownloads(filteredDownloads, isQueueFilter ? queueSortConfig! : sortConfig), + [filteredDownloads, isQueueFilter, queueSortConfig, sortConfig]); + sortedDownloadsRef.current = sortedDownloads; - if (aUnfinished && !bUnfinished) return -1; - if (!aUnfinished && bUnfinished) return 1; + useEffect(() => { + const visibleIds = new Set(sortedDownloads.map(download => download.id)); + setSelectedIds(current => { + const next = new Set(Array.from(current).filter(id => visibleIds.has(id))); + return next.size === current.size ? current : next; + }); + setLastSelectedId(current => current && visibleIds.has(current) ? current : null); + }, [sortedDownloads]); - let comparison = 0; - switch (sortConfig.column) { - case 'File Name': - comparison = (a.fileName || a.url || '').localeCompare(b.fileName || b.url || ''); - break; - case 'Size': - comparison = parseInt(a.size || '0', 10) - parseInt(b.size || '0', 10); - break; - case 'Status': - comparison = a.status.localeCompare(b.status); - break; - case 'Speed': - comparison = parseSpeed(a.speed) - parseSpeed(b.speed); - break; - case 'ETA': - comparison = parseEta(a.eta) - parseEta(b.eta); - break; - case 'Date Added': - comparison = new Date(a.dateAdded || 0).getTime() - new Date(b.dateAdded || 0).getTime(); - break; - } - return sortConfig.direction === 'asc' ? comparison : -comparison; - }) - : filteredDownloads; + useEffect(() => { + setQueueSortConfig(null); + }, [filter, isQueueFilter]); const handleItemClick = (e: React.MouseEvent, item: DownloadItem) => { if (e.detail === 2) { handleDownloadDoubleClick(item); @@ -312,15 +288,17 @@ export const DownloadTable: React.FC = ({ filter }) => { setContextMenu(menu); }; - const handleSort = (column: string) => { - if (filter.startsWith('queue:')) return; // Disable custom sorting in queues - setSortConfig(current => { - if (current?.column === column) { - if (current.direction === 'desc') return null; // Reset sort - return { column, direction: 'desc' }; - } - return { column, direction: 'asc' }; - }); + const handleSort = (column: DownloadSortColumn) => { + const update = (current: DownloadSortConfig | null): DownloadSortConfig => + current?.column === column + ? { column, direction: current.direction === 'asc' ? 'desc' : 'asc' } + : { column, direction: 'asc' }; + + if (isQueueFilter) { + setQueueSortConfig(update); + } else { + setSortConfig(current => update(current)); + } }; @@ -356,8 +334,25 @@ export const DownloadTable: React.FC = ({ filter }) => { } }; - const handleResume = (item: DownloadItem) => { - useDownloadStore.getState().resumeDownload(item.id); + const handleResume = async (item: DownloadItem) => { + try { + const resumed = await useDownloadStore.getState().resumeDownload(item.id); + if (!resumed) { + throw new Error('The backend rejected the start/resume request.'); + } + } catch (error) { + console.error("Failed to resume:", error); + showInteractionError(`Could not resume ${item.fileName}`, error); + } + }; + + const resumeItemsSequentially = async (items: DownloadItem[]) => { + for (const item of items) { + const current = useDownloadStore.getState().downloads.find(download => download.id === item.id); + if (current && canStartDownload(current.status)) { + await handleResume(current); + } + } }; const handleDelete = (ids: string | string[]) => { @@ -449,9 +444,7 @@ export const DownloadTable: React.FC = ({ filter }) => { className="main-control-button" disabled={sortedDownloads.length === 0} onClick={() => { - sortedDownloads - .filter(d => canStartDownload(d.status)) - .forEach(d => handleResume(d)); + void resumeItemsSequentially(sortedDownloads.filter(d => canStartDownload(d.status))); }} title="Resume All" > @@ -497,13 +490,15 @@ export const DownloadTable: React.FC = ({ filter }) => { {['File Name', 'Size', 'Status', 'Speed', 'ETA', 'Date Added'].map((label, index) => (
handleSort(label)} + className={`${index === 5 ? 'download-cell-right' : ''} cursor-pointer hover:text-text-primary transition-colors flex items-center justify-between`} + onClick={() => handleSort(label as DownloadSortColumn)} >
{label} - {sortConfig?.column === label && ( - sortConfig.direction === 'asc' ? : + {(isQueueFilter ? queueSortConfig : sortConfig)?.column === label && ( + (isQueueFilter ? queueSortConfig : sortConfig)?.direction === 'asc' + ? + : )}
= ({ filter }) => { {sortedDownloads.length === 0 ? (
) : ( @@ -576,27 +581,17 @@ export const DownloadTable: React.FC = ({ filter }) => { const itemsToResume = selectedDownloads.filter(d => canStartDownload(d.status)); const itemsToPause = selectedDownloads.filter(d => canPauseDownload(d.status)); + const itemsToQueue = selectedDownloads.filter(d => d.status !== 'completed'); return ( <> {/* Multi-Select Context Menu */} {itemsToResume.length > 0 && ( -
- {queues.map(q => ( - - ))} + {itemsToQueue.length > 0 && ( +
+ +
+ {queues.map(q => ( + + ))} +
-
+ )}
diff --git a/src/components/Sidebar.tsx b/src/components/Sidebar.tsx index 6047159..9b21c36 100644 --- a/src/components/Sidebar.tsx +++ b/src/components/Sidebar.tsx @@ -10,6 +10,7 @@ import { useDownloadStore, DownloadCategory, Queue } from '../store/useDownloadS import { ActiveView, useSettingsStore } from '../store/useSettingsStore'; import { WindowDragRegion } from './WindowDragRegion'; import { useToast } from '../contexts/ToastContext'; +import { isTransferActiveStatus } from '../utils/downloads'; export type SidebarFilter = 'all' | 'active' | 'completed' | 'unfinished' | DownloadCategory | 'settings' | string; @@ -100,7 +101,7 @@ export const Sidebar: React.FC = (props) => { } switch (filter) { case 'all': return downloads.length; - case 'active': return downloads.filter(d => d.status === 'downloading').length; + case 'active': return downloads.filter(d => isTransferActiveStatus(d.status)).length; case 'completed': return downloads.filter(d => d.status === 'completed').length; case 'unfinished': return downloads.filter(d => d.status !== 'completed').length; default: return downloads.filter(d => d.category === filter as DownloadCategory).length; @@ -330,14 +331,34 @@ export const Sidebar: React.FC = (props) => { >