use crate::clock;
use crate::control::{NodeControlMsg, WakeupRevision, WakeupSlot};
use crate::entity_context::current_node_telemetry_handle;
use crate::indexed_min_heap::IndexedMinHeap;
use otel_arrow_dfe_telemetry::error::Error as TelemetryError;
use otel_arrow_dfe_telemetry::instrument::{Counter, Gauge};
use otel_arrow_dfe_telemetry::metrics::MetricSet;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use otel_arrow_dfe_telemetry_macros::metric_set;
use std::cmp::Ordering;
use std::collections::{BinaryHeap, VecDeque};
use std::sync::{Arc, Mutex};
use std::time::Instant;
use tokio::sync::Notify;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum WakeupError {
Unsupported,
ShuttingDown,
Capacity,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum WakeupSetOutcome {
Inserted {
revision: WakeupRevision,
},
Replaced {
revision: WakeupRevision,
},
}
impl WakeupSetOutcome {
#[must_use]
pub const fn revision(self) -> WakeupRevision {
match self {
Self::Inserted { revision } | Self::Replaced { revision } => revision,
}
}
}
#[metric_set(name = "node.local_wakeup.scheduler")]
#[derive(Debug, Default, Clone)]
pub(crate) struct NodeLocalWakeupSchedulerMetrics {
#[metric(name = "set.inserted", unit = "{wakeup}")]
set_inserted: Counter<u64>,
#[metric(name = "set.replaced", unit = "{wakeup}")]
set_replaced: Counter<u64>,
#[metric(name = "set.error_capacity", unit = "{attempt}")]
set_error_capacity: Counter<u64>,
#[metric(name = "set.error_shutdown", unit = "{attempt}")]
set_error_shutdown: Counter<u64>,
#[metric(name = "cancel.removed", unit = "{wakeup}")]
cancel_removed: Counter<u64>,
#[metric(name = "cancel.missed", unit = "{attempt}")]
cancel_missed: Counter<u64>,
#[metric(name = "cancel.ignored_shutdown", unit = "{attempt}")]
cancel_ignored_shutdown: Counter<u64>,
#[metric(name = "pop.due", unit = "{wakeup}")]
pop_due: Counter<u64>,
#[metric(name = "shutdown.cleared", unit = "{wakeup}")]
shutdown_cleared: Counter<u64>,
#[metric(name = "live", unit = "{wakeup}")]
live: Gauge<u64>,
#[metric(name = "capacity", unit = "{wakeup}")]
capacity: Gauge<u64>,
}
#[derive(Debug)]
struct ScheduledResume<PData> {
when: Instant,
sequence: u64,
data: Box<PData>,
}
impl<PData> Ord for ScheduledResume<PData> {
fn cmp(&self, other: &Self) -> Ordering {
other
.when
.cmp(&self.when)
.then_with(|| other.sequence.cmp(&self.sequence))
}
}
impl<PData> PartialOrd for ScheduledResume<PData> {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl<PData> PartialEq for ScheduledResume<PData> {
fn eq(&self, other: &Self) -> bool {
self.when == other.when && self.sequence == other.sequence
}
}
impl<PData> Eq for ScheduledResume<PData> {}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct WakeupPriority {
when: Instant,
revision: WakeupRevision,
}
impl Ord for WakeupPriority {
fn cmp(&self, other: &Self) -> Ordering {
self.when
.cmp(&other.when)
.then_with(|| self.revision.cmp(&other.revision))
}
}
impl PartialOrd for WakeupPriority {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
struct NodeLocalScheduler<PData> {
delayed_resume_capacity: usize,
wakeup_capacity: usize,
next_sequence: u64,
next_revision: WakeupRevision,
delayed_resumes: BinaryHeap<ScheduledResume<PData>>,
due_now: VecDeque<NodeControlMsg<PData>>,
heap: IndexedMinHeap<WakeupSlot, WakeupPriority>,
shutting_down: bool,
metrics: Option<MetricSet<NodeLocalWakeupSchedulerMetrics>>,
}
impl<PData> NodeLocalScheduler<PData> {
#[cfg(test)]
fn new(delayed_resume_capacity: usize, wakeup_capacity: usize) -> Self {
Self::new_with_metrics(delayed_resume_capacity, wakeup_capacity, None)
}
fn new_with_metrics(
delayed_resume_capacity: usize,
wakeup_capacity: usize,
metrics: Option<MetricSet<NodeLocalWakeupSchedulerMetrics>>,
) -> Self {
Self {
delayed_resume_capacity,
wakeup_capacity,
next_sequence: 0,
next_revision: 0,
delayed_resumes: BinaryHeap::new(),
due_now: VecDeque::new(),
heap: IndexedMinHeap::new(),
shutting_down: false,
metrics,
}
}
fn next_sequence(&mut self) -> u64 {
let next = self.next_sequence;
self.next_sequence = self.next_sequence.saturating_add(1);
next
}
fn next_revision(&mut self) -> WakeupRevision {
let next = self.next_revision;
self.next_revision = self.next_revision.saturating_add(1);
next
}
fn refresh_gauges(&mut self) {
if let Some(metrics) = &mut self.metrics {
metrics.live.set(self.heap.len() as u64);
metrics.capacity.set(self.wakeup_capacity as u64);
}
}
fn requeue_later(&mut self, when: Instant, data: Box<PData>) -> Result<(), Box<PData>> {
if self.shutting_down || self.delayed_resumes.len() >= self.delayed_resume_capacity {
return Err(data);
}
let sequence = self.next_sequence();
self.delayed_resumes.push(ScheduledResume {
when,
sequence,
data,
});
Ok(())
}
fn set_wakeup(
&mut self,
slot: WakeupSlot,
when: Instant,
) -> Result<WakeupSetOutcome, WakeupError> {
if self.wakeup_capacity == 0 {
return Err(WakeupError::Unsupported);
}
if self.shutting_down {
if let Some(metrics) = &mut self.metrics {
metrics.set_error_shutdown.inc();
}
return Err(WakeupError::ShuttingDown);
}
if self.heap.contains_key(&slot) {
let revision = self.next_revision();
let priority = WakeupPriority { when, revision };
let _ = self.heap.insert(slot, priority);
if let Some(metrics) = &mut self.metrics {
metrics.set_replaced.inc();
}
self.refresh_gauges();
Ok(WakeupSetOutcome::Replaced { revision })
} else {
if self.heap.len() >= self.wakeup_capacity {
if let Some(metrics) = &mut self.metrics {
metrics.set_error_capacity.inc();
}
return Err(WakeupError::Capacity);
}
let revision = self.next_revision();
let priority = WakeupPriority { when, revision };
let _ = self.heap.insert(slot, priority);
if let Some(metrics) = &mut self.metrics {
metrics.set_inserted.inc();
}
self.refresh_gauges();
Ok(WakeupSetOutcome::Inserted { revision })
}
}
fn cancel_wakeup(&mut self, slot: WakeupSlot) -> bool {
if self.shutting_down {
if let Some(metrics) = &mut self.metrics {
metrics.cancel_ignored_shutdown.inc();
}
return false;
}
let removed = self.heap.remove(&slot).is_some();
if let Some(metrics) = &mut self.metrics {
if removed {
metrics.cancel_removed.inc();
} else {
metrics.cancel_missed.inc();
}
}
self.refresh_gauges();
removed
}
fn next_expiry(&mut self) -> Option<Instant> {
if !self.due_now.is_empty() {
return Some(clock::now());
}
#[cfg(debug_assertions)]
self.heap.assert_consistent();
match (
self.delayed_resumes.peek().map(|resume| resume.when),
self.heap.peek().map(|(_, priority)| priority.when),
) {
(Some(a), Some(b)) => Some(a.min(b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
}
}
fn pop_due(&mut self, now: Instant) -> Option<NodeControlMsg<PData>> {
if let Some(msg) = self.due_now.pop_front() {
return Some(msg);
}
#[cfg(debug_assertions)]
self.heap.assert_consistent();
let next_resume = self.delayed_resumes.peek().map(|resume| resume.when);
let next_wakeup = self.heap.peek().map(|(_, priority)| priority.when);
let take_resume = match (next_resume, next_wakeup) {
(Some(resume_when), Some(wakeup_when)) => {
resume_when <= now && (wakeup_when > now || resume_when <= wakeup_when)
}
(Some(resume_when), None) => resume_when <= now,
(None, Some(_)) => false,
(None, None) => return None,
};
if take_resume {
let resume = self
.delayed_resumes
.pop()
.expect("resume must exist");
return Some(NodeControlMsg::ResumeData {
when: resume.when,
data: resume.data,
});
}
let next_due = self.heap.peek().map(|(_, priority)| priority.when)?;
if next_due > now {
return None;
}
let (slot, priority) = self
.heap
.pop()
.expect("due wakeup should exist");
if let Some(metrics) = &mut self.metrics {
metrics.pop_due.inc();
}
self.refresh_gauges();
Some(NodeControlMsg::Wakeup {
slot,
when: priority.when,
revision: priority.revision,
})
}
#[cfg(any(test, feature = "test-utils"))]
fn pop_next(&mut self) -> Option<(WakeupSlot, Instant, WakeupRevision)> {
#[cfg(debug_assertions)]
self.heap.assert_consistent();
let (slot, priority) = self.heap.pop()?;
if let Some(metrics) = &mut self.metrics {
metrics.pop_due.inc();
}
self.refresh_gauges();
Some((slot, priority.when, priority.revision))
}
fn begin_shutdown(&mut self, now: Instant) {
if self.shutting_down {
return;
}
self.shutting_down = true;
while let Some(resume) = self.delayed_resumes.pop() {
self.due_now.push_back(NodeControlMsg::ResumeData {
when: now,
data: resume.data,
});
}
let cleared = self.heap.len();
if let Some(metrics) = &mut self.metrics {
metrics.shutdown_cleared.add(cleared as u64);
}
self.heap.clear();
self.refresh_gauges();
}
fn is_drained(&self) -> bool {
self.due_now.is_empty() && self.delayed_resumes.is_empty() && self.heap.is_empty()
}
fn report_metrics(
&mut self,
metrics_reporter: &mut MetricsReporter,
) -> Result<(), TelemetryError> {
self.refresh_gauges();
if let Some(metrics) = &mut self.metrics {
metrics_reporter.report(metrics)?;
}
Ok(())
}
}
pub(crate) struct NodeLocalSchedulerHandle<PData> {
inner: Arc<Mutex<NodeLocalScheduler<PData>>>,
notify: Arc<Notify>,
}
impl<PData> Clone for NodeLocalSchedulerHandle<PData> {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
notify: Arc::clone(&self.notify),
}
}
}
impl<PData> NodeLocalSchedulerHandle<PData> {
pub(crate) fn new(delayed_resume_capacity: usize, wakeup_capacity: usize) -> Self {
let metrics = current_node_telemetry_handle()
.map(|telemetry| telemetry.register_metric_set::<NodeLocalWakeupSchedulerMetrics>());
Self {
inner: Arc::new(Mutex::new(NodeLocalScheduler::new_with_metrics(
delayed_resume_capacity,
wakeup_capacity,
metrics,
))),
notify: Arc::new(Notify::new()),
}
}
fn with_scheduler<R>(&self, f: impl FnOnce(&mut NodeLocalScheduler<PData>) -> R) -> R {
let mut guard = self
.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
f(&mut guard)
}
pub(crate) fn requeue_later(&self, when: Instant, data: Box<PData>) -> Result<(), Box<PData>> {
let result = self.with_scheduler(|scheduler| scheduler.requeue_later(when, data));
if result.is_ok() {
self.notify.notify_one();
}
result
}
pub(crate) fn set_wakeup(
&self,
slot: WakeupSlot,
when: Instant,
) -> Result<WakeupSetOutcome, WakeupError> {
let result = self.with_scheduler(|scheduler| scheduler.set_wakeup(slot, when));
if result.is_ok() {
self.notify.notify_one();
}
result
}
#[must_use]
pub(crate) fn cancel_wakeup(&self, slot: WakeupSlot) -> bool {
let changed = self.with_scheduler(|scheduler| scheduler.cancel_wakeup(slot));
if changed {
self.notify.notify_one();
}
changed
}
pub(crate) fn next_expiry(&self) -> Option<Instant> {
self.with_scheduler(NodeLocalScheduler::next_expiry)
}
pub(crate) fn pop_due(&self, now: Instant) -> Option<NodeControlMsg<PData>> {
self.with_scheduler(|scheduler| scheduler.pop_due(now))
}
#[cfg(any(test, feature = "test-utils"))]
pub(crate) fn pop_next(&self) -> Option<(WakeupSlot, Instant, WakeupRevision)> {
self.with_scheduler(|scheduler| scheduler.pop_next())
}
pub(crate) fn begin_shutdown(&self, now: Instant) {
self.with_scheduler(|scheduler| scheduler.begin_shutdown(now));
self.notify.notify_waiters();
}
pub(crate) fn is_drained(&self) -> bool {
self.with_scheduler(|scheduler| scheduler.is_drained())
}
pub(crate) fn report_metrics(
&self,
metrics_reporter: &mut MetricsReporter,
) -> Result<(), TelemetryError> {
self.with_scheduler(|scheduler| scheduler.report_metrics(metrics_reporter))
}
pub(crate) async fn wait_for_change(&self) {
self.notify.notified().await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
fn assert_heap_bound<PData>(scheduler: &NodeLocalScheduler<PData>) {
#[cfg(debug_assertions)]
scheduler.heap.assert_consistent();
}
fn expect_resume(msg: Option<NodeControlMsg<i32>>, expected_when: Instant, expected_data: i32) {
match msg {
Some(NodeControlMsg::ResumeData { when, data }) => {
assert_eq!(when, expected_when);
assert_eq!(*data, expected_data);
}
_ => panic!("expected resume data"),
}
}
fn expect_wakeup(
msg: Option<NodeControlMsg<i32>>,
expected_slot: WakeupSlot,
expected_when: Instant,
expected_revision: WakeupRevision,
) {
match msg {
Some(NodeControlMsg::Wakeup {
slot,
when,
revision,
}) => {
assert_eq!(slot, expected_slot);
assert_eq!(when, expected_when);
assert_eq!(revision, expected_revision);
}
_ => panic!("expected wakeup"),
}
}
#[test]
fn requeue_later_emits_the_stored_payload() {
let mut scheduler = NodeLocalScheduler::<i32>::new(2, 2);
let when = Instant::now() + Duration::from_secs(1);
assert_eq!(scheduler.requeue_later(when, Box::new(17)), Ok(()));
expect_resume(scheduler.pop_due(when), when, 17);
assert_eq!(scheduler.next_expiry(), None);
}
#[test]
fn delayed_resumes_preserve_due_time_ordering() {
let mut scheduler = NodeLocalScheduler::new(4, 2);
let now = Instant::now();
let later = now + Duration::from_secs(3);
let sooner = now + Duration::from_secs(1);
let same_time_a = now + Duration::from_secs(2);
let same_time_b = same_time_a;
assert_eq!(scheduler.requeue_later(later, Box::new(3)), Ok(()));
assert_eq!(scheduler.requeue_later(same_time_a, Box::new(1)), Ok(()));
assert_eq!(scheduler.requeue_later(same_time_b, Box::new(2)), Ok(()));
assert_eq!(scheduler.requeue_later(sooner, Box::new(0)), Ok(()));
expect_resume(scheduler.pop_due(sooner), sooner, 0);
expect_resume(scheduler.pop_due(same_time_a), same_time_a, 1);
expect_resume(scheduler.pop_due(same_time_b), same_time_b, 2);
expect_resume(scheduler.pop_due(later), later, 3);
}
#[test]
fn delayed_resume_capacity_is_enforced() {
let mut scheduler = NodeLocalScheduler::new(1, 1);
let when = Instant::now() + Duration::from_secs(1);
assert_eq!(scheduler.requeue_later(when, Box::new(1)), Ok(()));
let rejected = scheduler
.requeue_later(when, Box::new(2))
.expect_err("capacity should reject");
assert_eq!(*rejected, 2);
}
#[test]
fn rejected_requeue_returns_the_original_payload() {
let mut scheduler = NodeLocalScheduler::new(2, 1);
let now = Instant::now();
scheduler.begin_shutdown(now);
let rejected = scheduler
.requeue_later(now + Duration::from_secs(1), Box::new(99))
.expect_err("shutdown should reject");
assert_eq!(*rejected, 99);
}
#[test]
fn shutdown_makes_pending_delayed_resumes_due_immediately() {
let mut scheduler = NodeLocalScheduler::new(4, 2);
let now = Instant::now();
let later = now + Duration::from_secs(30);
assert_eq!(scheduler.requeue_later(later, Box::new(11)), Ok(()));
assert_eq!(
scheduler.requeue_later(later + Duration::from_secs(1), Box::new(12)),
Ok(())
);
scheduler.begin_shutdown(now);
expect_resume(scheduler.pop_due(now), now, 11);
expect_resume(scheduler.pop_due(now), now, 12);
assert!(scheduler.pop_due(now).is_none());
}
#[test]
fn set_wakeup_schedules_a_wakeup() {
let mut scheduler = NodeLocalScheduler::<i32>::new(2, 2);
let now = Instant::now();
let when = now + Duration::from_secs(1);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(7), when),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_heap_bound(&scheduler);
assert_eq!(scheduler.next_expiry(), Some(when));
assert!(scheduler.pop_due(now).is_none());
expect_wakeup(scheduler.pop_due(when), WakeupSlot(7), when, 0);
assert_heap_bound(&scheduler);
assert_eq!(scheduler.next_expiry(), None);
}
#[test]
fn setting_same_slot_replaces_previous_due_time() {
let mut scheduler = NodeLocalScheduler::<i32>::new(2, 2);
let now = Instant::now();
let later = now + Duration::from_secs(10);
let sooner = now + Duration::from_secs(1);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(3), later),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(3), sooner),
Ok(WakeupSetOutcome::Replaced { revision: 1 })
);
assert_heap_bound(&scheduler);
assert_eq!(scheduler.heap.len(), 1);
assert_eq!(scheduler.next_expiry(), Some(sooner));
expect_wakeup(scheduler.pop_due(sooner), WakeupSlot(3), sooner, 1);
assert_heap_bound(&scheduler);
assert!(scheduler.pop_due(later).is_none());
}
#[test]
fn cancel_wakeup_removes_pending_wakeup() {
let mut scheduler = NodeLocalScheduler::<i32>::new(2, 2);
let when = Instant::now() + Duration::from_secs(1);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(5), when),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_heap_bound(&scheduler);
assert!(scheduler.cancel_wakeup(WakeupSlot(5)));
assert_heap_bound(&scheduler);
assert!(!scheduler.cancel_wakeup(WakeupSlot(5)));
assert_eq!(scheduler.next_expiry(), None);
assert!(scheduler.pop_due(when).is_none());
}
#[test]
fn scheduler_metrics_count_wakeup_lifecycle() {
let registry = otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle::default();
let metrics: MetricSet<NodeLocalWakeupSchedulerMetrics> =
registry.register_metric_set(otel_arrow_dfe_telemetry::testing::EmptyAttributes());
let mut scheduler = NodeLocalScheduler::<i32>::new_with_metrics(1, 1, Some(metrics));
let now = Instant::now();
let later = now + Duration::from_secs(1);
let sooner = now + Duration::from_millis(10);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(0), later),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(0), sooner),
Ok(WakeupSetOutcome::Replaced { revision: 1 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(1), later),
Err(WakeupError::Capacity)
);
assert!(!scheduler.cancel_wakeup(WakeupSlot(9)));
assert!(scheduler.cancel_wakeup(WakeupSlot(0)));
assert!(!scheduler.cancel_wakeup(WakeupSlot(0)));
assert_eq!(
scheduler.set_wakeup(WakeupSlot(0), later),
Ok(WakeupSetOutcome::Inserted { revision: 2 })
);
assert!(scheduler.pop_due(now).is_none());
expect_wakeup(scheduler.pop_due(later), WakeupSlot(0), later, 2);
let metrics = scheduler.metrics.as_ref().expect("metrics enabled");
assert_eq!(metrics.set_inserted.get(), 2);
assert_eq!(metrics.set_replaced.get(), 1);
assert_eq!(metrics.set_error_capacity.get(), 1);
assert_eq!(metrics.cancel_removed.get(), 1);
assert_eq!(metrics.cancel_missed.get(), 2);
assert_eq!(metrics.pop_due.get(), 1);
assert_eq!(metrics.live.get(), 0);
assert_eq!(metrics.capacity.get(), 1);
}
#[test]
fn cancel_after_reschedule_removes_the_moved_entry() {
let mut scheduler = NodeLocalScheduler::<i32>::new(4, 4);
let now = Instant::now();
let first = now + Duration::from_secs(1);
let second = now + Duration::from_secs(10);
let third = now + Duration::from_secs(20);
let fourth = now + Duration::from_secs(30);
let moved = now + Duration::from_secs(2);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(1), first),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(2), second),
Ok(WakeupSetOutcome::Inserted { revision: 1 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(3), third),
Ok(WakeupSetOutcome::Inserted { revision: 2 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(4), fourth),
Ok(WakeupSetOutcome::Inserted { revision: 3 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(3), moved),
Ok(WakeupSetOutcome::Replaced { revision: 4 })
);
assert!(scheduler.heap.contains_key(&WakeupSlot(3)));
assert_eq!(scheduler.heap.peek().map(|(k, _)| *k), Some(WakeupSlot(1)));
assert!(scheduler.cancel_wakeup(WakeupSlot(3)));
assert_heap_bound(&scheduler);
expect_wakeup(scheduler.pop_due(first), WakeupSlot(1), first, 0);
assert!(scheduler.pop_due(moved).is_none());
expect_wakeup(scheduler.pop_due(second), WakeupSlot(2), second, 1);
expect_wakeup(scheduler.pop_due(fourth), WakeupSlot(4), fourth, 3);
assert_eq!(scheduler.next_expiry(), None);
}
#[test]
fn wakeup_capacity_is_enforced_on_distinct_live_slots() {
let mut scheduler = NodeLocalScheduler::<i32>::new(1, 1);
let when = Instant::now() + Duration::from_secs(1);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(0), when),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(1), when),
Err(WakeupError::Capacity)
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(0), when + Duration::from_secs(1)),
Ok(WakeupSetOutcome::Replaced { revision: 1 })
);
assert_heap_bound(&scheduler);
}
#[test]
fn wakeup_is_unsupported_when_capacity_is_zero() {
let mut scheduler = NodeLocalScheduler::<i32>::new(1, 0);
let when = Instant::now() + Duration::from_secs(1);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(0), when),
Err(WakeupError::Unsupported)
);
}
#[test]
fn repeated_reschedules_keep_single_heap_entry() {
let mut scheduler = NodeLocalScheduler::<i32>::new(2, 2);
let now = Instant::now();
for offset in (1..=32).rev() {
let when = now + Duration::from_secs(offset);
let outcome = scheduler
.set_wakeup(WakeupSlot(9), when)
.expect("wakeup should schedule");
let expected_revision: WakeupRevision = 32 - offset;
assert_eq!(outcome.revision(), expected_revision);
assert_heap_bound(&scheduler);
assert_eq!(scheduler.heap.len(), 1);
assert_eq!(scheduler.next_expiry(), Some(when));
}
let expected = now + Duration::from_secs(1);
expect_wakeup(scheduler.pop_due(expected), WakeupSlot(9), expected, 31);
assert_eq!(scheduler.next_expiry(), None);
}
#[test]
fn equal_deadlines_follow_schedule_sequence() {
let mut scheduler = NodeLocalScheduler::<i32>::new(4, 4);
let when = Instant::now() + Duration::from_secs(1);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(1), when),
Ok(WakeupSetOutcome::Inserted { revision: 0 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(2), when),
Ok(WakeupSetOutcome::Inserted { revision: 1 })
);
assert_eq!(
scheduler.set_wakeup(WakeupSlot(3), when),
Ok(WakeupSetOutcome::Inserted { revision: 2 })
);
assert_heap_bound(&scheduler);
expect_wakeup(scheduler.pop_due(when), WakeupSlot(1), when, 0);
expect_wakeup(scheduler.pop_due(when), WakeupSlot(2), when, 1);
expect_wakeup(scheduler.pop_due(when), WakeupSlot(3), when, 2);
assert_heap_bound(&scheduler);
}
}