Skip to main content

detcore/scheduler/
runqueue.rs

1/*
2 * Copyright (c) Meta Platforms, Inc. and affiliates.
3 * All rights reserved.
4 *
5 * This source code is licensed under the BSD-style license found in the
6 * LICENSE file in the root directory of this source tree.
7 */
8
9//! The main queue for runnable tasks.
10//!
11//! Tasks are selected from the queue based on 2 factors:
12//! 1. Their priority
13//! 2. Their round-robin order.
14//!
15//! The factors are compared in that order. Round-robin orders monotonically increase across
16//! the entire queue; a task is assigned an order at insertion time.
17//!
18//! Round-robin orders can also be negative when the `push_front` method is
19//! used. This "skips the line" within the priority level. This is mostly
20//! relevant in non-chaos modes, where all threads have the same priority. In
21//! this case, time-based events should skip the line and end-up in the front of
22//! the queue.
23//!
24//! In addition, some priority values are reserved, such as a high priority
25//! for IO eager polling.
26//!
27//! # Polling strategy
28//!
29//! We employ a polling strategy for guest threads that *would* blocked, but where we
30//! don't model precisely what conditions they're waiting for.  Anywhere we have a precise
31//! model of inter-thread dependencies (e.g. futexes), we can sleep a thread until we
32//! encounter the matching event that will wake it.  But for Linux features that we don't
33//! model 100% precisely, polling is a way to remain agnostic as to the exact
34//! dependencies, but still support these blocking behaviors deterministically.
35//!
36//! The older DetTrace system used polling, but it would poll every time through the
37//! round robin queue, which can create extremely bad performance with many threads
38//! polling an unbounded number of times.  We can greatly improve the performance by
39//! polling only at less-frequent, but still deterministically-defined intervals, such
40//! as when we think we're out of "productive" work to do.
41//!
42//! Thererefore we have special handling for scheduling polling tasks. When
43//! initially queued, there is exponential backoff in priority with the number
44//! of attempts. After enough queueing operations are performed, however,
45//! polling tasks are upgraded to their original priority to prevent complete
46//! starvation. The frequency of upgrades is controlled by `POLLING_UPGRADE_INTERVAL`.
47
48use std::cmp::Ordering;
49use std::collections::BTreeMap;
50use std::fmt;
51use std::fmt::Display;
52
53use rand::RngExt as _;
54use rand::SeedableRng;
55use rand::distr::uniform::SampleUniform;
56use rand_pcg::Pcg64Mcg;
57
58use crate::config::SchedHeuristic;
59use crate::detlog;
60use crate::types::DetTid;
61
62/// The user-accessible priority of a thread. Lowest runs first.
63pub type Priority = u64;
64
65const EAGER_IO_REPOLL_PRIORITY: Priority = Priority::MIN;
66
67/// The lowest/highest priority a thread can have.
68pub const FIRST_PRIORITY: Priority = EAGER_IO_REPOLL_PRIORITY + 1;
69
70/// The last/lowest (numerically largest) priority a thread can have.
71pub const LAST_PRIORITY: Priority = 10000;
72
73/// A high priority for threads the replayer DOES want to run.
74pub const REPLAY_FOREGROUND_PRIORITY: Priority = FIRST_PRIORITY;
75
76/// A low priority given to threads the replayer does NOT want to run.
77pub const REPLAY_DEFERRED_PRIORITY: Priority = LAST_PRIORITY - 1;
78
79/// The default priority for a thread. If chaos mode is not enabled, all threads have this
80/// priority.
81pub const DEFAULT_PRIORITY: Priority = 1000;
82
83/// Whether the priority is user-accessible. We use some values for special
84/// purposes; these shouldn't be set by the user.
85pub fn is_ordinary_priority(prio: Priority) -> bool {
86    (FIRST_PRIORITY..=LAST_PRIORITY).contains(&prio)
87}
88
89/// Deterministically transform 64 bits of entropy into a random user-settable
90/// priority.
91pub fn entropy_to_priority(entropy: u64) -> Priority {
92    let range = LAST_PRIORITY - FIRST_PRIORITY + 1;
93    let offset = entropy % range;
94    FIRST_PRIORITY + offset
95}
96
97/// The round robin turn of threads within a given priority level. Lowest runs
98/// first. Both negative and positive values are used to allow insertion of a
99/// thread at both the "front" and "back" of a priority level.
100type RoundRobinTurn = i64;
101
102/// The key into the priority queue that uniquely determines what to run next.
103/// Priorities that compare lower run first.
104#[derive(Debug, Copy, Clone)]
105pub struct PrioritizedOrder {
106    priority: Priority,
107    turn: RoundRobinTurn,
108}
109
110// These match the derived definitions, but clearly display our intention:
111impl Ord for PrioritizedOrder {
112    fn cmp(&self, other: &Self) -> Ordering {
113        self.priority
114            .cmp(&other.priority)
115            .then(self.turn.cmp(&other.turn))
116    }
117}
118impl PartialOrd for PrioritizedOrder {
119    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
120        Some(self.cmp(other))
121    }
122}
123impl PartialEq for PrioritizedOrder {
124    fn eq(&self, other: &Self) -> bool {
125        self.cmp(other).is_eq()
126    }
127}
128impl Eq for PrioritizedOrder {}
129
130impl fmt::Display for PrioritizedOrder {
131    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
132        write!(f, "(p: {}, t: {})", self.priority, self.turn)
133    }
134}
135
136/// After queueing this many new tasks, perform a "poll upgrade," in which
137/// we upgrade all outstanding polling tasks their original priority levels,
138/// temporarily negating backoff behavior.
139const POLLING_UPGRADE_INTERVAL: u64 = 200;
140
141#[derive(Debug, Copy, Clone)]
142struct QueueValue {
143    tid: DetTid,
144    /// Upgrade to this priority during polling upgrades
145    poll_upgrade: Option<Priority>,
146}
147
148/// One suspended queue entry. Global yield/random-selection state stays live.
149#[derive(Debug, Clone)]
150pub(super) struct SuspendedRunQueueEntry {
151    key: PrioritizedOrder,
152    value: QueueValue,
153    persistent_priority: Priority,
154}
155
156#[derive(Debug, Clone)]
157pub struct RunQueue {
158    /// We use a "flattened" queue (rather than a Priority -> Vec<DetTid> map)
159    /// to simplify peek/pop logic: there's no need to ignore clear empty
160    /// from unused priority levels vectors. This could also reduce allocator
161    /// pressure. Also, each thread having a clear global key makes it easier to
162    /// change their priorities after they are in the queue.
163    ///
164    /// Additionally, we use a TreeMap rather than a Heap to ease removing /
165    /// inserting random values for poll upgrades. std::BinaryHeap would require
166    /// destroying/re-allocating the entire structure to do this.
167    queue: BTreeMap<PrioritizedOrder, QueueValue>,
168
169    // We use global turn counters across all priority levels. This foregoes
170    // the need for an extra data structure to track them, while also ensuring
171    // unique keys for every insertion. Because of this, we never need to alter
172    // a turn value when altering the priority level of a thread.
173    last_back_turn: RoundRobinTurn,
174    last_front_turn: RoundRobinTurn,
175
176    /// Used to lock the queue from other changes while we are tentatively popping from it, and also
177    /// cache the result.
178    tentative_selection: Option<DetTid>,
179    tentative_selection_is_exact: bool,
180
181    /// A thread that explicitly yielded must not be selected again until some
182    /// other runnable thread receives a turn.
183    yielded_skip: Option<DetTid>,
184
185    // TODO: The following fields need to be properly abstracted into separate types of run queues.
186    /// Which scheduling strategy shall we use.
187    sched_strategy: SchedHeuristic,
188    prng: Pcg64Mcg,
189
190    sticky_random_param: f64,
191    sticky_random_selection: Option<DetTid>,
192}
193
194/// A multi-line print of the runqueue.
195impl fmt::Display for RunQueue {
196    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
197        writeln!(
198            f,
199            "Run queue, size={}, last_back_turn={}, last_front_turn={}:",
200            self.queue.len(),
201            self.last_back_turn,
202            self.last_front_turn,
203        )?;
204        for x in self.queue.iter() {
205            writeln!(f, "    {:.500?}", x)?;
206        }
207        Ok(())
208    }
209}
210
211impl RunQueue {
212    /// Create a new RunQueue.
213    pub fn new(ss: SchedHeuristic, seed: u64, srp: f64) -> Self {
214        detlog!("SCHEDRAND: seeding scheduler runqueue with seed {}", seed);
215        Self {
216            queue: BTreeMap::new(),
217            // For clarity, 0 is unused so that positive/negative == back/front:
218            last_back_turn: 0,
219            last_front_turn: 0,
220            sched_strategy: ss,
221            tentative_selection: None,
222            tentative_selection_is_exact: false,
223            yielded_skip: None,
224            prng: Pcg64Mcg::seed_from_u64(seed),
225            sticky_random_param: srp,
226            sticky_random_selection: None,
227        }
228    }
229
230    fn push_safety_check(&self, tid: DetTid) {
231        if cfg!(debug_assertions) {
232            // Expensive.
233            for qv in self.queue.values() {
234                if qv.tid == tid {
235                    panic!(
236                        "Invariant violation! Tried to add {} to runqueue, but it's already present:\n {:?}",
237                        tid, self
238                    );
239                }
240            }
241        }
242    }
243
244    // Return the numerically least Priority value in the run_queue, or None if the queue is empty.
245    pub fn first_priority(&self) -> Option<Priority> {
246        let (k, _) = self.queue.first_key_value()?;
247        Some(k.priority)
248    }
249
250    /// True if any thread other than `exclude` is runnable at ordinary
251    /// (non-poller) priority. This is the "deterministic work still runnable"
252    /// test used to decide whether an asynchronous signal delivery must defer to
253    /// guest work that was already scheduled. Read-only, so it is safe to call
254    /// while a tentative_pop selection is in progress.
255    pub fn has_runnable_besides(&self, exclude: DetTid) -> bool {
256        self.queue
257            .iter()
258            .any(|(k, v)| v.tid != exclude && k.priority < LAST_PRIORITY)
259    }
260
261    /// True while a `tentative_pop`/commit transaction is underway, i.e. the
262    /// daemon has peeked a selection and may have released the scheduler lock
263    /// across an await. Fatal backend publication uses this read-only query to
264    /// close the transaction once before waking consuming cleanup. Ordinary
265    /// run-queue admissions and removals are always buffered and
266    /// applied at the deterministic `step2` drain (window guaranteed closed), so
267    /// a handler never needs to ask whether a window is live — see
268    /// `Scheduler::admit_to_run_queue` / `deschedule_or_defer`. Read-only.
269    pub fn tentative_pop_in_progress(&self) -> bool {
270        self.tentative_selection.is_some()
271    }
272
273    /// Push a thread to the back of the specified priority. Return the
274    /// resulting overall position in the queue.
275    ///
276    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
277    pub fn push_back(&mut self, tid: DetTid, priority: Priority) -> PrioritizedOrder {
278        assert!(self.tentative_selection.is_none());
279        self.push_safety_check(tid);
280        if !is_ordinary_priority(priority) {
281            panic!("This is not an acceptable priority value: {}", priority);
282        }
283        self.push_back_inner(tid, priority, None)
284    }
285
286    /// Requeue an explicitly yielding thread at its persistent priority while
287    /// excluding it from the next selection. The exclusion, rather than the
288    /// queue key, makes this a one-turn operation under every heuristic.
289    pub fn push_yielded(&mut self, tid: DetTid, priority: Priority) -> PrioritizedOrder {
290        assert!(self.yielded_skip.is_none());
291        self.yielded_skip = Some(tid);
292        self.push_back(tid, priority)
293    }
294
295    /// Push a polling thread. The priority level is an exponential backoff from
296    /// the given `normal_priority` value. The pushed thread will also
297    /// participatein "poll upgrades" in which periodically polling threads are
298    /// re-boosted to their original `normal_priority` values.
299    ///
300    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
301    pub fn push_poller(
302        &mut self,
303        tid: DetTid,
304        normal_priority: Priority,
305        poll_attempt: u32,
306    ) -> PrioritizedOrder {
307        assert!(self.tentative_selection.is_none());
308        self.push_safety_check(tid);
309        // Exponential backoff in priority, up to LAST_PRIORITY:
310        let priority = 1u64
311            .checked_shl(poll_attempt)
312            .and_then(|f| f.checked_mul(normal_priority))
313            .unwrap_or(Priority::MAX)
314            .min(LAST_PRIORITY);
315        // Upgrade back to original priority:
316        self.push_back_inner(tid, priority, Some(normal_priority))
317    }
318
319    fn push_back_inner(
320        &mut self,
321        tid: DetTid,
322        priority: Priority,
323        poll_upgrade: Option<Priority>,
324    ) -> PrioritizedOrder {
325        self.last_back_turn += 1;
326        let turn = self.last_back_turn;
327        let prio = PrioritizedOrder { priority, turn };
328        self.push_inner(tid, prio, poll_upgrade)
329    }
330
331    /// Push a thread to the front of the specified priority. `push_back` should
332    /// be used unless special circumstances call for `push_front`. Return the
333    /// resulting overall position in the queue.
334    ///
335    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
336    pub fn push_front(&mut self, tid: DetTid, priority: Priority) -> PrioritizedOrder {
337        assert!(self.tentative_selection.is_none());
338        self.push_safety_check(tid);
339        assert!(is_ordinary_priority(priority));
340        self.push_front_inner(tid, priority, None)
341    }
342
343    /// Workaround for eager io repolling: this will send the thread to the
344    /// absolute front of the queue. Return the resulting overall position in
345    /// the queue.
346    ///
347    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
348    pub fn push_eager_io_repoll(&mut self, tid: DetTid) -> PrioritizedOrder {
349        assert!(self.tentative_selection.is_none());
350        self.push_safety_check(tid);
351        let priority = EAGER_IO_REPOLL_PRIORITY;
352        self.push_front_inner(tid, priority, None)
353    }
354
355    fn push_front_inner(
356        &mut self,
357        tid: DetTid,
358        priority: Priority,
359        poll_upgrade: Option<Priority>,
360    ) -> PrioritizedOrder {
361        self.last_front_turn -= 1;
362        let turn = self.last_front_turn;
363        let prio = PrioritizedOrder { priority, turn };
364        self.push_inner(tid, prio, poll_upgrade)
365    }
366
367    fn push_inner(
368        &mut self,
369        tid: DetTid,
370        prio: PrioritizedOrder,
371        poll_upgrade: Option<Priority>,
372    ) -> PrioritizedOrder {
373        let qval = QueueValue { tid, poll_upgrade };
374        let old = self.queue.insert(prio, qval);
375        assert!(old.is_none()); // last_*_turn should be monotonic
376        self.check_poll_upgrade();
377        prio
378    }
379
380    /// Read-only: this is ok while locked by tentative_pop.
381    pub fn is_empty(&self) -> bool {
382        self.queue.is_empty()
383    }
384
385    /// Read-only: this is ok while locked by tentative_pop.
386    pub fn len(&self) -> usize {
387        self.queue.len()
388    }
389
390    /// Read-only: this is ok while locked by tentative_pop.
391    pub fn tids(&self) -> impl Iterator<Item = &DetTid> {
392        self.queue.values().map(|v| &v.tid)
393    }
394
395    /// Read-only: this is ok while locked by tentative_pop.
396    pub fn contains_tid(&self, tid: DetTid) -> bool {
397        self.tids().any(|t| t == &tid)
398    }
399
400    pub(super) fn suspend(
401        &mut self,
402        tid: DetTid,
403        persistent_priority: Priority,
404    ) -> Option<SuspendedRunQueueEntry> {
405        assert!(self.tentative_selection.is_none());
406        let key = *self.queue.iter().find(|(_, v)| v.tid == tid)?.0;
407        let value = self.queue.remove(&key).expect("located queue entry");
408        Some(SuspendedRunQueueEntry {
409            key,
410            value,
411            persistent_priority,
412        })
413    }
414
415    pub(super) fn restore(
416        &mut self,
417        mut entry: SuspendedRunQueueEntry,
418        current_priority: Priority,
419    ) {
420        assert!(self.tentative_selection.is_none());
421        assert!(!self.contains_tid(entry.value.tid));
422        if entry.persistent_priority != current_priority {
423            entry.key.priority = current_priority;
424            if entry.value.poll_upgrade.is_some() {
425                entry.value.poll_upgrade = Some(current_priority);
426            }
427        }
428        assert!(self.queue.insert(entry.key, entry.value).is_none());
429    }
430
431    /// Remove `tid` from the queue, returning true if removal ocurred.
432    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
433    pub fn remove_tid(&mut self, tid: DetTid) -> bool {
434        assert!(self.tentative_selection.is_none());
435
436        // This is O(N), but could be faster if we also stored a thread -> priority mapping.
437        let mut kept_all = true;
438        self.queue.retain(|_k, v| {
439            let ret = v.tid != tid;
440            kept_all = kept_all && ret;
441            ret
442        });
443        if self.yielded_skip == Some(tid) {
444            self.yielded_skip = None;
445        }
446        // A successful nonleader exec can remove one thread incarnation and
447        // admit its replacement under the same raw TID in a single scheduler
448        // drain. Do not let the replacement inherit the destroyed leader's
449        // sticky-random selection.
450        if self.sticky_random_selection == Some(tid) {
451            self.sticky_random_selection = None;
452        }
453        !kept_all
454    }
455
456    // Helper function for logging purposes.
457    fn random_range<T>(&mut self, start: T, end: T) -> T
458    where
459        T: SampleUniform + Display + PartialOrd + Copy,
460    {
461        let r = self.prng.random_range(start..end);
462        detlog!("SCHEDRAND: [{},{}) => {}", start, end, r);
463        r
464    }
465
466    /// Begin, but do not complete, a pop_transaction.  This can be committed or undone later.  But
467    /// one of those must happen before other modification operations can occur on the RunQueue.
468    ///
469    /// Postcondition: if return a `Some` value, the RunQueue enters a *locked* state where
470    /// commit or undo must happen before any other mutations to the structure.
471    pub fn tentative_pop_next(&mut self) -> Option<DetTid> {
472        assert!(!self.tentative_selection_is_exact);
473        let skip = self
474            .yielded_skip
475            .filter(|tid| self.queue.len() > 1 && self.contains_tid(*tid));
476        self.tentative_selection = match self.sched_strategy {
477            SchedHeuristic::None | SchedHeuristic::ConnectBind => self
478                .queue
479                .values()
480                .find(|value| Some(value.tid) != skip)
481                .map(|value| value.tid),
482            SchedHeuristic::Random => {
483                if self.queue.is_empty() {
484                    return None;
485                }
486
487                // If there is not Tid picked from a previous operation, let's pick one now.
488                if self.tentative_selection.is_none() {
489                    let eligible = self.queue.len() - usize::from(skip.is_some());
490                    let random_idx = self.random_range(0, eligible);
491                    self.tentative_selection = self
492                        .queue
493                        .values()
494                        .filter(|value| Some(value.tid) != skip)
495                        .nth(random_idx)
496                        .map(|value| value.tid);
497                };
498
499                self.tentative_selection
500            }
501            SchedHeuristic::StickyRandom => {
502                if self.queue.is_empty() {
503                    return None;
504                }
505
506                if self.sticky_random_selection == skip {
507                    self.sticky_random_selection = None;
508                }
509                if self.sticky_random_selection.is_none()
510                    || !self.contains_tid(self.sticky_random_selection.unwrap())
511                {
512                    let eligible = self.queue.len() - usize::from(skip.is_some());
513                    let random_idx = self.random_range(0, eligible);
514                    self.sticky_random_selection = self
515                        .queue
516                        .values()
517                        .filter(|value| Some(value.tid) != skip)
518                        .nth(random_idx)
519                        .map(|value| value.tid);
520                }
521
522                self.sticky_random_selection
523            }
524        };
525
526        self.tentative_selection
527    }
528
529    // TODO-HUMAN-REVIEW(PR-868): Review exact run-queue selection for vfork barriers.
530    /// Begin a pop transaction for one specific queued thread, bypassing the
531    /// configured scheduling heuristic without changing the thread's priority.
532    pub fn tentative_pop_tid(&mut self, tid: DetTid) -> Option<DetTid> {
533        assert!(self.tentative_selection.is_none());
534        if self.contains_tid(tid) {
535            self.tentative_selection = Some(tid);
536            self.tentative_selection_is_exact = true;
537        }
538        self.tentative_selection
539    }
540
541    /// Complete the tentative pop operation, readying the RunQueue for future operations.  This
542    /// operation is only permissible when the queue is locked, i.e. the tentative_pop has
543    /// previously returned `Some`.
544    pub fn commit_tentative_pop(&mut self) -> DetTid {
545        // Check that queue is locked and unlock it.
546        let tentative_selection = self
547            .tentative_selection
548            .take()
549            .expect("tentative_pop to already returned a `Some`");
550        let exact = std::mem::take(&mut self.tentative_selection_is_exact);
551
552        let ret = if exact {
553            let key = *self
554                .queue
555                .iter()
556                .find(|(_key, value)| value.tid == tentative_selection)
557                .map(|(key, _value)| key)
558                .unwrap();
559            self.queue.remove(&key).map(|value| value.tid)
560        } else {
561            match self.sched_strategy {
562                SchedHeuristic::None | SchedHeuristic::ConnectBind | SchedHeuristic::Random => {
563                    let key = *self
564                        .queue
565                        .iter()
566                        .find(|(_k, v)| v.tid == tentative_selection)
567                        .map(|(k, _v)| k)
568                        .unwrap();
569                    self.queue.remove(&key).map(|v| v.tid)
570                }
571                SchedHeuristic::StickyRandom => {
572                    let tid = self.sticky_random_selection.unwrap();
573                    // Probability of staying to our current thread on the next round.
574                    // If the generated random number is smaller than what we set, we switch threads.
575                    if self.random_range(0f64, 1f64) <= 1.0 - self.sticky_random_param {
576                        self.sticky_random_selection = None;
577                    }
578
579                    let key = *self
580                        .queue
581                        .iter()
582                        .find(|(_k, v)| v.tid == tid)
583                        .map(|(k, _v)| k)
584                        .unwrap();
585
586                    self.queue.remove(&key).map(|v| v.tid)
587                }
588            }
589        }
590        .expect("to always return a DetTid");
591        // The above should always return a DetTid as we peeked right before.
592        // If this invariant is violated, then it's a bug, or the queue is modified
593        // between the peek and tentative_pop and commit_tentative_pop.
594        debug_assert!(ret == tentative_selection);
595        ret
596    }
597
598    /// Commit a tentative pop for a guest turn that will actually run. A
599    /// yielded thread's one-turn exclusion is consumed only here, not by
600    /// scheduler bookkeeping turns that never unblock a guest.
601    pub fn commit_tentative_pop_completed_turn(&mut self) -> DetTid {
602        let tid = self.commit_tentative_pop();
603        self.consume_yield_exclusion();
604        tid
605    }
606
607    /// Mark that a different guest received execution after an explicit yield.
608    pub fn consume_yield_exclusion(&mut self) {
609        self.yielded_skip = None;
610    }
611
612    /// Forget the tentative pop as though it never happened.
613    pub fn undo_tentative_pop(&mut self) {
614        assert!(self.tentative_selection.is_some());
615        self.tentative_selection = None;
616        self.tentative_selection_is_exact = false;
617    }
618
619    /// Return how many things have been queued.
620    fn turn_counter(&self) -> u64 {
621        debug_assert!(self.last_back_turn >= 0);
622        debug_assert!(self.last_front_turn <= 0);
623        self.last_back_turn as u64 + self.last_front_turn.unsigned_abs()
624    }
625
626    fn check_poll_upgrade(&mut self) {
627        if self.turn_counter().is_multiple_of(POLLING_UPGRADE_INTERVAL) {
628            self.do_poll_upgrade()
629        }
630    }
631
632    /// Upgrade polled tasks to their specified normal priority.
633    #[cold]
634    fn do_poll_upgrade(&mut self) {
635        let upgrades = self
636            .queue
637            .iter()
638            .filter_map(|(k, v)| v.poll_upgrade.map(|upgd| (*k, upgd)))
639            .collect::<Vec<(PrioritizedOrder, Priority)>>();
640        // TODO(T100400409): if all polling threads are below a certain priority, this
641        // can use a range query rather than iterating over all threads in the
642        // run queue:
643        for (key, upgrade_prio) in upgrades {
644            let mut new_key = key;
645            new_key.priority = upgrade_prio;
646            let mut qval = self.queue.remove(&key).unwrap();
647            qval.poll_upgrade = None; // there's no need to upgrade to the same priority twice
648            let old = self.queue.insert(new_key, qval);
649            assert!(old.is_none()); // round robin turns should ensure uniqueness
650        }
651    }
652}
653
654impl Default for RunQueue {
655    fn default() -> Self {
656        Self::new(SchedHeuristic::None, 0, 0.0)
657    }
658}
659
660#[cfg(test)]
661mod tests {
662    use super::*;
663
664    #[test]
665    fn transport_suspend_restore_preserves_all_queue_state() {
666        for strategy in [
667            SchedHeuristic::None,
668            SchedHeuristic::Random,
669            SchedHeuristic::StickyRandom,
670        ] {
671            let tid = DetTid::from_raw(1);
672            let peer = DetTid::from_raw(2);
673            let mut queue = RunQueue::new(strategy, 391, 0.5);
674            queue.push_yielded(tid, DEFAULT_PRIORITY);
675            queue.push_back(peer, DEFAULT_PRIORITY);
676            let key = *queue
677                .queue
678                .iter()
679                .find(|(_, value)| value.tid == tid)
680                .unwrap()
681                .0;
682            queue.queue.get_mut(&key).unwrap().poll_upgrade = Some(DEFAULT_PRIORITY);
683            queue.sticky_random_selection = Some(tid);
684            let before = format!("{queue:?}");
685            let saved = queue.suspend(tid, DEFAULT_PRIORITY).unwrap();
686            assert_eq!(queue.yielded_skip, Some(tid));
687            assert_eq!(queue.sticky_random_selection, Some(tid));
688            queue.restore(saved, DEFAULT_PRIORITY);
689            assert_eq!(format!("{queue:?}"), before, "{strategy:?}");
690        }
691    }
692
693    #[test]
694    fn observation_restore_keeps_real_selection_progress_and_new_priority() {
695        for strategy in [
696            SchedHeuristic::None,
697            SchedHeuristic::Random,
698            SchedHeuristic::StickyRandom,
699        ] {
700            let tid = DetTid::from_raw(1);
701            let peer = DetTid::from_raw(2);
702            let mut queue = RunQueue::new(strategy, 391, 0.5);
703            queue.push_yielded(tid, DEFAULT_PRIORITY);
704            queue.push_back(peer, DEFAULT_PRIORITY);
705            let key = *queue
706                .queue
707                .iter()
708                .find(|(_, value)| value.tid == tid)
709                .unwrap()
710                .0;
711            queue.queue.get_mut(&key).unwrap().poll_upgrade = Some(DEFAULT_PRIORITY);
712            let saved = queue.suspend(tid, DEFAULT_PRIORITY).unwrap();
713            assert_eq!(queue.tentative_pop_next(), Some(peer));
714            assert_eq!(queue.commit_tentative_pop_completed_turn(), peer);
715            assert_eq!(queue.yielded_skip, None);
716            queue.push_back(peer, DEFAULT_PRIORITY);
717            let random_after_turn = format!("{:?}", queue.prng);
718            let sticky_after_turn = queue.sticky_random_selection;
719            let turns_after_turn = (queue.last_back_turn, queue.last_front_turn);
720            queue.restore(saved, DEFAULT_PRIORITY + 3);
721            let (restored, value) = queue
722                .queue
723                .iter()
724                .find(|(_, value)| value.tid == tid)
725                .unwrap();
726            assert_eq!(restored.turn, key.turn);
727            assert_eq!(restored.priority, DEFAULT_PRIORITY + 3);
728            assert_eq!(value.poll_upgrade, Some(DEFAULT_PRIORITY + 3));
729            assert_eq!(queue.yielded_skip, None);
730            assert_eq!(queue.sticky_random_selection, sticky_after_turn);
731            assert_eq!(
732                (queue.last_back_turn, queue.last_front_turn),
733                turns_after_turn
734            );
735            assert_eq!(format!("{:?}", queue.prng), random_after_turn);
736        }
737    }
738
739    #[test]
740    fn yielded_thread_cedes_exactly_one_turn_under_every_heuristic() {
741        for strategy in [
742            SchedHeuristic::None,
743            SchedHeuristic::ConnectBind,
744            SchedHeuristic::Random,
745            SchedHeuristic::StickyRandom,
746        ] {
747            let yielding = DetTid::from_raw(1);
748            let peer = DetTid::from_raw(2);
749            let mut queue = RunQueue::new(strategy, 0, 1.0);
750
751            queue.push_back(yielding, DEFAULT_PRIORITY - 1);
752            assert_eq!(queue.tentative_pop_next(), Some(yielding));
753            assert_eq!(queue.commit_tentative_pop(), yielding);
754
755            queue.push_yielded(yielding, DEFAULT_PRIORITY - 1);
756            queue.push_back(peer, LAST_PRIORITY);
757            assert_eq!(queue.tentative_pop_next(), Some(peer), "{strategy:?}");
758            assert_eq!(
759                queue.commit_tentative_pop_completed_turn(),
760                peer,
761                "{strategy:?}"
762            );
763
764            assert_eq!(queue.yielded_skip, None, "{strategy:?}");
765            let restored_priority = queue
766                .queue
767                .iter()
768                .find(|(_key, value)| value.tid == yielding)
769                .map(|(key, _value)| key.priority);
770            assert_eq!(
771                restored_priority,
772                Some(DEFAULT_PRIORITY - 1),
773                "{strategy:?}"
774            );
775
776            if matches!(strategy, SchedHeuristic::None | SchedHeuristic::ConnectBind) {
777                queue.push_back(peer, LAST_PRIORITY);
778                assert_eq!(queue.tentative_pop_next(), Some(yielding), "{strategy:?}");
779            }
780        }
781    }
782
783    #[test]
784    fn exact_selection_bypasses_priority_and_heuristic() {
785        for strategy in [
786            SchedHeuristic::None,
787            SchedHeuristic::ConnectBind,
788            SchedHeuristic::Random,
789            SchedHeuristic::StickyRandom,
790        ] {
791            let higher_priority = DetTid::from_raw(1);
792            let selected = DetTid::from_raw(2);
793            let mut queue = RunQueue::new(strategy, 0, 1.0);
794            queue.push_back(higher_priority, FIRST_PRIORITY);
795            queue.push_back(selected, LAST_PRIORITY);
796
797            assert_eq!(queue.tentative_pop_tid(selected), Some(selected));
798            assert_eq!(queue.commit_tentative_pop(), selected);
799            assert!(queue.contains_tid(higher_priority));
800            assert!(!queue.contains_tid(selected));
801        }
802    }
803
804    #[test]
805    fn scheduler_only_commit_does_not_consume_yield_exclusion() {
806        let yielding = DetTid::from_raw(1);
807        let peer = DetTid::from_raw(2);
808        let mut queue = RunQueue::default();
809
810        queue.push_back(yielding, DEFAULT_PRIORITY);
811        assert_eq!(queue.tentative_pop_next(), Some(yielding));
812        assert_eq!(queue.commit_tentative_pop(), yielding);
813
814        queue.push_yielded(yielding, DEFAULT_PRIORITY);
815        queue.push_back(peer, DEFAULT_PRIORITY);
816        assert_eq!(queue.tentative_pop_next(), Some(peer));
817        assert_eq!(queue.commit_tentative_pop(), peer);
818
819        queue.push_back(peer, DEFAULT_PRIORITY);
820        assert_eq!(queue.yielded_skip, Some(yielding));
821        assert_eq!(queue.tentative_pop_next(), Some(peer));
822        assert_eq!(queue.commit_tentative_pop_completed_turn(), peer);
823        assert_eq!(queue.yielded_skip, None);
824    }
825
826    #[test]
827    fn tentative_pop_in_progress_tracks_the_transaction() {
828        let a = DetTid::from_raw(1);
829        let b = DetTid::from_raw(2);
830        let mut queue = RunQueue::default();
831
832        // No selection: safe to push.
833        assert!(!queue.tentative_pop_in_progress());
834        queue.push_back(a, DEFAULT_PRIORITY);
835        queue.push_back(b, DEFAULT_PRIORITY);
836        assert!(!queue.tentative_pop_in_progress());
837
838        // Peeking a tentative selection opens the transaction; this is exactly
839        // the window in which a concurrent handler must defer its admission
840        // rather than push (a push here trips the tentative-selection guard).
841        assert_eq!(queue.tentative_pop_next(), Some(a));
842        assert!(queue.tentative_pop_in_progress());
843
844        // Committing closes it again.
845        assert_eq!(queue.commit_tentative_pop(), a);
846        assert!(!queue.tentative_pop_in_progress());
847
848        // The exact-selection form and undo path behave the same way.
849        assert_eq!(queue.tentative_pop_tid(b), Some(b));
850        assert!(queue.tentative_pop_in_progress());
851        queue.undo_tentative_pop();
852        assert!(!queue.tentative_pop_in_progress());
853    }
854
855    #[test]
856    fn removal_clears_per_incarnation_selection_state_before_tid_reuse() {
857        let tid = DetTid::from_raw(7);
858        let mut queue = RunQueue::new(SchedHeuristic::StickyRandom, 0x5107, 1.0);
859        queue.push_back(tid, DEFAULT_PRIORITY);
860
861        assert_eq!(queue.tentative_pop_next(), Some(tid));
862        queue.undo_tentative_pop();
863        assert_eq!(queue.sticky_random_selection, Some(tid));
864        // Model an earlier explicit yield by the old incarnation. Both caches
865        // are keyed only by raw TID and therefore must be cleared together.
866        queue.yielded_skip = Some(tid);
867
868        assert!(queue.remove_tid(tid));
869        assert_eq!(queue.sticky_random_selection, None);
870        assert_eq!(queue.yielded_skip, None);
871
872        queue.push_back(tid, DEFAULT_PRIORITY);
873        assert_eq!(queue.sticky_random_selection, None);
874        assert_eq!(queue.yielded_skip, None);
875    }
876}