Skip to main content

ixa/
plan_queue.rs

1//! Implementation details for Ixa's plan queues.
2//!
3//! [`PlanId`] is the only public item in this module. The queues store regular
4//! time-ordered plans and shutdown-time plans, sharing a single `PlanId`
5//! allocator and cancellation map.
6
7use std::cmp::Ordering;
8use std::collections::BinaryHeap;
9
10use crate::context::{Context, ExecutionPhase};
11use crate::{trace, HashMap, HashMapExt};
12
13type Callback = dyn FnOnce(&mut Context);
14type BoxedCallback = Box<Callback>;
15
16struct QueuedPlan {
17    callback: BoxedCallback,
18    is_passive: bool,
19}
20
21/// A priority queue that stores scheduled plans.
22///
23/// Regular plans are ordered by simulation time, execution phase, and plan ID.
24/// Shutdown-time plans are ordered by execution phase and plan ID; their stored
25/// time is only an internal constant and has no simulation-time meaning.
26pub(crate) struct PlanQueue {
27    queue: BinaryHeap<PlanSchedule>,
28    shutdown_queue: BinaryHeap<PlanSchedule>,
29    data_map: HashMap<u64, QueuedPlan>,
30    /// Count of scheduled plans excluding passive and shutdown-time plans.
31    regular_plan_count: usize,
32    /// The next plan ID that will be issued.
33    next_plan_id: u64,
34    /// Tracks the high water mark of plans in flight (scheduled but not yet executed).
35    /// This is the max of the two heap lengths, not of `self.data_map.len()`.
36    #[cfg(feature = "profiling")]
37    pub(crate) max_plans_in_flight: u64,
38    #[cfg(feature = "profiling")]
39    pub(crate) max_memory_in_use: u64,
40}
41
42impl PlanQueue {
43    /// Create a new empty `PlanQueue`.
44    #[must_use]
45    pub(crate) fn new() -> PlanQueue {
46        PlanQueue {
47            queue: BinaryHeap::new(),
48            shutdown_queue: BinaryHeap::new(),
49            data_map: HashMap::new(),
50            regular_plan_count: 0,
51            next_plan_id: 0,
52            #[cfg(feature = "profiling")]
53            max_plans_in_flight: 0,
54            #[cfg(feature = "profiling")]
55            max_memory_in_use: 0,
56        }
57    }
58
59    /// Add a regular plan to the queue at the specified time.
60    ///
61    /// Returns a [`PlanId`] for the newly-added plan that can be used to cancel it
62    /// if needed.
63    pub(crate) fn add_plan(
64        &mut self,
65        time: f64,
66        callback: BoxedCallback,
67        phase: ExecutionPhase,
68        is_passive: bool,
69    ) -> PlanId {
70        trace!("adding plan at {time}");
71        let plan_id = self.next_plan_id;
72        self.queue.push(PlanSchedule {
73            plan_id,
74            time,
75            phase,
76        });
77        self.data_map.insert(
78            plan_id,
79            QueuedPlan {
80                callback,
81                is_passive,
82            },
83        );
84        if !is_passive {
85            self.regular_plan_count += 1;
86        }
87        self.next_plan_id += 1;
88        self.update_profiling_high_water_marks();
89
90        PlanId(plan_id)
91    }
92
93    /// Add a shutdown-time plan.
94    ///
95    /// Shutdown-time plans have no simulation time. They are ordered by phase and
96    /// plan ID.
97    pub(crate) fn add_shutdown_plan(
98        &mut self,
99        callback: BoxedCallback,
100        phase: ExecutionPhase,
101    ) -> PlanId {
102        trace!("adding shutdown-time plan");
103        let plan_id = self.next_plan_id;
104        self.shutdown_queue.push(PlanSchedule {
105            plan_id,
106            time: 0.0,
107            phase,
108        });
109        self.data_map.insert(
110            plan_id,
111            QueuedPlan {
112                callback,
113                is_passive: true,
114            },
115        );
116        self.next_plan_id += 1;
117        self.update_profiling_high_water_marks();
118
119        PlanId(plan_id)
120    }
121
122    /// Cancel a plan that has been added to either queue.
123    pub(crate) fn cancel_plan(&mut self, plan_id: &PlanId) -> Option<BoxedCallback> {
124        trace!("cancel plan {plan_id:?}");
125        // Delete the plan from the map, but leave in the heap. It will be skipped
126        // when its heap entry reaches the root.
127        self.data_map.remove(&plan_id.0).map(|queued_plan| {
128            if !queued_plan.is_passive {
129                self.regular_plan_count -= 1;
130            }
131            queued_plan.callback
132        })
133    }
134
135    /// Return the time the next plan is scheduled for, if there is one.
136    #[must_use]
137    pub(crate) fn next_time(&mut self) -> Option<f64> {
138        while let Some(entry) = self.queue.peek() {
139            // We only want to report the time if the plan has not been canceled.
140            if self.data_map.contains_key(&entry.plan_id) {
141                return Some(entry.time);
142            }
143            // Trim the canceled plan.
144            self.queue.pop();
145        }
146        None
147    }
148
149    /// Completely empties the queue, including the plans scheduled at shutdown time.
150    #[allow(dead_code)]
151    pub(crate) fn clear(&mut self) {
152        self.data_map.clear();
153        self.queue.clear();
154        self.shutdown_queue.clear();
155        self.regular_plan_count = 0;
156        self.next_plan_id = 0;
157    }
158
159    /// Retrieve the earliest regular plan only if the regular queue is active.
160    ///
161    /// The queue is active while at least one non-passive regular plan is
162    /// scheduled. The returned plan is simply the next regular plan by time,
163    /// phase, and plan ID; it may be passive. If no live non-passive regular
164    /// plan is scheduled, this returns `None` without removing passive plans
165    /// from the regular queue.
166    pub(crate) fn pop_next_if_active(&mut self) -> Option<Plan> {
167        if self.regular_plan_count == 0 {
168            return None;
169        }
170
171        trace!("getting next plan");
172        loop {
173            // The `pop` should be infallible when the plan count is positive unless the
174            // queue invariants have been violated.
175            let entry = self
176                .queue
177                .pop()
178                .expect("Ixa internal error: plan count was positive but no plan was available");
179
180            // Discard any cancelled plans we encounter.
181            if let Some(queued_plan) = self.data_map.remove(&entry.plan_id) {
182                if !queued_plan.is_passive {
183                    self.regular_plan_count -= 1;
184                }
185                return Some(Plan {
186                    time: entry.time,
187                    data: queued_plan.callback,
188                });
189            }
190        }
191    }
192
193    /// Retrieve the earliest regular plan only if it is scheduled at `time`.
194    ///
195    /// Returns `None` without removing a future plan if the next regular plan is
196    /// later than `time`.
197    pub(crate) fn pop_next_at(&mut self, time: f64) -> Option<Plan> {
198        loop {
199            match self.queue.peek() {
200                // Trim any cancelled plans
201                Some(entry) if !self.data_map.contains_key(&entry.plan_id) => {
202                    self.queue.pop();
203                }
204
205                // Return only if the plan is scheduled for the given time
206                Some(entry) if entry.time == time => {
207                    // The preceding `peek` proves that `pop` will return an entry.
208                    let entry = self.queue.pop().unwrap();
209                    let queued_plan = self
210                        .data_map
211                        .remove(&entry.plan_id)
212                        .expect("Ixa internal error: live plan has no callback");
213                    if !queued_plan.is_passive {
214                        self.regular_plan_count -= 1;
215                    }
216                    return Some(Plan {
217                        time,
218                        data: queued_plan.callback,
219                    });
220                }
221
222                // There are no plans scheduled at the given time
223                _ => return None,
224            }
225        }
226    }
227
228    /// Retrieve the next shutdown-time plan.
229    ///
230    /// Returns the next shutdown-time plan if it exists or else `None` if the
231    /// shutdown-time queue is empty.
232    pub(crate) fn pop_next_shutdown(&mut self) -> Option<Plan> {
233        trace!("getting next shutdown-time plan");
234        std::iter::from_fn(|| self.shutdown_queue.pop()).find_map(|entry| {
235            // If there's no `data_map` entry, the plan has been canceled, so discard
236            // and pop another plan.
237            let queued_plan = self.data_map.remove(&entry.plan_id)?;
238            if !queued_plan.is_passive {
239                self.regular_plan_count -= 1;
240            }
241            Some(Plan {
242                time: entry.time,
243                data: queued_plan.callback,
244            })
245        })
246    }
247
248    fn update_profiling_high_water_marks(&mut self) {
249        #[cfg(feature = "profiling")]
250        {
251            let plans_in_flight = self.queue.len() + self.shutdown_queue.len();
252            self.max_plans_in_flight = self.max_plans_in_flight.max(plans_in_flight as u64);
253            self.max_memory_in_use = self
254                .max_memory_in_use
255                .max(self.estimated_memory_in_use() as u64);
256        }
257    }
258
259    #[cfg(feature = "profiling")]
260    fn estimated_memory_in_use(&self) -> usize {
261        let queue_bytes =
262            (self.queue.capacity() + self.shutdown_queue.capacity()) * size_of::<PlanSchedule>();
263
264        let map_entry_bytes = self.data_map.capacity() * size_of::<(u64, QueuedPlan)>();
265
266        queue_bytes + map_entry_bytes
267    }
268}
269
270impl Default for PlanQueue {
271    fn default() -> Self {
272        Self::new()
273    }
274}
275
276/// A time, id, and phase object used to order plans in a [`PlanQueue`].
277///
278/// Regular [`PlanSchedule`] objects are sorted in increasing order of time,
279/// phase, and then plan id. Shutdown-time schedules all have the same internal
280/// time and are therefore sorted by phase and then plan id.
281#[derive(PartialEq, Debug, Clone, Copy)]
282pub(crate) struct PlanSchedule {
283    pub plan_id: u64,
284    pub time: f64,
285    pub phase: ExecutionPhase,
286}
287
288impl Eq for PlanSchedule {}
289
290impl PartialOrd for PlanSchedule {
291    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
292        Some(self.cmp(other))
293    }
294}
295
296/// Entry objects are ordered in increasing order by time, phase, and then plan id.
297impl Ord for PlanSchedule {
298    fn cmp(&self, other: &Self) -> Ordering {
299        let time_ordering = self
300            .time
301            .partial_cmp(&other.time)
302            .expect("Ixa internal error: plan queue contains a NaN time")
303            .reverse();
304        match time_ordering {
305            Ordering::Equal => {
306                let phase_ordering = self.phase.cmp(&other.phase).reverse();
307                match phase_ordering {
308                    Ordering::Equal => self.plan_id.cmp(&other.plan_id).reverse(),
309                    _ => phase_ordering,
310                }
311            }
312            _ => time_ordering,
313        }
314    }
315}
316
317/// A unique identifier for a plan scheduled on a [`Context`].
318///
319/// `PlanId` values are returned by [`Context::add_plan`] and related scheduling
320/// methods, and can be passed to [`Context::cancel_plan`].
321///
322/// # Examples
323///
324/// ```
325/// use ixa::{Context, PlanId};
326///
327/// let mut context = Context::new();
328/// let plan_id: PlanId = context.add_plan(1.0, |_| {});
329/// context.cancel_plan(&plan_id);
330/// ```
331#[derive(Clone, Copy, Debug, Hash, Eq, PartialEq)]
332pub struct PlanId(pub(crate) u64);
333
334/// A plan that holds a callback intended to be executed at the specified time.
335pub(crate) struct Plan {
336    pub time: f64,
337    pub data: BoxedCallback,
338}
339
340#[cfg(test)]
341#[allow(clippy::float_cmp)]
342mod tests {
343    use std::cell::RefCell;
344    use std::rc::Rc;
345
346    use super::PlanQueue;
347    use crate::context::{Context, ExecutionPhase};
348
349    fn callback(value: u32, observed: Rc<RefCell<Vec<u32>>>) -> Box<dyn FnOnce(&mut Context)> {
350        Box::new(move |_| observed.borrow_mut().push(value))
351    }
352
353    fn run_plan(plan: super::Plan, context: &mut Context) {
354        (plan.data)(context);
355    }
356
357    #[test]
358    fn empty_queue() {
359        let mut plan_queue = PlanQueue::new();
360        assert!(plan_queue.pop_next_if_active().is_none());
361    }
362
363    #[test]
364    fn add_plans() {
365        let observed = Rc::new(RefCell::new(Vec::new()));
366        let mut context = Context::new();
367        let mut plan_queue = PlanQueue::new();
368        plan_queue.add_plan(
369            1.0,
370            callback(1, Rc::clone(&observed)),
371            ExecutionPhase::Normal,
372            false,
373        );
374        plan_queue.add_plan(
375            3.0,
376            callback(3, Rc::clone(&observed)),
377            ExecutionPhase::Normal,
378            false,
379        );
380        plan_queue.add_plan(
381            2.0,
382            callback(2, Rc::clone(&observed)),
383            ExecutionPhase::Normal,
384            false,
385        );
386        assert_eq!(plan_queue.next_time(), Some(1.0));
387
388        let next_plan = plan_queue.pop_next_if_active().unwrap();
389        assert_eq!(next_plan.time, 1.0);
390        run_plan(next_plan, &mut context);
391
392        assert_eq!(plan_queue.next_time(), Some(2.0));
393        let next_plan = plan_queue.pop_next_if_active().unwrap();
394        assert_eq!(next_plan.time, 2.0);
395        run_plan(next_plan, &mut context);
396
397        assert_eq!(plan_queue.next_time(), Some(3.0));
398        let next_plan = plan_queue.pop_next_if_active().unwrap();
399        assert_eq!(next_plan.time, 3.0);
400        run_plan(next_plan, &mut context);
401
402        assert!(plan_queue.pop_next_if_active().is_none());
403        assert_eq!(*observed.borrow(), vec![1, 2, 3]);
404    }
405
406    #[test]
407    fn add_plans_at_same_time_with_same_phase() {
408        let observed = Rc::new(RefCell::new(Vec::new()));
409        let mut context = Context::new();
410        let mut plan_queue = PlanQueue::new();
411        plan_queue.add_plan(
412            1.0,
413            callback(1, Rc::clone(&observed)),
414            ExecutionPhase::Normal,
415            false,
416        );
417        plan_queue.add_plan(
418            1.0,
419            callback(2, Rc::clone(&observed)),
420            ExecutionPhase::Normal,
421            false,
422        );
423
424        let next_plan = plan_queue.pop_next_if_active().unwrap();
425        assert_eq!(next_plan.time, 1.0);
426        run_plan(next_plan, &mut context);
427        let next_plan = plan_queue.pop_next_if_active().unwrap();
428        assert_eq!(next_plan.time, 1.0);
429        run_plan(next_plan, &mut context);
430
431        assert!(plan_queue.pop_next_if_active().is_none());
432        assert_eq!(*observed.borrow(), vec![1, 2]);
433    }
434
435    #[test]
436    fn add_plans_at_same_time_with_different_phase() {
437        let observed = Rc::new(RefCell::new(Vec::new()));
438        let mut context = Context::new();
439        let mut plan_queue = PlanQueue::new();
440        plan_queue.add_plan(
441            1.0,
442            callback(1, Rc::clone(&observed)),
443            ExecutionPhase::Normal,
444            false,
445        );
446        plan_queue.add_plan(
447            1.0,
448            callback(2, Rc::clone(&observed)),
449            ExecutionPhase::First,
450            false,
451        );
452
453        let next_plan = plan_queue.pop_next_if_active().unwrap();
454        assert_eq!(next_plan.time, 1.0);
455        run_plan(next_plan, &mut context);
456        let next_plan = plan_queue.pop_next_if_active().unwrap();
457        assert_eq!(next_plan.time, 1.0);
458        run_plan(next_plan, &mut context);
459
460        assert!(plan_queue.pop_next_if_active().is_none());
461        assert_eq!(*observed.borrow(), vec![2, 1]);
462    }
463
464    #[test]
465    fn cancel_plan() {
466        let observed = Rc::new(RefCell::new(Vec::new()));
467        let mut context = Context::new();
468        let mut plan_queue = PlanQueue::new();
469        plan_queue.add_plan(
470            1.0,
471            callback(1, Rc::clone(&observed)),
472            ExecutionPhase::Normal,
473            false,
474        );
475        let plan_to_cancel = plan_queue.add_plan(
476            2.0,
477            callback(2, Rc::clone(&observed)),
478            ExecutionPhase::Normal,
479            false,
480        );
481        plan_queue.add_plan(
482            3.0,
483            callback(3, Rc::clone(&observed)),
484            ExecutionPhase::Normal,
485            false,
486        );
487        plan_queue.cancel_plan(&plan_to_cancel);
488
489        let next_plan = plan_queue.pop_next_if_active().unwrap();
490        assert_eq!(next_plan.time, 1.0);
491        run_plan(next_plan, &mut context);
492
493        let next_plan = plan_queue.pop_next_if_active().unwrap();
494        assert_eq!(next_plan.time, 3.0);
495        run_plan(next_plan, &mut context);
496
497        assert!(plan_queue.pop_next_if_active().is_none());
498        assert_eq!(*observed.borrow(), vec![1, 3]);
499    }
500
501    #[test]
502    fn passive_only_plans_do_not_pop_during_active_execution() {
503        let observed = Rc::new(RefCell::new(Vec::new()));
504        let mut plan_queue = PlanQueue::new();
505        plan_queue.add_plan(
506            1.0,
507            callback(1, Rc::clone(&observed)),
508            ExecutionPhase::Normal,
509            true,
510        );
511
512        assert_eq!(plan_queue.regular_plan_count, 0);
513        assert!(plan_queue.pop_next_if_active().is_none());
514        assert_eq!(plan_queue.next_time(), Some(1.0));
515        assert!(observed.borrow().is_empty());
516    }
517
518    #[test]
519    fn passive_plan_can_pop_before_later_non_passive_plan() {
520        let observed = Rc::new(RefCell::new(Vec::new()));
521        let mut context = Context::new();
522        let mut plan_queue = PlanQueue::new();
523        plan_queue.add_plan(
524            1.0,
525            callback(1, Rc::clone(&observed)),
526            ExecutionPhase::Normal,
527            true,
528        );
529        plan_queue.add_plan(
530            2.0,
531            callback(2, Rc::clone(&observed)),
532            ExecutionPhase::Normal,
533            false,
534        );
535
536        assert_eq!(plan_queue.regular_plan_count, 1);
537        let next_plan = plan_queue.pop_next_if_active().unwrap();
538        assert_eq!(next_plan.time, 1.0);
539        run_plan(next_plan, &mut context);
540        assert_eq!(plan_queue.regular_plan_count, 1);
541
542        let next_plan = plan_queue.pop_next_if_active().unwrap();
543        assert_eq!(next_plan.time, 2.0);
544        run_plan(next_plan, &mut context);
545        assert_eq!(plan_queue.regular_plan_count, 0);
546
547        assert_eq!(*observed.borrow(), vec![1, 2]);
548    }
549
550    #[test]
551    fn canceling_non_passive_plan_decrements_regular_count() {
552        let observed = Rc::new(RefCell::new(Vec::new()));
553        let mut plan_queue = PlanQueue::new();
554        let non_passive = plan_queue.add_plan(
555            1.0,
556            callback(1, Rc::clone(&observed)),
557            ExecutionPhase::Normal,
558            false,
559        );
560        let passive = plan_queue.add_plan(
561            2.0,
562            callback(2, Rc::clone(&observed)),
563            ExecutionPhase::Normal,
564            true,
565        );
566
567        assert_eq!(plan_queue.regular_plan_count, 1);
568        plan_queue.cancel_plan(&passive);
569        assert_eq!(plan_queue.regular_plan_count, 1);
570        plan_queue.cancel_plan(&non_passive);
571        assert_eq!(plan_queue.regular_plan_count, 0);
572    }
573
574    #[test]
575    fn shutdown_plans_do_not_affect_regular_count() {
576        let observed = Rc::new(RefCell::new(Vec::new()));
577        let mut context = Context::new();
578        let mut plan_queue = PlanQueue::new();
579        plan_queue.add_shutdown_plan(callback(1, Rc::clone(&observed)), ExecutionPhase::Normal);
580
581        assert_eq!(plan_queue.regular_plan_count, 0);
582        let next_plan = plan_queue.pop_next_shutdown().unwrap();
583        run_plan(next_plan, &mut context);
584        assert_eq!(plan_queue.regular_plan_count, 0);
585        assert_eq!(*observed.borrow(), vec![1]);
586    }
587
588    #[test]
589    fn next_time_ignores_canceled_root() {
590        let observed = Rc::new(RefCell::new(Vec::new()));
591        let mut plan_queue = PlanQueue::new();
592        let plan_to_cancel = plan_queue.add_plan(
593            1.0,
594            callback(1, Rc::clone(&observed)),
595            ExecutionPhase::Normal,
596            false,
597        );
598        plan_queue.add_plan(
599            2.0,
600            callback(2, Rc::clone(&observed)),
601            ExecutionPhase::Normal,
602            false,
603        );
604
605        plan_queue.cancel_plan(&plan_to_cancel);
606
607        assert_eq!(plan_queue.next_time(), Some(2.0));
608    }
609
610    #[test]
611    fn pop_next_at_leaves_future_plan_in_queue() {
612        let observed = Rc::new(RefCell::new(Vec::new()));
613        let mut context = Context::new();
614        let mut plan_queue = PlanQueue::new();
615        plan_queue.add_plan(
616            2.0,
617            callback(2, Rc::clone(&observed)),
618            ExecutionPhase::Normal,
619            false,
620        );
621
622        assert!(plan_queue.pop_next_at(1.0).is_none());
623
624        let next_plan = plan_queue.pop_next_if_active().unwrap();
625        assert_eq!(next_plan.time, 2.0);
626        run_plan(next_plan, &mut context);
627        assert_eq!(*observed.borrow(), vec![2]);
628    }
629
630    #[test]
631    fn shutdown_plans_use_phase_and_fifo_order() {
632        let observed = Rc::new(RefCell::new(Vec::new()));
633        let mut context = Context::new();
634        let mut plan_queue = PlanQueue::new();
635        plan_queue.add_shutdown_plan(callback(3, Rc::clone(&observed)), ExecutionPhase::Last);
636        plan_queue.add_shutdown_plan(callback(1, Rc::clone(&observed)), ExecutionPhase::First);
637        plan_queue.add_shutdown_plan(callback(2, Rc::clone(&observed)), ExecutionPhase::Normal);
638        plan_queue.add_shutdown_plan(callback(4, Rc::clone(&observed)), ExecutionPhase::Last);
639
640        while let Some(plan) = plan_queue.pop_next_shutdown() {
641            run_plan(plan, &mut context);
642        }
643
644        assert_eq!(*observed.borrow(), vec![1, 2, 3, 4]);
645    }
646
647    #[test]
648    fn plan_ids_are_shared_between_regular_and_shutdown_queues() {
649        let observed = Rc::new(RefCell::new(Vec::new()));
650        let mut plan_queue = PlanQueue::new();
651        let regular_id = plan_queue.add_plan(
652            1.0,
653            callback(1, Rc::clone(&observed)),
654            ExecutionPhase::Normal,
655            false,
656        );
657        let shutdown_id =
658            plan_queue.add_shutdown_plan(callback(2, Rc::clone(&observed)), ExecutionPhase::Normal);
659
660        assert_ne!(regular_id, shutdown_id);
661        assert_eq!(regular_id.0, 0);
662        assert_eq!(shutdown_id.0, 1);
663    }
664}