fix(downloads): harden queue admission and live speed controls

This commit is contained in:
NimBold
2026-07-23 04:13:07 +03:30
parent b60818d3af
commit 3587fb0c0d
17 changed files with 1298 additions and 93 deletions
+12
View File
@@ -70,6 +70,18 @@ pub struct Queue {
pub id: String,
pub name: String,
pub is_main: bool,
#[serde(default)]
#[ts(optional)]
pub max_concurrent: Option<usize>,
}
#[derive(Clone, Debug, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export, export_to = "../../src/bindings/")]
pub struct QueueConcurrencyConfig {
pub id: String,
#[serde(default)]
pub max_concurrent: Option<usize>,
}
#[derive(Clone, Debug, Serialize, Deserialize, TS)]
+112 -30
View File
@@ -4498,7 +4498,12 @@ async fn resume_download(
app_handle: tauri::AppHandle,
state: tauri::State<'_, AppState>,
id: String,
queue_id: String,
) -> Result<bool, String> {
let queue_id = queue_id.trim().to_string();
if queue_id.is_empty() {
return Err("Queue id cannot be empty".to_string());
}
let control_guard = state.queue_manager.acquire_aria2_control(&id).await;
let Some(gid) = state.queue_manager.aria2_gid_for_download(&id) else {
log::info!(
@@ -4521,6 +4526,11 @@ async fn resume_download(
match status.as_str() {
"paused" => {
let control_epoch = state.queue_manager.next_aria2_control_epoch(&id).await;
let lifecycle_generation = state
.queue_manager
.registered_lifecycle_generation(&id)
.await
.unwrap_or_default();
state.queue_manager.allow_aria2_retries(&id).await;
state.queue_manager.reset_aria2_retry_strikes(&id).await;
use tauri::Emitter;
@@ -4542,7 +4552,13 @@ async fn resume_download(
let permit_candidate = if had_permit {
None
} else {
queue_manager.acquire_aria2_permit_candidate().await
queue_manager
.acquire_aria2_permit_candidate_for_queue(
&id_clone,
&queue_id,
lifecycle_generation,
)
.await
};
if permit_candidate.is_none() && !had_permit {
return;
@@ -4557,12 +4573,22 @@ async fn resume_download(
!= Some(gid_clone.as_str())
|| !queue_manager.is_registered(&id_clone).await
{
if permit_candidate.is_some() {
queue_manager
.release_aria2_permit_candidate(&id_clone, lifecycle_generation)
.await;
}
return;
}
if !queue_manager
.rebind_aria2_gid_epoch(&id_clone, &gid_clone, control_epoch)
.await
{
if permit_candidate.is_some() {
queue_manager
.release_aria2_permit_candidate(&id_clone, lifecycle_generation)
.await;
}
log::warn!(
"aria2 resume [{}]: gid {} disappeared before lifecycle rebind",
id_clone,
@@ -4572,7 +4598,12 @@ async fn resume_download(
}
if let Some(permit) = permit_candidate {
let _ = queue_manager
.park_aria2_permit_if_missing(&id_clone, permit)
.park_aria2_permit_if_missing_for_queue(
&id_clone,
&queue_id,
lifecycle_generation,
permit,
)
.await;
}
let _ = app_handle_clone.emit(
@@ -4702,47 +4733,82 @@ async fn resume_download(
}
"active" | "waiting" => {
let resume_epoch = state.queue_manager.current_aria2_control_epoch(&id).await;
let lifecycle_generation = state
.queue_manager
.registered_lifecycle_generation(&id)
.await
.unwrap_or_default();
drop(control_guard);
state.queue_manager.allow_aria2_retries(&id).await;
let had_permit = state.queue_manager.has_active_permit(&id).await;
let permit_candidate = if had_permit {
None
} else {
state.queue_manager.acquire_aria2_permit_candidate().await
};
if permit_candidate.is_none() && !had_permit {
return Ok(true);
}
let _control_guard = state.queue_manager.acquire_aria2_control(&id).await;
let still_current = state.queue_manager.is_registered(&id).await
&& !state.queue_manager.is_aria2_retry_cancelled(&id).await
&& state
.queue_manager
.is_aria2_control_epoch_current(&id, resume_epoch)
.await
&& state.queue_manager.aria2_gid_for_download(&id).as_deref() == Some(gid.as_str());
if still_current {
let queue_manager = state.queue_manager.clone();
let id_clone = id.clone();
let gid_clone = gid.clone();
let queue_id_clone = queue_id.clone();
let app_handle_clone = app_handle.clone();
let status_for_log = status.to_string();
tauri::async_runtime::spawn(async move {
let had_permit = queue_manager.has_active_permit(&id_clone).await;
let permit_candidate = if had_permit {
None
} else {
queue_manager
.acquire_aria2_permit_candidate_for_queue(
&id_clone,
&queue_id_clone,
lifecycle_generation,
)
.await
};
let had_permit_candidate = permit_candidate.is_some();
if permit_candidate.is_none() && !had_permit {
return;
}
let _control_guard = queue_manager.acquire_aria2_control(&id_clone).await;
let still_current = queue_manager.is_registered(&id_clone).await
&& !queue_manager.is_aria2_retry_cancelled(&id_clone).await
&& queue_manager
.is_aria2_control_epoch_current(&id_clone, resume_epoch)
.await
&& queue_manager.aria2_gid_for_download(&id_clone).as_deref()
== Some(gid_clone.as_str());
if !still_current {
if had_permit_candidate {
queue_manager
.release_aria2_permit_candidate(&id_clone, lifecycle_generation)
.await;
}
return;
}
if let Some(permit) = permit_candidate {
let _ = state
.queue_manager
.park_aria2_permit_if_missing(&id, permit)
let parked = queue_manager
.park_aria2_permit_if_missing_for_queue(
&id_clone,
&queue_id_clone,
lifecycle_generation,
permit,
)
.await;
if !parked && !queue_manager.has_active_permit(&id_clone).await {
return;
}
}
log::info!(
"aria2 resume [{}]: gid {} already {}; no duplicate job created",
id,
gid,
status
id_clone,
gid_clone,
status_for_log
);
use tauri::Emitter;
let _ = app_handle.emit(
let _ = app_handle_clone.emit(
"download-state",
crate::ipc::DownloadStateEvent::new(
id,
&id_clone,
crate::ipc::DownloadStatus::Downloading,
),
);
}
});
Ok(true)
}
"complete" | "error" | "removed" => {
@@ -5687,6 +5753,22 @@ async fn set_concurrent_limit(
Ok(())
}
#[tauri::command]
async fn set_queue_concurrency_limits(
state: tauri::State<'_, AppState>,
limits: Vec<crate::ipc::QueueConcurrencyConfig>,
) -> Result<(), String> {
state
.queue_manager
.replace_queue_limits(
limits
.into_iter()
.map(|limit| (limit.id, limit.max_concurrent))
.collect(),
)
.await
}
#[tauri::command]
async fn set_download_speed_limit(
state: tauri::State<'_, AppState>,
@@ -9327,7 +9409,7 @@ pub fn run() {
authorize_keychain_access,
acknowledge_pairing_token_change,
check_file_exists, toggle_tray_icon, set_extension_pairing_token,
get_extension_server_port, set_extension_frontend_ready, ack_extension_download, set_concurrent_limit, set_download_speed_limit, set_global_speed_limit, remove_download,
get_extension_server_port, set_extension_frontend_ready, ack_extension_download, set_concurrent_limit, set_queue_concurrency_limits, set_download_speed_limit, set_global_speed_limit, remove_download,
detach_download_for_reconfigure,
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,
+410 -53
View File
@@ -14,6 +14,7 @@ use ts_rs::TS;
/// Default capacity when no setting is read yet.
pub const DEFAULT_MAX_CONCURRENT: usize = 3;
pub const MAX_QUEUE_CONCURRENT: usize = 12;
pub const MEDIA_RUN_CANCELLED: &str = "__firelink_media_run_cancelled__";
pub const DOWNLOAD_CONNECTIONS_MIN: i32 = 1;
pub const DOWNLOAD_CONNECTIONS_MAX: i32 = 16;
@@ -122,6 +123,13 @@ pub struct QueuedTask {
pub lifecycle_generation: u64,
}
#[derive(Debug, Clone)]
struct QueuePermitOwnership {
queue_id: String,
lifecycle_generation: u64,
active: bool,
}
/// Args mirroring start_download / start_media_download. Kept untyped-loose
/// (String/Option) to match the existing command signatures exactly.
#[derive(Debug, Clone, Default)]
@@ -197,6 +205,19 @@ pub struct QueueManager<R: tauri::Runtime = tauri::Wry> {
active_permits: Mutex<HashMap<String, OwnedSemaphorePermit>>,
active_permit_generations: Mutex<HashMap<String, u64>>,
active_kinds: Mutex<HashMap<String, TaskKind>>,
/// Queue overrides are stored only for queues with an explicit limit.
/// Missing entries inherit the global target capacity.
queue_limits: Mutex<HashMap<String, usize>>,
/// One entry represents either a queued dispatch reservation or an active
/// transfer. Keeping both phases in one map makes queue-slot ownership
/// exactly-once across the async addUri handoff.
queue_permit_ownership: Mutex<HashMap<String, QueuePermitOwnership>>,
/// Serializes queue-slot selection with global permit acquisition and
/// ownership transitions.
admission_gate: Mutex<()>,
/// Last queue selected by the dispatcher. Selection starts after this
/// queue when multiple queues have eligible work.
dispatch_cursor: Mutex<Option<String>>,
target_capacity: AtomicUsize,
slots_to_retire: AtomicUsize,
notify: Notify,
@@ -275,6 +296,10 @@ impl<R: tauri::Runtime> QueueManager<R> {
active_permits: Mutex::new(HashMap::new()),
active_permit_generations: Mutex::new(HashMap::new()),
active_kinds: Mutex::new(HashMap::new()),
queue_limits: Mutex::new(HashMap::new()),
queue_permit_ownership: Mutex::new(HashMap::new()),
admission_gate: Mutex::new(()),
dispatch_cursor: Mutex::new(None),
target_capacity: AtomicUsize::new(capacity),
slots_to_retire: AtomicUsize::new(0),
notify: Notify::new(),
@@ -315,6 +340,7 @@ impl<R: tauri::Runtime> QueueManager<R> {
// Epoch checks remain the authoritative guard; removing this marker
// prevents terminal downloads from accumulating cancellation entries.
self.aria2_retry_cancelled.lock().await.remove(id);
self.notify.notify_waiters();
}
pub async fn is_registered(&self, id: &str) -> bool {
@@ -463,6 +489,38 @@ impl<R: tauri::Runtime> QueueManager<R> {
self.push_with_generation(task, 0).await
}
/// Replace the explicit per-queue concurrency overrides. A missing queue
/// entry inherits the global limit; every effective queue limit is capped
/// by the global target at admission time.
pub async fn replace_queue_limits(
&self,
limits: Vec<(String, Option<usize>)>,
) -> Result<(), String> {
let mut next = HashMap::with_capacity(limits.len());
for (queue_id, limit) in limits {
let queue_id = queue_id.trim();
if queue_id.is_empty() {
return Err("Queue id cannot be empty".to_string());
}
if next.contains_key(queue_id) {
return Err(format!("Duplicate queue id '{queue_id}'"));
}
if let Some(limit) = limit {
if !(1..=MAX_QUEUE_CONCURRENT).contains(&limit) {
return Err(format!(
"Queue concurrency must be between 1 and {MAX_QUEUE_CONCURRENT}"
));
}
next.insert(queue_id.to_string(), limit);
}
}
let _admission_gate = self.admission_gate.lock().await;
*self.queue_limits.lock().await = next;
self.notify.notify_waiters();
Ok(())
}
pub async fn next_aria2_control_epoch(&self, id: &str) -> u64 {
let mut epochs = self.aria2_control_epochs.lock().await;
let epoch = epochs.get(id).copied().unwrap_or_default().wrapping_add(1);
@@ -627,6 +685,17 @@ impl<R: tauri::Runtime> QueueManager<R> {
self.semaphore.clone().acquire_owned().await.ok()
}
fn try_acquire_permit_after_retirement(&self) -> Option<OwnedSemaphorePermit> {
loop {
let permit = self.semaphore.clone().try_acquire_owned().ok()?;
if self.retire_slot_if_needed() {
permit.forget();
continue;
}
return Some(permit);
}
}
async fn acquire_permit_after_retirement(&self) -> Option<OwnedSemaphorePermit> {
loop {
let permit = self.acquire_permit().await?;
@@ -645,6 +714,37 @@ impl<R: tauri::Runtime> QueueManager<R> {
self.acquire_permit_after_retirement().await
}
/// Acquire and reserve both the global slot and a queue slot for a resume
/// that may still be outside the per-download control lock. The reservation
/// is generation-stamped so a concurrent pause cannot be overwritten by a
/// late resume worker.
pub async fn acquire_aria2_permit_candidate_for_queue(
&self,
id: &str,
queue_id: &str,
lifecycle_generation: u64,
) -> Option<OwnedSemaphorePermit> {
loop {
if !self.is_registered_generation(id, lifecycle_generation).await {
return None;
}
let notified = self.notify.notified();
if let Some(permit) = self
.try_reserve_queue_slot(id, queue_id, lifecycle_generation)
.await
{
if self.is_registered_generation(id, lifecycle_generation).await {
return Some(permit);
}
self.release_queue_reservation_for_generation(id, lifecycle_generation)
.await;
drop(permit);
return None;
}
notified.await;
}
}
fn retire_slot_if_needed(&self) -> bool {
let mut debt = self.slots_to_retire.load(Ordering::Relaxed);
while debt > 0 {
@@ -661,12 +761,180 @@ impl<R: tauri::Runtime> QueueManager<R> {
false
}
async fn try_reserve_queue_slot(
&self,
id: &str,
queue_id: &str,
lifecycle_generation: u64,
) -> Option<OwnedSemaphorePermit> {
let _admission_gate = self.admission_gate.lock().await;
let mut ownership = self.queue_permit_ownership.lock().await;
if ownership.contains_key(id)
|| self.active_permits.lock().await.contains_key(id)
{
return None;
}
let global_target = self.target_capacity.load(Ordering::Relaxed);
let queue_limit = self
.queue_limits
.lock()
.await
.get(queue_id)
.copied()
.unwrap_or(global_target)
.min(global_target);
let queue_active = ownership
.values()
.filter(|active| active.queue_id == queue_id)
.count();
if queue_active >= queue_limit {
return None;
}
let permit = self.try_acquire_permit_after_retirement()?;
ownership.insert(
id.to_string(),
QueuePermitOwnership {
queue_id: queue_id.to_string(),
lifecycle_generation,
active: false,
},
);
Some(permit)
}
async fn release_queue_reservation_for_generation(&self, id: &str, generation: u64) {
let _admission_gate = self.admission_gate.lock().await;
let removed = self
.queue_permit_ownership
.lock()
.await
.get(id)
.is_some_and(|ownership| {
ownership.lifecycle_generation == generation && !ownership.active
});
if removed {
self.queue_permit_ownership.lock().await.remove(id);
self.notify.notify_waiters();
}
}
async fn activate_admitted_permit(
&self,
id: &str,
lifecycle_generation: u64,
permit: OwnedSemaphorePermit,
) -> bool {
let _admission_gate = self.admission_gate.lock().await;
let mut ownership = self.queue_permit_ownership.lock().await;
let owned = ownership
.get(id)
.is_some_and(|entry| entry.lifecycle_generation == lifecycle_generation && !entry.active);
let active_already_owned = self.active_permits.lock().await.contains_key(id);
if !owned || active_already_owned {
if owned {
ownership.remove(id);
self.notify.notify_waiters();
}
return false;
}
if let Some(entry) = ownership.get_mut(id) {
entry.active = true;
}
drop(ownership);
self.active_permits.lock().await.insert(id.to_string(), permit);
self.active_permit_generations
.lock()
.await
.insert(id.to_string(), lifecycle_generation);
true
}
async fn try_admit_next_task(&self) -> Option<(OwnedSemaphorePermit, QueuedTask)> {
let _admission_gate = self.admission_gate.lock().await;
let mut pending = self.pending.lock().await;
if pending.is_empty() {
return None;
}
let ownership = self.queue_permit_ownership.lock().await;
let mut queue_ids = Vec::new();
let mut seen = HashSet::new();
for task in pending.iter() {
if !ownership.contains_key(&task.id) && seen.insert(task.queue_id.clone()) {
queue_ids.push(task.queue_id.clone());
}
}
if queue_ids.is_empty() {
return None;
}
let cursor = self.dispatch_cursor.lock().await.clone();
let start = cursor
.as_ref()
.and_then(|queue_id| queue_ids.iter().position(|candidate| candidate == queue_id))
.map_or(0, |position| (position + 1) % queue_ids.len());
let global_target = self.target_capacity.load(Ordering::Relaxed);
let queue_limits = self.queue_limits.lock().await;
let selected_queue = (0..queue_ids.len()).find_map(|offset| {
let queue_id = &queue_ids[(start + offset) % queue_ids.len()];
let queue_limit = queue_limits
.get(queue_id)
.copied()
.unwrap_or(global_target)
.min(global_target);
let active = ownership
.values()
.filter(|active| active.queue_id == *queue_id)
.count();
(active < queue_limit).then_some(queue_id.clone())
})?;
let task_index = pending
.iter()
.position(|task| task.queue_id == selected_queue && !ownership.contains_key(&task.id))?;
drop(queue_limits);
drop(ownership);
let permit = self.try_acquire_permit_after_retirement()?;
let task = pending.remove(task_index)?;
self.queue_permit_ownership.lock().await.insert(
task.id.clone(),
QueuePermitOwnership {
queue_id: selected_queue.clone(),
lifecycle_generation: task.lifecycle_generation,
active: false,
},
);
*self.dispatch_cursor.lock().await = Some(selected_queue);
Some((permit, task))
}
/// Park an already-acquired permit under `id`.
pub async fn park_permit(&self, id: &str, permit: OwnedSemaphorePermit) {
self.park_permit_for_queue(id, "main", 0, permit).await;
}
async fn park_permit_for_queue(
&self,
id: &str,
queue_id: &str,
lifecycle_generation: u64,
permit: OwnedSemaphorePermit,
) {
let _admission_gate = self.admission_gate.lock().await;
self.active_permits
.lock()
.await
.insert(id.to_string(), permit);
self.queue_permit_ownership.lock().await.insert(
id.to_string(),
QueuePermitOwnership {
queue_id: queue_id.to_string(),
lifecycle_generation,
active: true,
},
);
}
/// Park a candidate only when no newer lifecycle has already claimed the
@@ -676,12 +944,72 @@ impl<R: tauri::Runtime> QueueManager<R> {
id: &str,
permit: OwnedSemaphorePermit,
) -> bool {
self.park_aria2_permit_if_missing_for_queue(id, "main", 0, permit)
.await
}
/// Activate a queue-aware resume reservation. The caller owns the permit
/// while it revalidates the lifecycle under the per-download control lock.
pub async fn park_aria2_permit_if_missing_for_queue(
&self,
id: &str,
queue_id: &str,
lifecycle_generation: u64,
permit: OwnedSemaphorePermit,
) -> bool {
let _admission_gate = self.admission_gate.lock().await;
let registered_generation = self
.registered_lifecycle_generations
.lock()
.await
.get(id)
.copied();
let mut ownership = self.queue_permit_ownership.lock().await;
let has_matching_reservation = ownership.get(id).is_some_and(|entry| {
entry.queue_id == queue_id
&& entry.lifecycle_generation == lifecycle_generation
&& !entry.active
&& (registered_generation == Some(lifecycle_generation)
|| (registered_generation.is_none() && lifecycle_generation == 0))
});
if ownership.contains_key(id) && !has_matching_reservation {
let remove_reservation = ownership
.get(id)
.is_some_and(|entry| entry.lifecycle_generation == lifecycle_generation && !entry.active);
if remove_reservation {
ownership.remove(id);
self.notify.notify_waiters();
}
return false;
}
if ownership.get(id).is_none()
&& registered_generation.is_some_and(|current| current != lifecycle_generation)
{
return false;
}
let mut permits = self.active_permits.lock().await;
if permits.contains_key(id) {
return false;
}
permits.insert(id.to_string(), permit);
drop(permits);
if let Some(entry) = ownership.get_mut(id) {
entry.active = true;
} else {
ownership.insert(
id.to_string(),
QueuePermitOwnership {
queue_id: queue_id.to_string(),
lifecycle_generation,
active: true,
},
);
}
drop(ownership);
self.active_permit_generations
.lock()
.await
.insert(id.to_string(), lifecycle_generation);
self.active_kinds
.lock()
.await
@@ -689,11 +1017,9 @@ impl<R: tauri::Runtime> QueueManager<R> {
true
}
async fn tag_permit_generation(&self, id: &str, generation: u64) {
self.active_permit_generations
.lock()
.await
.insert(id.to_string(), generation);
pub async fn release_aria2_permit_candidate(&self, id: &str, lifecycle_generation: u64) {
self.release_queue_reservation_for_generation(id, lifecycle_generation)
.await;
}
pub async fn active_kind(&self, id: &str) -> Option<TaskKind> {
@@ -704,52 +1030,84 @@ impl<R: tauri::Runtime> QueueManager<R> {
/// when this call acquired and parked the permit, false when one was
/// already parked.
pub async fn ensure_aria2_permit(&self, id: &str) -> bool {
if self.active_permits.lock().await.contains_key(id) {
return false;
}
self.ensure_aria2_permit_for_queue(id, "main").await
}
let permit = match self.acquire_permit_after_retirement().await {
Some(p) => p,
None => return false,
};
let mut permits = self.active_permits.lock().await;
if permits.contains_key(id) {
drop(permits);
drop(permit);
return false;
/// Ensure a paused or externally recovered aria2 transfer owns one global
/// permit and one slot in its queue. This waits without holding a
/// per-download control lock and wakes on either kind of capacity change.
pub async fn ensure_aria2_permit_for_queue(&self, id: &str, queue_id: &str) -> bool {
loop {
if self.has_active_permit(id).await
|| self.queue_permit_ownership.lock().await.contains_key(id)
{
return false;
}
let notified = self.notify.notified();
let generation = self
.registered_lifecycle_generation(id)
.await
.unwrap_or_default();
if let Some(permit) = self
.try_reserve_queue_slot(id, queue_id, generation)
.await
{
if self.park_aria2_permit_if_missing_for_queue(
id,
queue_id,
generation,
permit,
)
.await
{
self.active_kinds
.lock()
.await
.insert(id.to_string(), TaskKind::Aria2);
return true;
}
return false;
}
notified.await;
}
permits.insert(id.to_string(), permit);
drop(permits);
self.active_kinds
.lock()
.await
.insert(id.to_string(), TaskKind::Aria2);
true
}
pub async fn release_permit(&self, id: &str) {
let _admission_gate = self.admission_gate.lock().await;
let queue_removed = self.queue_permit_ownership.lock().await.remove(id).is_some();
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();
if removed || queue_removed {
self.notify.notify_waiters();
}
}
async fn release_permit_for_generation(&self, id: &str, generation: u64) {
let _admission_gate = self.admission_gate.lock().await;
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) {
let queue_owned = self
.queue_permit_ownership
.lock()
.await
.get(id)
.is_some_and(|ownership| ownership.lifecycle_generation == generation);
if queue_owned || generations.get(id).copied() == Some(generation) {
generations.remove(id);
permits.remove(id).is_some()
let active_removed = permits.remove(id).is_some();
if queue_owned {
self.queue_permit_ownership.lock().await.remove(id);
}
active_removed || queue_owned
} else {
false
}
};
if removed {
self.active_kinds.lock().await.remove(id);
self.notify.notify_one();
self.notify.notify_waiters();
}
}
@@ -821,10 +1179,11 @@ impl<R: tauri::Runtime> QueueManager<R> {
if delta > 0 {
self.semaphore.add_permits(delta);
}
self.notify.notify_one();
self.notify.notify_waiters();
} else {
let delta = prev_target - new_target;
self.slots_to_retire.fetch_add(delta, Ordering::Relaxed);
self.notify.notify_waiters();
}
}
@@ -834,29 +1193,20 @@ impl<R: tauri::Runtime> QueueManager<R> {
}
/// The long-running dispatcher. One instance is spawned in setup().
/// Idle-parks on Notify; CAS-honors retirement debt; re-pops under lock.
/// It scans for a queue with capacity before reserving the global slot, so
/// a saturated front queue cannot block later eligible queues.
pub async fn run_dispatcher(self: Arc<Self>) {
loop {
// (1) Idle-park: avoid busy-spin when pending is empty.
if self.pending.lock().await.is_empty() {
self.notify.notified().await;
continue;
let notified = self.notify.notified();
if let Some((permit, task)) = self.try_admit_next_task().await {
Arc::clone(&self).dispatch_one(permit, task).await;
} else {
// This covers both an empty pending list and the case where
// all queue/global capacity is occupied. The notification
// future is created before inspection to close the lost-wake
// window without polling or sleeping.
notified.await;
}
// (2) Acquire a slot.
let permit = match self.acquire_permit_after_retirement().await {
Some(p) => p,
None => break, // Semaphore closed, exit dispatcher
};
// (4) Re-pop under lock — guards against racing removals between
// waking from Notify and acquiring the permit.
let task = match self.pending.lock().await.pop_front() {
Some(t) => t,
None => {
drop(permit);
continue;
}
};
Arc::clone(&self).dispatch_one(permit, task).await;
}
}
@@ -872,16 +1222,23 @@ impl<R: tauri::Runtime> QueueManager<R> {
.is_registered_generation(&id, lifecycle_generation)
.await
{
self.release_queue_reservation_for_generation(&id, lifecycle_generation)
.await;
drop(control_guard);
drop(permit);
return;
}
// Park the permit BEFORE spawning. Uniform parking:
// Activate the reservation 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;
if !self
.activate_admitted_permit(&id, lifecycle_generation, permit)
.await
{
drop(control_guard);
return;
}
self.active_kinds
.lock()
.await
+245
View File
@@ -783,6 +783,245 @@ async fn aria2_resume_waits_for_shrunk_capacity() {
.expect("resume task should not panic"));
}
#[tokio::test]
async fn dispatcher_skips_a_full_front_queue_for_later_eligible_work() {
let (manager, spawner) = make_manager(2);
let manager = Arc::new(manager);
manager
.replace_queue_limits(vec![
("queue-a".to_string(), Some(1)),
("queue-b".to_string(), Some(1)),
])
.await
.unwrap();
manager.push(aria2_task_in_queue("a1", "queue-a")).await.unwrap();
manager.push(aria2_task_in_queue("a2", "queue-a")).await.unwrap();
manager.push(aria2_task_in_queue("b1", "queue-b")).await.unwrap();
let dispatcher = {
let manager = Arc::clone(&manager);
tokio::spawn(async move { manager.run_dispatcher().await })
};
timeout(Duration::from_secs(1), async {
loop {
if spawner.add_uri_calls.load(Ordering::SeqCst) == 2 {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("later eligible queue should dispatch");
assert!(manager.aria2_gid_for_download("a1").is_some());
assert!(manager.aria2_gid_for_download("b1").is_some());
assert!(manager.aria2_gid_for_download("a2").is_none());
dispatcher.abort();
}
#[tokio::test]
async fn dispatcher_rotates_eligible_queues_without_starving_later_work() {
let (manager, spawner) = make_manager(1);
let manager = Arc::new(manager);
manager
.replace_queue_limits(vec![
("queue-a".to_string(), Some(1)),
("queue-b".to_string(), Some(1)),
])
.await
.unwrap();
manager.push(aria2_task_in_queue("a1", "queue-a")).await.unwrap();
manager.push(aria2_task_in_queue("a2", "queue-a")).await.unwrap();
manager.push(aria2_task_in_queue("b1", "queue-b")).await.unwrap();
let dispatcher = {
let manager = Arc::clone(&manager);
tokio::spawn(async move { manager.run_dispatcher().await })
};
timeout(Duration::from_secs(1), async {
while spawner.add_uri_calls.load(Ordering::SeqCst) < 1 {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("first queue task should dispatch");
manager.release_permit("a1").await;
timeout(Duration::from_secs(1), async {
while manager.aria2_gid_for_download("b1").is_none() {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("later queue should be selected before another task from queue-a");
assert!(manager.aria2_gid_for_download("a2").is_none());
dispatcher.abort();
}
#[tokio::test]
async fn queue_limit_increase_wakes_waiting_work_and_decrease_keeps_active_tasks() {
let (manager, spawner) = make_manager(3);
let manager = Arc::new(manager);
manager
.replace_queue_limits(vec![("queue-a".to_string(), Some(2))])
.await
.unwrap();
for id in ["a1", "a2", "a3"] {
manager.push(aria2_task_in_queue(id, "queue-a")).await.unwrap();
}
let dispatcher = {
let manager = Arc::clone(&manager);
tokio::spawn(async move { manager.run_dispatcher().await })
};
timeout(Duration::from_secs(1), async {
while spawner.add_uri_calls.load(Ordering::SeqCst) < 2 {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("queue override should allow two active transfers");
manager
.replace_queue_limits(vec![("queue-a".to_string(), Some(1))])
.await
.unwrap();
assert_eq!(spawner.add_uri_calls.load(Ordering::SeqCst), 2);
manager.release_permit("a1").await;
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
spawner.add_uri_calls.load(Ordering::SeqCst),
2,
"reducing a queue limit must not admit work while one active transfer remains"
);
manager.release_permit("a2").await;
timeout(Duration::from_secs(1), async {
while spawner.add_uri_calls.load(Ordering::SeqCst) < 3 {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("queue work should resume after active count falls below the reduced limit");
dispatcher.abort();
}
#[tokio::test]
async fn queue_overrides_never_exceed_the_global_ceiling() {
let (manager, spawner) = make_manager(2);
let manager = Arc::new(manager);
manager
.replace_queue_limits(vec![
("queue-a".to_string(), Some(12)),
("queue-b".to_string(), Some(12)),
])
.await
.unwrap();
for (id, queue_id) in [
("a1", "queue-a"),
("a2", "queue-a"),
("b1", "queue-b"),
("b2", "queue-b"),
] {
manager.push(aria2_task_in_queue(id, queue_id)).await.unwrap();
}
let dispatcher = {
let manager = Arc::clone(&manager);
tokio::spawn(async move { manager.run_dispatcher().await })
};
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(spawner.add_uri_calls.load(Ordering::SeqCst), 2);
assert_eq!(manager.available_permits(), 0);
dispatcher.abort();
}
#[tokio::test]
async fn queue_limit_configuration_rejects_malformed_input() {
let (manager, _spawner) = make_manager(2);
assert!(manager
.replace_queue_limits(vec![(String::new(), Some(1))])
.await
.is_err());
assert!(manager
.replace_queue_limits(vec![("queue-a".to_string(), Some(0))])
.await
.is_err());
assert!(manager
.replace_queue_limits(vec![
("queue-a".to_string(), Some(1)),
("queue-a".to_string(), None),
])
.await
.is_err());
manager
.replace_queue_limits(vec![("queue-a".to_string(), None)])
.await
.unwrap();
}
#[tokio::test]
async fn stale_queue_resume_reservation_is_removed_when_activation_fails() {
let (manager, _spawner) = make_manager(1);
let previous = manager
.reserve_enqueue_generation("candidate", 1)
.await
.unwrap();
assert!(previous.is_none());
let candidate = manager
.acquire_aria2_permit_candidate_for_queue("candidate", "queue-a", 1)
.await
.expect("candidate should reserve the queue slot");
manager.release_registered_id("candidate").await;
assert!(!manager
.park_aria2_permit_if_missing_for_queue("candidate", "queue-a", 1, candidate)
.await);
assert_eq!(manager.available_permits(), 1);
assert!(manager.ensure_aria2_permit_for_queue("replacement", "queue-a").await);
manager.release_permit("replacement").await;
}
#[tokio::test]
async fn duplicate_pending_id_cannot_replace_an_existing_queue_ownership() {
let (manager, spawner) = make_manager(2);
let manager = Arc::new(manager);
manager
.replace_queue_limits(vec![
("queue-a".to_string(), Some(1)),
("queue-b".to_string(), Some(1)),
])
.await
.unwrap();
assert!(manager
.ensure_aria2_permit_for_queue("duplicate", "queue-b")
.await);
manager
.reserve_enqueue_generation("duplicate", 1)
.await
.unwrap();
manager
.commit_reserved_enqueue(aria2_task_in_queue("duplicate", "queue-a"), 1)
.await
.unwrap();
let dispatcher = {
let manager = Arc::clone(&manager);
tokio::spawn(async move { manager.run_dispatcher().await })
};
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(spawner.add_uri_calls.load(Ordering::SeqCst), 0);
assert_eq!(manager.available_permits(), 1);
manager.release_permit("duplicate").await;
dispatcher.abort();
}
fn aria2_task(id: &str) -> QueuedTask {
QueuedTask {
id: id.to_string(),
@@ -793,6 +1032,12 @@ fn aria2_task(id: &str) -> QueuedTask {
}
}
fn aria2_task_in_queue(id: &str, queue_id: &str) -> QueuedTask {
let mut task = aria2_task(id);
task.queue_id = queue_id.to_string();
task
}
fn media_task(id: &str) -> QueuedTask {
QueuedTask {
id: id.to_string(),