From 38e65a3791a15997483e558fa9434c58a7bde202 Mon Sep 17 00:00:00 2001 From: YoursFunny Date: Thu, 24 Sep 2026 02:49:10 +0800 Subject: [PATCH] fix(queue): close the lost-wakeup window on shutdown MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit stop() flags the shutdown and fires notify_waiters, but a worker parked between its loop-top stop check and its notified() registration — i.e. inside earliest_run_after's DB await — was not registered when the notification fired, so it slept until the next enqueue that never comes; stop() then blocked until main's 30s shutdown timeout force-killed the drain. The sweep had the same window before recover_expired's await and would sit out a full 30s tick. Both loops now enable() the waiter first and re-check the stop flag: either the stop already happened (recheck returns) or the waiter is registered (notify_waiters reaches it) — no gap. The window itself is a scheduling race with no test seam, so this is pinned by reasoning rather than a regression test; the existing stop tests cover the ordinary path. --- crates/xmedia-bot/src/queue.rs | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/crates/xmedia-bot/src/queue.rs b/crates/xmedia-bot/src/queue.rs index 3837f33..5ff58b0 100644 --- a/crates/xmedia-bot/src/queue.rs +++ b/crates/xmedia-bot/src/queue.rs @@ -179,6 +179,15 @@ impl PersistentTaskQueue { loop { let notified = sweep_notify.notified(); tokio::pin!(notified); + // Register before re-checking stop (same rule as the worker + // loop): a stop() landing between the previous iteration and + // here would wake nobody, and the sweep would sit out a full + // 30s tick — long enough for the shutdown timeout to treat the + // drain as stuck. + notified.as_mut().enable(); + if sweep_stop.load(Ordering::Relaxed) { + break; + } tokio::select! { _ = &mut notified => {} _ = interval.tick() => {} @@ -379,6 +388,18 @@ impl QueueWorker { let wait_until = self.earliest_run_after().await; let notified = self.notify.notified(); tokio::pin!(notified); + // Register before re-checking stop: a stop() between the + // loop-top check and this point fires notify_waiters with + // no waiter registered, and this worker would then sleep + // until the next enqueue (the same hazard enqueue's + // notify_one comment names). enable() closes it — either + // the stop already happened and the recheck below returns, + // or the waiter is registered and stop's notify_waiters + // reaches it. + notified.as_mut().enable(); + if self.stop.load(Ordering::Relaxed) { + return; + } match wait_until { Some(until) => { let delay = (until - now_f64()).max(0.0);