hermit-detcore 0.4.0

Detcore: the deterministic scheduler and syscall determinization core of the Hermit execution engine.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
/*
 * Copyright (c) Meta Platforms, Inc. and affiliates.
 * All rights reserved.
 *
 * This source code is licensed under the BSD-style license found in the
 * LICENSE file in the root directory of this source tree.
 */

//! The main queue for runnable tasks.
//!
//! Tasks are selected from the queue based on 2 factors:
//! 1. Their priority
//! 2. Their round-robin order.
//!
//! The factors are compared in that order. Round-robin orders monotonically increase across
//! the entire queue; a task is assigned an order at insertion time.
//!
//! Round-robin orders can also be negative when the `push_front` method is
//! used. This "skips the line" within the priority level. This is mostly
//! relevant in non-chaos modes, where all threads have the same priority. In
//! this case, time-based events should skip the line and end-up in the front of
//! the queue.
//!
//! In addition, some priority values are reserved, such as a high priority
//! for IO eager polling.
//!
//! # Polling strategy
//!
//! We employ a polling strategy for guest threads that *would* blocked, but where we
//! don't model precisely what conditions they're waiting for.  Anywhere we have a precise
//! model of inter-thread dependencies (e.g. futexes), we can sleep a thread until we
//! encounter the matching event that will wake it.  But for Linux features that we don't
//! model 100% precisely, polling is a way to remain agnostic as to the exact
//! dependencies, but still support these blocking behaviors deterministically.
//!
//! The older DetTrace system used polling, but it would poll every time through the
//! round robin queue, which can create extremely bad performance with many threads
//! polling an unbounded number of times.  We can greatly improve the performance by
//! polling only at less-frequent, but still deterministically-defined intervals, such
//! as when we think we're out of "productive" work to do.
//!
//! Thererefore we have special handling for scheduling polling tasks. When
//! initially queued, there is exponential backoff in priority with the number
//! of attempts. After enough queueing operations are performed, however,
//! polling tasks are upgraded to their original priority to prevent complete
//! starvation. The frequency of upgrades is controlled by `POLLING_UPGRADE_INTERVAL`.

use std::cmp::Ordering;
use std::collections::BTreeMap;
use std::fmt;
use std::fmt::Display;

use rand::RngExt as _;
use rand::SeedableRng;
use rand::distr::uniform::SampleUniform;
use rand_pcg::Pcg64Mcg;

use crate::config::SchedHeuristic;
use crate::detlog;
use crate::types::DetTid;

/// The user-accessible priority of a thread. Lowest runs first.
pub type Priority = u64;

const EAGER_IO_REPOLL_PRIORITY: Priority = Priority::MIN;

/// The lowest/highest priority a thread can have.
pub const FIRST_PRIORITY: Priority = EAGER_IO_REPOLL_PRIORITY + 1;

/// The last/lowest (numerically largest) priority a thread can have.
pub const LAST_PRIORITY: Priority = 10000;

/// A high priority for threads the replayer DOES want to run.
pub const REPLAY_FOREGROUND_PRIORITY: Priority = FIRST_PRIORITY;

/// A low priority given to threads the replayer does NOT want to run.
pub const REPLAY_DEFERRED_PRIORITY: Priority = LAST_PRIORITY - 1;

/// The default priority for a thread. If chaos mode is not enabled, all threads have this
/// priority.
pub const DEFAULT_PRIORITY: Priority = 1000;

/// Whether the priority is user-accessible. We use some values for special
/// purposes; these shouldn't be set by the user.
pub fn is_ordinary_priority(prio: Priority) -> bool {
    (FIRST_PRIORITY..=LAST_PRIORITY).contains(&prio)
}

/// Deterministically transform 64 bits of entropy into a random user-settable
/// priority.
pub fn entropy_to_priority(entropy: u64) -> Priority {
    let range = LAST_PRIORITY - FIRST_PRIORITY + 1;
    let offset = entropy % range;
    FIRST_PRIORITY + offset
}

/// The round robin turn of threads within a given priority level. Lowest runs
/// first. Both negative and positive values are used to allow insertion of a
/// thread at both the "front" and "back" of a priority level.
type RoundRobinTurn = i64;

/// The key into the priority queue that uniquely determines what to run next.
/// Priorities that compare lower run first.
#[derive(Debug, Copy, Clone)]
pub struct PrioritizedOrder {
    priority: Priority,
    turn: RoundRobinTurn,
}

// These match the derived definitions, but clearly display our intention:
impl Ord for PrioritizedOrder {
    fn cmp(&self, other: &Self) -> Ordering {
        self.priority
            .cmp(&other.priority)
            .then(self.turn.cmp(&other.turn))
    }
}
impl PartialOrd for PrioritizedOrder {
    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
        Some(self.cmp(other))
    }
}
impl PartialEq for PrioritizedOrder {
    fn eq(&self, other: &Self) -> bool {
        self.cmp(other).is_eq()
    }
}
impl Eq for PrioritizedOrder {}

impl fmt::Display for PrioritizedOrder {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        write!(f, "(p: {}, t: {})", self.priority, self.turn)
    }
}

/// After queueing this many new tasks, perform a "poll upgrade," in which
/// we upgrade all outstanding polling tasks their original priority levels,
/// temporarily negating backoff behavior.
const POLLING_UPGRADE_INTERVAL: u64 = 200;

#[derive(Debug, Copy, Clone)]
struct QueueValue {
    tid: DetTid,
    /// Upgrade to this priority during polling upgrades
    poll_upgrade: Option<Priority>,
}

/// One suspended queue entry. Global yield/random-selection state stays live.
#[derive(Debug, Clone)]
pub(super) struct SuspendedRunQueueEntry {
    key: PrioritizedOrder,
    value: QueueValue,
    persistent_priority: Priority,
}

#[derive(Debug, Clone)]
pub struct RunQueue {
    /// We use a "flattened" queue (rather than a Priority -> Vec<DetTid> map)
    /// to simplify peek/pop logic: there's no need to ignore clear empty
    /// from unused priority levels vectors. This could also reduce allocator
    /// pressure. Also, each thread having a clear global key makes it easier to
    /// change their priorities after they are in the queue.
    ///
    /// Additionally, we use a TreeMap rather than a Heap to ease removing /
    /// inserting random values for poll upgrades. std::BinaryHeap would require
    /// destroying/re-allocating the entire structure to do this.
    queue: BTreeMap<PrioritizedOrder, QueueValue>,

    // We use global turn counters across all priority levels. This foregoes
    // the need for an extra data structure to track them, while also ensuring
    // unique keys for every insertion. Because of this, we never need to alter
    // a turn value when altering the priority level of a thread.
    last_back_turn: RoundRobinTurn,
    last_front_turn: RoundRobinTurn,

    /// Used to lock the queue from other changes while we are tentatively popping from it, and also
    /// cache the result.
    tentative_selection: Option<DetTid>,
    tentative_selection_is_exact: bool,

    /// A thread that explicitly yielded must not be selected again until some
    /// other runnable thread receives a turn.
    yielded_skip: Option<DetTid>,

    // TODO: The following fields need to be properly abstracted into separate types of run queues.
    /// Which scheduling strategy shall we use.
    sched_strategy: SchedHeuristic,
    prng: Pcg64Mcg,

    sticky_random_param: f64,
    sticky_random_selection: Option<DetTid>,
}

/// A multi-line print of the runqueue.
impl fmt::Display for RunQueue {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        writeln!(
            f,
            "Run queue, size={}, last_back_turn={}, last_front_turn={}:",
            self.queue.len(),
            self.last_back_turn,
            self.last_front_turn,
        )?;
        for x in self.queue.iter() {
            writeln!(f, "    {:.500?}", x)?;
        }
        Ok(())
    }
}

impl RunQueue {
    /// Create a new RunQueue.
    pub fn new(ss: SchedHeuristic, seed: u64, srp: f64) -> Self {
        detlog!("SCHEDRAND: seeding scheduler runqueue with seed {}", seed);
        Self {
            queue: BTreeMap::new(),
            // For clarity, 0 is unused so that positive/negative == back/front:
            last_back_turn: 0,
            last_front_turn: 0,
            sched_strategy: ss,
            tentative_selection: None,
            tentative_selection_is_exact: false,
            yielded_skip: None,
            prng: Pcg64Mcg::seed_from_u64(seed),
            sticky_random_param: srp,
            sticky_random_selection: None,
        }
    }

    fn push_safety_check(&self, tid: DetTid) {
        if cfg!(debug_assertions) {
            // Expensive.
            for qv in self.queue.values() {
                if qv.tid == tid {
                    panic!(
                        "Invariant violation! Tried to add {} to runqueue, but it's already present:\n {:?}",
                        tid, self
                    );
                }
            }
        }
    }

    // Return the numerically least Priority value in the run_queue, or None if the queue is empty.
    pub fn first_priority(&self) -> Option<Priority> {
        let (k, _) = self.queue.first_key_value()?;
        Some(k.priority)
    }

    /// True if any thread other than `exclude` is runnable at ordinary
    /// (non-poller) priority. This is the "deterministic work still runnable"
    /// test used to decide whether an asynchronous signal delivery must defer to
    /// guest work that was already scheduled. Read-only, so it is safe to call
    /// while a tentative_pop selection is in progress.
    pub fn has_runnable_besides(&self, exclude: DetTid) -> bool {
        self.queue
            .iter()
            .any(|(k, v)| v.tid != exclude && k.priority < LAST_PRIORITY)
    }

    /// True while a `tentative_pop`/commit transaction is underway, i.e. the
    /// daemon has peeked a selection and may have released the scheduler lock
    /// across an await. Fatal backend publication uses this read-only query to
    /// close the transaction once before waking consuming cleanup. Ordinary
    /// run-queue admissions and removals are always buffered and
    /// applied at the deterministic `step2` drain (window guaranteed closed), so
    /// a handler never needs to ask whether a window is live — see
    /// `Scheduler::admit_to_run_queue` / `deschedule_or_defer`. Read-only.
    pub fn tentative_pop_in_progress(&self) -> bool {
        self.tentative_selection.is_some()
    }

    /// Push a thread to the back of the specified priority. Return the
    /// resulting overall position in the queue.
    ///
    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
    pub fn push_back(&mut self, tid: DetTid, priority: Priority) -> PrioritizedOrder {
        assert!(self.tentative_selection.is_none());
        self.push_safety_check(tid);
        if !is_ordinary_priority(priority) {
            panic!("This is not an acceptable priority value: {}", priority);
        }
        self.push_back_inner(tid, priority, None)
    }

    /// Requeue an explicitly yielding thread at its persistent priority while
    /// excluding it from the next selection. The exclusion, rather than the
    /// queue key, makes this a one-turn operation under every heuristic.
    pub fn push_yielded(&mut self, tid: DetTid, priority: Priority) -> PrioritizedOrder {
        assert!(self.yielded_skip.is_none());
        self.yielded_skip = Some(tid);
        self.push_back(tid, priority)
    }

    /// Push a polling thread. The priority level is an exponential backoff from
    /// the given `normal_priority` value. The pushed thread will also
    /// participatein "poll upgrades" in which periodically polling threads are
    /// re-boosted to their original `normal_priority` values.
    ///
    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
    pub fn push_poller(
        &mut self,
        tid: DetTid,
        normal_priority: Priority,
        poll_attempt: u32,
    ) -> PrioritizedOrder {
        assert!(self.tentative_selection.is_none());
        self.push_safety_check(tid);
        // Exponential backoff in priority, up to LAST_PRIORITY:
        let priority = 1u64
            .checked_shl(poll_attempt)
            .and_then(|f| f.checked_mul(normal_priority))
            .unwrap_or(Priority::MAX)
            .min(LAST_PRIORITY);
        // Upgrade back to original priority:
        self.push_back_inner(tid, priority, Some(normal_priority))
    }

    fn push_back_inner(
        &mut self,
        tid: DetTid,
        priority: Priority,
        poll_upgrade: Option<Priority>,
    ) -> PrioritizedOrder {
        self.last_back_turn += 1;
        let turn = self.last_back_turn;
        let prio = PrioritizedOrder { priority, turn };
        self.push_inner(tid, prio, poll_upgrade)
    }

    /// Push a thread to the front of the specified priority. `push_back` should
    /// be used unless special circumstances call for `push_front`. Return the
    /// resulting overall position in the queue.
    ///
    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
    pub fn push_front(&mut self, tid: DetTid, priority: Priority) -> PrioritizedOrder {
        assert!(self.tentative_selection.is_none());
        self.push_safety_check(tid);
        assert!(is_ordinary_priority(priority));
        self.push_front_inner(tid, priority, None)
    }

    /// Workaround for eager io repolling: this will send the thread to the
    /// absolute front of the queue. Return the resulting overall position in
    /// the queue.
    ///
    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
    pub fn push_eager_io_repoll(&mut self, tid: DetTid) -> PrioritizedOrder {
        assert!(self.tentative_selection.is_none());
        self.push_safety_check(tid);
        let priority = EAGER_IO_REPOLL_PRIORITY;
        self.push_front_inner(tid, priority, None)
    }

    fn push_front_inner(
        &mut self,
        tid: DetTid,
        priority: Priority,
        poll_upgrade: Option<Priority>,
    ) -> PrioritizedOrder {
        self.last_front_turn -= 1;
        let turn = self.last_front_turn;
        let prio = PrioritizedOrder { priority, turn };
        self.push_inner(tid, prio, poll_upgrade)
    }

    fn push_inner(
        &mut self,
        tid: DetTid,
        prio: PrioritizedOrder,
        poll_upgrade: Option<Priority>,
    ) -> PrioritizedOrder {
        let qval = QueueValue { tid, poll_upgrade };
        let old = self.queue.insert(prio, qval);
        assert!(old.is_none()); // last_*_turn should be monotonic
        self.check_poll_upgrade();
        prio
    }

    /// Read-only: this is ok while locked by tentative_pop.
    pub fn is_empty(&self) -> bool {
        self.queue.is_empty()
    }

    /// Read-only: this is ok while locked by tentative_pop.
    pub fn len(&self) -> usize {
        self.queue.len()
    }

    /// Read-only: this is ok while locked by tentative_pop.
    pub fn tids(&self) -> impl Iterator<Item = &DetTid> {
        self.queue.values().map(|v| &v.tid)
    }

    /// Read-only: this is ok while locked by tentative_pop.
    pub fn contains_tid(&self, tid: DetTid) -> bool {
        self.tids().any(|t| t == &tid)
    }

    pub(super) fn suspend(
        &mut self,
        tid: DetTid,
        persistent_priority: Priority,
    ) -> Option<SuspendedRunQueueEntry> {
        assert!(self.tentative_selection.is_none());
        let key = *self.queue.iter().find(|(_, v)| v.tid == tid)?.0;
        let value = self.queue.remove(&key).expect("located queue entry");
        Some(SuspendedRunQueueEntry {
            key,
            value,
            persistent_priority,
        })
    }

    pub(super) fn restore(
        &mut self,
        mut entry: SuspendedRunQueueEntry,
        current_priority: Priority,
    ) {
        assert!(self.tentative_selection.is_none());
        assert!(!self.contains_tid(entry.value.tid));
        if entry.persistent_priority != current_priority {
            entry.key.priority = current_priority;
            if entry.value.poll_upgrade.is_some() {
                entry.value.poll_upgrade = Some(current_priority);
            }
        }
        assert!(self.queue.insert(entry.key, entry.value).is_none());
    }

    /// Remove `tid` from the queue, returning true if removal ocurred.
    /// Mutating operation: this will error if a tentative_pop/commit transaction is underway.
    pub fn remove_tid(&mut self, tid: DetTid) -> bool {
        assert!(self.tentative_selection.is_none());

        // This is O(N), but could be faster if we also stored a thread -> priority mapping.
        let mut kept_all = true;
        self.queue.retain(|_k, v| {
            let ret = v.tid != tid;
            kept_all = kept_all && ret;
            ret
        });
        if self.yielded_skip == Some(tid) {
            self.yielded_skip = None;
        }
        // A successful nonleader exec can remove one thread incarnation and
        // admit its replacement under the same raw TID in a single scheduler
        // drain. Do not let the replacement inherit the destroyed leader's
        // sticky-random selection.
        if self.sticky_random_selection == Some(tid) {
            self.sticky_random_selection = None;
        }
        !kept_all
    }

    // Helper function for logging purposes.
    fn random_range<T>(&mut self, start: T, end: T) -> T
    where
        T: SampleUniform + Display + PartialOrd + Copy,
    {
        let r = self.prng.random_range(start..end);
        detlog!("SCHEDRAND: [{},{}) => {}", start, end, r);
        r
    }

    /// Begin, but do not complete, a pop_transaction.  This can be committed or undone later.  But
    /// one of those must happen before other modification operations can occur on the RunQueue.
    ///
    /// Postcondition: if return a `Some` value, the RunQueue enters a *locked* state where
    /// commit or undo must happen before any other mutations to the structure.
    pub fn tentative_pop_next(&mut self) -> Option<DetTid> {
        assert!(!self.tentative_selection_is_exact);
        let skip = self
            .yielded_skip
            .filter(|tid| self.queue.len() > 1 && self.contains_tid(*tid));
        self.tentative_selection = match self.sched_strategy {
            SchedHeuristic::None | SchedHeuristic::ConnectBind => self
                .queue
                .values()
                .find(|value| Some(value.tid) != skip)
                .map(|value| value.tid),
            SchedHeuristic::Random => {
                if self.queue.is_empty() {
                    return None;
                }

                // If there is not Tid picked from a previous operation, let's pick one now.
                if self.tentative_selection.is_none() {
                    let eligible = self.queue.len() - usize::from(skip.is_some());
                    let random_idx = self.random_range(0, eligible);
                    self.tentative_selection = self
                        .queue
                        .values()
                        .filter(|value| Some(value.tid) != skip)
                        .nth(random_idx)
                        .map(|value| value.tid);
                };

                self.tentative_selection
            }
            SchedHeuristic::StickyRandom => {
                if self.queue.is_empty() {
                    return None;
                }

                if self.sticky_random_selection == skip {
                    self.sticky_random_selection = None;
                }
                if self.sticky_random_selection.is_none()
                    || !self.contains_tid(self.sticky_random_selection.unwrap())
                {
                    let eligible = self.queue.len() - usize::from(skip.is_some());
                    let random_idx = self.random_range(0, eligible);
                    self.sticky_random_selection = self
                        .queue
                        .values()
                        .filter(|value| Some(value.tid) != skip)
                        .nth(random_idx)
                        .map(|value| value.tid);
                }

                self.sticky_random_selection
            }
        };

        self.tentative_selection
    }

    // TODO-HUMAN-REVIEW(PR-868): Review exact run-queue selection for vfork barriers.
    /// Begin a pop transaction for one specific queued thread, bypassing the
    /// configured scheduling heuristic without changing the thread's priority.
    pub fn tentative_pop_tid(&mut self, tid: DetTid) -> Option<DetTid> {
        assert!(self.tentative_selection.is_none());
        if self.contains_tid(tid) {
            self.tentative_selection = Some(tid);
            self.tentative_selection_is_exact = true;
        }
        self.tentative_selection
    }

    /// Complete the tentative pop operation, readying the RunQueue for future operations.  This
    /// operation is only permissible when the queue is locked, i.e. the tentative_pop has
    /// previously returned `Some`.
    pub fn commit_tentative_pop(&mut self) -> DetTid {
        // Check that queue is locked and unlock it.
        let tentative_selection = self
            .tentative_selection
            .take()
            .expect("tentative_pop to already returned a `Some`");
        let exact = std::mem::take(&mut self.tentative_selection_is_exact);

        let ret = if exact {
            let key = *self
                .queue
                .iter()
                .find(|(_key, value)| value.tid == tentative_selection)
                .map(|(key, _value)| key)
                .unwrap();
            self.queue.remove(&key).map(|value| value.tid)
        } else {
            match self.sched_strategy {
                SchedHeuristic::None | SchedHeuristic::ConnectBind | SchedHeuristic::Random => {
                    let key = *self
                        .queue
                        .iter()
                        .find(|(_k, v)| v.tid == tentative_selection)
                        .map(|(k, _v)| k)
                        .unwrap();
                    self.queue.remove(&key).map(|v| v.tid)
                }
                SchedHeuristic::StickyRandom => {
                    let tid = self.sticky_random_selection.unwrap();
                    // Probability of staying to our current thread on the next round.
                    // If the generated random number is smaller than what we set, we switch threads.
                    if self.random_range(0f64, 1f64) <= 1.0 - self.sticky_random_param {
                        self.sticky_random_selection = None;
                    }

                    let key = *self
                        .queue
                        .iter()
                        .find(|(_k, v)| v.tid == tid)
                        .map(|(k, _v)| k)
                        .unwrap();

                    self.queue.remove(&key).map(|v| v.tid)
                }
            }
        }
        .expect("to always return a DetTid");
        // The above should always return a DetTid as we peeked right before.
        // If this invariant is violated, then it's a bug, or the queue is modified
        // between the peek and tentative_pop and commit_tentative_pop.
        debug_assert!(ret == tentative_selection);
        ret
    }

    /// Commit a tentative pop for a guest turn that will actually run. A
    /// yielded thread's one-turn exclusion is consumed only here, not by
    /// scheduler bookkeeping turns that never unblock a guest.
    pub fn commit_tentative_pop_completed_turn(&mut self) -> DetTid {
        let tid = self.commit_tentative_pop();
        self.consume_yield_exclusion();
        tid
    }

    /// Mark that a different guest received execution after an explicit yield.
    pub fn consume_yield_exclusion(&mut self) {
        self.yielded_skip = None;
    }

    /// Forget the tentative pop as though it never happened.
    pub fn undo_tentative_pop(&mut self) {
        assert!(self.tentative_selection.is_some());
        self.tentative_selection = None;
        self.tentative_selection_is_exact = false;
    }

    /// Return how many things have been queued.
    fn turn_counter(&self) -> u64 {
        debug_assert!(self.last_back_turn >= 0);
        debug_assert!(self.last_front_turn <= 0);
        self.last_back_turn as u64 + self.last_front_turn.unsigned_abs()
    }

    fn check_poll_upgrade(&mut self) {
        if self.turn_counter().is_multiple_of(POLLING_UPGRADE_INTERVAL) {
            self.do_poll_upgrade()
        }
    }

    /// Upgrade polled tasks to their specified normal priority.
    #[cold]
    fn do_poll_upgrade(&mut self) {
        let upgrades = self
            .queue
            .iter()
            .filter_map(|(k, v)| v.poll_upgrade.map(|upgd| (*k, upgd)))
            .collect::<Vec<(PrioritizedOrder, Priority)>>();
        // TODO(T100400409): if all polling threads are below a certain priority, this
        // can use a range query rather than iterating over all threads in the
        // run queue:
        for (key, upgrade_prio) in upgrades {
            let mut new_key = key;
            new_key.priority = upgrade_prio;
            let mut qval = self.queue.remove(&key).unwrap();
            qval.poll_upgrade = None; // there's no need to upgrade to the same priority twice
            let old = self.queue.insert(new_key, qval);
            assert!(old.is_none()); // round robin turns should ensure uniqueness
        }
    }
}

impl Default for RunQueue {
    fn default() -> Self {
        Self::new(SchedHeuristic::None, 0, 0.0)
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn transport_suspend_restore_preserves_all_queue_state() {
        for strategy in [
            SchedHeuristic::None,
            SchedHeuristic::Random,
            SchedHeuristic::StickyRandom,
        ] {
            let tid = DetTid::from_raw(1);
            let peer = DetTid::from_raw(2);
            let mut queue = RunQueue::new(strategy, 391, 0.5);
            queue.push_yielded(tid, DEFAULT_PRIORITY);
            queue.push_back(peer, DEFAULT_PRIORITY);
            let key = *queue
                .queue
                .iter()
                .find(|(_, value)| value.tid == tid)
                .unwrap()
                .0;
            queue.queue.get_mut(&key).unwrap().poll_upgrade = Some(DEFAULT_PRIORITY);
            queue.sticky_random_selection = Some(tid);
            let before = format!("{queue:?}");
            let saved = queue.suspend(tid, DEFAULT_PRIORITY).unwrap();
            assert_eq!(queue.yielded_skip, Some(tid));
            assert_eq!(queue.sticky_random_selection, Some(tid));
            queue.restore(saved, DEFAULT_PRIORITY);
            assert_eq!(format!("{queue:?}"), before, "{strategy:?}");
        }
    }

    #[test]
    fn observation_restore_keeps_real_selection_progress_and_new_priority() {
        for strategy in [
            SchedHeuristic::None,
            SchedHeuristic::Random,
            SchedHeuristic::StickyRandom,
        ] {
            let tid = DetTid::from_raw(1);
            let peer = DetTid::from_raw(2);
            let mut queue = RunQueue::new(strategy, 391, 0.5);
            queue.push_yielded(tid, DEFAULT_PRIORITY);
            queue.push_back(peer, DEFAULT_PRIORITY);
            let key = *queue
                .queue
                .iter()
                .find(|(_, value)| value.tid == tid)
                .unwrap()
                .0;
            queue.queue.get_mut(&key).unwrap().poll_upgrade = Some(DEFAULT_PRIORITY);
            let saved = queue.suspend(tid, DEFAULT_PRIORITY).unwrap();
            assert_eq!(queue.tentative_pop_next(), Some(peer));
            assert_eq!(queue.commit_tentative_pop_completed_turn(), peer);
            assert_eq!(queue.yielded_skip, None);
            queue.push_back(peer, DEFAULT_PRIORITY);
            let random_after_turn = format!("{:?}", queue.prng);
            let sticky_after_turn = queue.sticky_random_selection;
            let turns_after_turn = (queue.last_back_turn, queue.last_front_turn);
            queue.restore(saved, DEFAULT_PRIORITY + 3);
            let (restored, value) = queue
                .queue
                .iter()
                .find(|(_, value)| value.tid == tid)
                .unwrap();
            assert_eq!(restored.turn, key.turn);
            assert_eq!(restored.priority, DEFAULT_PRIORITY + 3);
            assert_eq!(value.poll_upgrade, Some(DEFAULT_PRIORITY + 3));
            assert_eq!(queue.yielded_skip, None);
            assert_eq!(queue.sticky_random_selection, sticky_after_turn);
            assert_eq!(
                (queue.last_back_turn, queue.last_front_turn),
                turns_after_turn
            );
            assert_eq!(format!("{:?}", queue.prng), random_after_turn);
        }
    }

    #[test]
    fn yielded_thread_cedes_exactly_one_turn_under_every_heuristic() {
        for strategy in [
            SchedHeuristic::None,
            SchedHeuristic::ConnectBind,
            SchedHeuristic::Random,
            SchedHeuristic::StickyRandom,
        ] {
            let yielding = DetTid::from_raw(1);
            let peer = DetTid::from_raw(2);
            let mut queue = RunQueue::new(strategy, 0, 1.0);

            queue.push_back(yielding, DEFAULT_PRIORITY - 1);
            assert_eq!(queue.tentative_pop_next(), Some(yielding));
            assert_eq!(queue.commit_tentative_pop(), yielding);

            queue.push_yielded(yielding, DEFAULT_PRIORITY - 1);
            queue.push_back(peer, LAST_PRIORITY);
            assert_eq!(queue.tentative_pop_next(), Some(peer), "{strategy:?}");
            assert_eq!(
                queue.commit_tentative_pop_completed_turn(),
                peer,
                "{strategy:?}"
            );

            assert_eq!(queue.yielded_skip, None, "{strategy:?}");
            let restored_priority = queue
                .queue
                .iter()
                .find(|(_key, value)| value.tid == yielding)
                .map(|(key, _value)| key.priority);
            assert_eq!(
                restored_priority,
                Some(DEFAULT_PRIORITY - 1),
                "{strategy:?}"
            );

            if matches!(strategy, SchedHeuristic::None | SchedHeuristic::ConnectBind) {
                queue.push_back(peer, LAST_PRIORITY);
                assert_eq!(queue.tentative_pop_next(), Some(yielding), "{strategy:?}");
            }
        }
    }

    #[test]
    fn exact_selection_bypasses_priority_and_heuristic() {
        for strategy in [
            SchedHeuristic::None,
            SchedHeuristic::ConnectBind,
            SchedHeuristic::Random,
            SchedHeuristic::StickyRandom,
        ] {
            let higher_priority = DetTid::from_raw(1);
            let selected = DetTid::from_raw(2);
            let mut queue = RunQueue::new(strategy, 0, 1.0);
            queue.push_back(higher_priority, FIRST_PRIORITY);
            queue.push_back(selected, LAST_PRIORITY);

            assert_eq!(queue.tentative_pop_tid(selected), Some(selected));
            assert_eq!(queue.commit_tentative_pop(), selected);
            assert!(queue.contains_tid(higher_priority));
            assert!(!queue.contains_tid(selected));
        }
    }

    #[test]
    fn scheduler_only_commit_does_not_consume_yield_exclusion() {
        let yielding = DetTid::from_raw(1);
        let peer = DetTid::from_raw(2);
        let mut queue = RunQueue::default();

        queue.push_back(yielding, DEFAULT_PRIORITY);
        assert_eq!(queue.tentative_pop_next(), Some(yielding));
        assert_eq!(queue.commit_tentative_pop(), yielding);

        queue.push_yielded(yielding, DEFAULT_PRIORITY);
        queue.push_back(peer, DEFAULT_PRIORITY);
        assert_eq!(queue.tentative_pop_next(), Some(peer));
        assert_eq!(queue.commit_tentative_pop(), peer);

        queue.push_back(peer, DEFAULT_PRIORITY);
        assert_eq!(queue.yielded_skip, Some(yielding));
        assert_eq!(queue.tentative_pop_next(), Some(peer));
        assert_eq!(queue.commit_tentative_pop_completed_turn(), peer);
        assert_eq!(queue.yielded_skip, None);
    }

    #[test]
    fn tentative_pop_in_progress_tracks_the_transaction() {
        let a = DetTid::from_raw(1);
        let b = DetTid::from_raw(2);
        let mut queue = RunQueue::default();

        // No selection: safe to push.
        assert!(!queue.tentative_pop_in_progress());
        queue.push_back(a, DEFAULT_PRIORITY);
        queue.push_back(b, DEFAULT_PRIORITY);
        assert!(!queue.tentative_pop_in_progress());

        // Peeking a tentative selection opens the transaction; this is exactly
        // the window in which a concurrent handler must defer its admission
        // rather than push (a push here trips the tentative-selection guard).
        assert_eq!(queue.tentative_pop_next(), Some(a));
        assert!(queue.tentative_pop_in_progress());

        // Committing closes it again.
        assert_eq!(queue.commit_tentative_pop(), a);
        assert!(!queue.tentative_pop_in_progress());

        // The exact-selection form and undo path behave the same way.
        assert_eq!(queue.tentative_pop_tid(b), Some(b));
        assert!(queue.tentative_pop_in_progress());
        queue.undo_tentative_pop();
        assert!(!queue.tentative_pop_in_progress());
    }

    #[test]
    fn removal_clears_per_incarnation_selection_state_before_tid_reuse() {
        let tid = DetTid::from_raw(7);
        let mut queue = RunQueue::new(SchedHeuristic::StickyRandom, 0x5107, 1.0);
        queue.push_back(tid, DEFAULT_PRIORITY);

        assert_eq!(queue.tentative_pop_next(), Some(tid));
        queue.undo_tentative_pop();
        assert_eq!(queue.sticky_random_selection, Some(tid));
        // Model an earlier explicit yield by the old incarnation. Both caches
        // are keyed only by raw TID and therefore must be cleared together.
        queue.yielded_skip = Some(tid);

        assert!(queue.remove_tid(tid));
        assert_eq!(queue.sticky_random_selection, None);
        assert_eq!(queue.yielded_skip, None);

        queue.push_back(tid, DEFAULT_PRIORITY);
        assert_eq!(queue.sticky_random_selection, None);
        assert_eq!(queue.yielded_skip, None);
    }
}