From 2263e92ef462ff9d5083bb5501c2a8b87b1387ef Mon Sep 17 00:00:00 2001 From: Wez Furlong Date: Wed, 28 Aug 2024 08:30:50 -0700 Subject: [PATCH] queue: fix non-singleton maintainer wakeups I recently broke when adding the singleton maintainer task. The nature of the issue is that the scheduled queue maintainers are created and won't wake up to do any work for 1 day, or until they get notified by an action like a rebind operation. Why didn't the integration tests catch this obvious regression? Because they didn't assert that we saw a specific number of attempts, only that the mathmatical relationship between an arbitrary number of records was correct. As part of this, I found a similar sort of wakeup issue with the skiplist queue, where adding the initial entry didn't wakeup the queue because of the way that the Ord trait works on Option. This commit addresses that as well. --- crates/integration-tests/src/main.rs | 3 +++ crates/kumod/src/queue.rs | 4 ++-- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/crates/integration-tests/src/main.rs b/crates/integration-tests/src/main.rs index 6ebf7f6c..023768f3 100644 --- a/crates/integration-tests/src/main.rs +++ b/crates/integration-tests/src/main.rs @@ -1134,6 +1134,9 @@ DeliverySummary { }) .collect(); + println!("***** event_times: {event_times:?}"); + assert!(event_times.len() > 1); + let mut last = None; let mut intervals: Vec<_> = event_times .iter() diff --git a/crates/kumod/src/queue.rs b/crates/kumod/src/queue.rs index 41227a77..92fc56d5 100644 --- a/crates/kumod/src/queue.rs +++ b/crates/kumod/src/queue.rs @@ -738,7 +738,7 @@ impl QueueStructure { // we do not want to wake up for every message insertion, // as that would generally be a waste of effort and bog // down the system without gain. - should_notify: now_due < due, + should_notify: if due.is_none() { true } else { now_due < due }, } } Self::SingletonTimerWheel(q) => { @@ -1888,7 +1888,7 @@ impl QueueManager { #[instrument(skip(q))] async fn maintain_named_queue(q: &QueueHandle) -> anyhow::Result<()> { let mut shutdown = ShutdownSubcription::get(); - let mut next_item_due = Instant::now() + ONE_DAY; + let mut next_item_due = Instant::now(); loop { let sleeping = Instant::now();