diff --git a/desktop/crates/cloudreve-sync/src/tasks/queue.rs b/desktop/crates/cloudreve-sync/src/tasks/queue.rs index 745416a4..7a9a796e 100644 --- a/desktop/crates/cloudreve-sync/src/tasks/queue.rs +++ b/desktop/crates/cloudreve-sync/src/tasks/queue.rs @@ -58,6 +58,18 @@ struct FileChangeSignature { modified: SystemTime, } +/// Decrements `inflight` and wakes `wait_for_idle` when the spawned task +/// exits — including when it is aborted by cancel_by_path or shutdown, +/// where the normal tail-decrement would never run. +struct InflightGuard(Arc); + +impl Drop for InflightGuard { + fn drop(&mut self) { + self.0.inflight.fetch_sub(1, Ordering::SeqCst); + self.0.idle_notify.notify_waiters(); + } +} + impl TaskQueue { pub async fn new( drive_id: impl Into, @@ -328,46 +340,32 @@ impl TaskQueue { } async fn launch_task(self: &Arc, task: QueuedTask) { - let permit = match self.semaphore.clone().acquire_owned().await { - Ok(permit) => permit, - Err(err) => { - error!( - target: "tasks::queue", - drive = %self.drive_id, - error = %err, - "Failed to acquire semaphore permit" - ); - if let Err(update_err) = self.inventory.update_task( - &task.task_id, - TaskUpdate { - status: Some(TaskStatus::Failed), - error: Some(Some("Failed to schedule task".to_string())), - ..Default::default() - }, - ) { - warn!( - target: "tasks::queue", - drive = %self.drive_id, - error = %update_err, - "Failed to persist scheduling failure" - ); - } - return; - } - }; - + // Acquire the semaphore inside the spawned task, not in the dispatch + // loop: blocking the loop here starves every queued task whenever + // `max_concurrent` long-running tasks hold the permits, leaving all + // enqueued work "pending" until a running task happens to finish. self.inflight.fetch_add(1, Ordering::SeqCst); - let queue_for_execute = Arc::clone(self); - let queue_for_notify = Arc::clone(self); + let queue = Arc::clone(self); let task_id = task.task_id.clone(); let handle_task_id = task_id.clone(); + let semaphore = self.semaphore.clone(); let handle = tokio::spawn(async move { - queue_for_execute.execute_task(task).await; + // Guarded so an aborted task still decrements inflight and wakes + // wait_for_idle. + let _inflight_guard = InflightGuard(queue.clone()); + let permit = match semaphore.acquire_owned().await { + Ok(permit) => permit, + Err(_) => { + // Semaphore closed during shutdown; leave the row pending + // so resume_incomplete_tasks picks it up on next launch. + queue.task_handles.remove(&handle_task_id); + return; + } + }; + queue.clone().execute_task(task).await; drop(permit); - queue_for_notify.inflight.fetch_sub(1, Ordering::SeqCst); - queue_for_notify.idle_notify.notify_waiters(); - queue_for_notify.task_handles.remove(&handle_task_id); + queue.task_handles.remove(&handle_task_id); }); self.task_handles.insert(task_id, handle); @@ -383,6 +381,7 @@ impl TaskQueue { task_id = %task.task_id, "Task was cancelled before execution, skipping" ); + self.cleanup_task_entry(&task.task_id).await; return; } Ok(Some(status)) if !status.is_active() => { @@ -393,6 +392,7 @@ impl TaskQueue { status = ?status, "Task is no longer active, skipping" ); + self.cleanup_task_entry(&task.task_id).await; return; } Err(err) => {